【发布时间】:2018-04-11 23:34:39
【问题描述】:
我有一个类,我在 updateLiveSockets() 方法内每 30 秒从单个后台线程填充一个映射 liveSocketsByDatacenter,然后我有一个方法 getNextSocket() 将由多个读取器线程调用以获取可用的实时套接字,它使用相同的映射来获取此信息。
public class SocketManager {
private static final Random random = new Random();
private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
private final AtomicReference<Map<Datacenters, List<SocketHolder>>> liveSocketsByDatacenter =
new AtomicReference<>(Collections.unmodifiableMap(new HashMap<>()));
private final ZContext ctx = new ZContext();
// Lazy Loaded Singleton Pattern
private static class Holder {
private static final SocketManager instance = new SocketManager();
}
public static SocketManager getInstance() {
return Holder.instance;
}
private SocketManager() {
connectToZMQSockets();
scheduler.scheduleAtFixedRate(new Runnable() {
public void run() {
updateLiveSockets();
}
}, 30, 30, TimeUnit.SECONDS);
}
// during startup, making a connection and populate once
private void connectToZMQSockets() {
Map<Datacenters, ImmutableList<String>> socketsByDatacenter = Utils.SERVERS;
// The map in which I put all the live sockets
Map<Datacenters, List<SocketHolder>> updatedLiveSocketsByDatacenter = new HashMap<>();
for (Map.Entry<Datacenters, ImmutableList<String>> entry : socketsByDatacenter.entrySet()) {
List<SocketHolder> addedColoSockets = connect(entry.getKey(), entry.getValue(), ZMQ.PUSH);
updatedLiveSocketsByDatacenter.put(entry.getKey(),
Collections.unmodifiableList(addedColoSockets));
}
// Update the map content
this.liveSocketsByDatacenter.set(Collections.unmodifiableMap(updatedLiveSocketsByDatacenter));
}
private List<SocketHolder> connect(Datacenters colo, List<String> addresses, int socketType) {
List<SocketHolder> socketList = new ArrayList<>();
for (String address : addresses) {
try {
Socket client = ctx.createSocket(socketType);
// Set random identity to make tracing easier
String identity = String.format("%04X-%04X", random.nextInt(), random.nextInt());
client.setIdentity(identity.getBytes(ZMQ.CHARSET));
client.setTCPKeepAlive(1);
client.setSendTimeOut(7);
client.setLinger(0);
client.connect(address);
SocketHolder zmq = new SocketHolder(client, ctx, address, true);
socketList.add(zmq);
} catch (Exception ex) {
// log error
}
}
return socketList;
}
// this method will be called by multiple threads to get the next live socket
// is there any concurrency or thread safety issue or race condition here?
public Optional<SocketHolder> getNextSocket() {
// For the sake of consistency make sure to use the same map instance
// in the whole implementation of my method by getting my entries
// from the local variable instead of the member variable
Map<Datacenters, List<SocketHolder>> liveSocketsByDatacenter =
this.liveSocketsByDatacenter.get();
Optional<SocketHolder> liveSocket = Optional.absent();
List<Datacenters> dcs = Datacenters.getOrderedDatacenters();
for (Datacenters dc : dcs) {
liveSocket = getLiveSocket(liveSocketsByDatacenter.get(dc));
if (liveSocket.isPresent()) {
break;
}
}
return liveSocket;
}
// is there any concurrency or thread safety issue or race condition here?
private Optional<SocketHolder> getLiveSocketX(final List<SocketHolder> endpoints) {
if (!CollectionUtils.isEmpty(endpoints)) {
// The list of live sockets
List<SocketHolder> liveOnly = new ArrayList<>(endpoints.size());
for (SocketHolder obj : endpoints) {
if (obj.isLive()) {
liveOnly.add(obj);
}
}
if (!liveOnly.isEmpty()) {
// The list is not empty so we shuffle it an return the first element
Collections.shuffle(liveOnly);
return Optional.of(liveOnly.get(0));
}
}
return Optional.absent();
}
// Added the modifier synchronized to prevent concurrent modification
// it is needed because to build the new map we first need to get the
// old one so both must be done atomically to prevent concistency issues
private synchronized void updateLiveSockets() {
Map<Datacenters, ImmutableList<String>> socketsByDatacenter = Utils.SERVERS;
// Initialize my new map with the current map content
Map<Datacenters, List<SocketHolder>> liveSocketsByDatacenter =
new HashMap<>(this.liveSocketsByDatacenter.get());
for (Entry<Datacenters, ImmutableList<String>> entry : socketsByDatacenter.entrySet()) {
List<SocketHolder> liveSockets = liveSocketsByDatacenter.get(entry.getKey());
List<SocketHolder> liveUpdatedSockets = new ArrayList<>();
for (SocketHolder liveSocket : liveSockets) { // LINE A
Socket socket = liveSocket.getSocket();
String endpoint = liveSocket.getEndpoint();
Map<byte[], byte[]> holder = populateMap();
Message message = new Message(holder, Partition.COMMAND);
boolean status = SendToSocket.getInstance().execute(message.getAdd(), holder, socket);
boolean isLive = (status) ? true : false;
// is there any problem the way I am using `SocketHolder` class?
SocketHolder zmq = new SocketHolder(socket, liveSocket.getContext(), endpoint, isLive);
liveUpdatedSockets.add(zmq);
}
liveSocketsByDatacenter.put(entry.getKey(),
Collections.unmodifiableList(liveUpdatedSockets));
}
this.liveSocketsByDatacenter.set(Collections.unmodifiableMap(liveSocketsByDatacenter));
}
}
正如你在我的课堂上看到的那样:
- 从每 30 秒运行一次的单个后台线程,我用
updateLiveSockets()方法中的所有活动套接字填充liveSocketsByDatacenter映射。 - 然后从多个线程中,我调用
getNextSocket()方法给我一个可用的活动套接字,它使用liveSocketsByDatacenter映射来获取所需的信息。
我的代码运行良好,没有任何问题,我想看看是否有更好或更有效的方法来编写它。我还想就线程安全问题或任何竞争条件(如果有的话)发表意见,但到目前为止我还没有看到任何问题,但我可能是错的。
我最担心的是updateLiveSockets() 方法和getLiveSocketX() 方法。我正在迭代liveSockets,这是 LINE A 的 SocketHolder 的 List,然后创建一个新的 SocketHolder 对象并添加到另一个新列表中。这里可以吗?
注意:SocketHolder 是一个不可变类。
【问题讨论】:
-
好吧,代码 C。洗牌 N 个元素以随机选择 1 是没有意义的。
-
代码的线程安全取决于您对代码中其他地方的
endpoints列表所做的其他操作。如果在您执行此方法时另一个线程能够修改它,则它不是线程安全的。 -
我觉得我的问题写得太快了。
endpoints列表可以包含不活动的套接字,所以在我的Code B中我正在迭代原始列表并制作一个包含所有活动套接字的新列表,然后调用 shufle 并从中获取第一个索引。在我的Code C中,我直接在列表中获取随机索引。所以看起来我需要添加那个代码来创建一个新的活动套接字列表,然后使用 ThreadLocalRandom 来获取随机索引吗? -
我再次更新了我的
Code C。现在Code B和Code C应该做同样的事情,我相信。 -
您的编辑添加了太多代码;现在很不清楚你在问什么。请将其缩减为minimal reproducible example。 (如果您以线程安全的方式调用
getLiveSocketX是不可能的 - 如果您实际上没有在代码中调用getLiveSocketX,那么如何调用它与其中的内容一样重要)。
标签: java multithreading arraylist concurrency thread-safety