【问题标题】:Return an object from an array list in a thread safe way?以线程安全的方式从数组列表中返回一个对象?
【发布时间】: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


【解决方案1】:

代码 B 或 C 都不是线程安全的。

代码 B

当您在 enpoints 列表上进行迭代以复制它时,没有什么能阻止另一个线程进行修改,即元素被添加和/或删除。

代码 C

假设endpoints 不为空,您将对列表对象进行三个调用:isEmpty、size 和get。从并发的角度来看有几个问题:

  1. 基于参数的类型List&lt;SocketHolder&gt;,不能保证这些方法强制对列表的内部更改传播到其他线程(内存可见性),不考虑竞争条件(如果列表是在您的线程执行此函数之一时修改)。

  2. 假设列表 endpoints 提供了之前描述的保证 - 例如它已被Collections.synchronizedList() 包裹。在这种情况下,仍然缺少线程安全性,因为在对isEmpty、size 和get 的每次调用之间,可以在您的线程执行getLiveSocketX 方法时修改该列表。这可能会使您的代码使用列表的过时状态。例如,您可以使用endpoints.size() 返回的大小,该大小不再有效,因为已将元素添加到列表中或从列表中删除。

编辑 - 代码更新后

在您提供的代码中,乍一看似乎是:

    1234563 987654338@.
  1. 您使用AtomicReference 将Datacenters 映射到感兴趣的套接字列表。这个AtomicReference 的后果是强制从这个映射到所有列表及其元素的内存可见性。这意味着,副作用,您可以避免“生产者”和“消费者”线程之间的内存不一致(分别执行 updateLiveSockets 和 getLiveSocket)。不过,您仍然面临竞争条件 - 让我们想象一下 updateLiveSockets 和 getLiveSocket 同时运行。考虑一个套接字S,它的状态只是从活动切换到关闭。 updateLiveSockets 将看到套接字S 的状态为非活动状态,并相应地创建一个新的SocketHolder。但是,同时运行的getLiveSocket 将看到S 的过时状态 - 因为它仍将使用updateLiveSockets 正在重新创建的套接字列表。

  2. updateLiveSockets方法上使用的synchronized关键字在这里不给你任何保证,因为代码的其他部分也不是synchronized。

总而言之,我想说:

  1. getLiveSocketX 的代码不是 固有线程安全的;
  2. 但是,复制列表的方式会阻止并发修改;并且您将从AtomicReference 的副作用 中受益,从而在内存可见性方面获得最低限度的保证,人们期望在从另一个线程生成之后,在getNextSocket 中获得一致的套接字列表;
  3. 您仍会遇到 (2) 中所述的竞争条件,但这可能没问题,具体取决于您希望赋予 getLiveSocket 和 getNextSocket 方法的规范 - 您可以接受 @ 返回的一个套接字987654361@ 不可用并具有重试机制。

话虽如此,我会彻底审查和重构代码,以展示更易读和更明确的线程安全消费者/生产者模式。使用 AtomicReference 和单个 synchronized 时应格外小心,在我看来,它们使用不当 - 尽管 很好 AtomicReference 确实对您有所帮助,如前所述。

【讨论】:

  • hmmm 我没有理解我们在这里遇到的比赛条件是什么?您认为解决这个生产者消费者问题的最佳方法是什么?
  • @user1950349:今晚(10-12 小时后)会尝试为您提供一些代码
  • @user1950349:所以我查看了您的代码并尝试查看可以更改的内容。我在这里可能是错的,但由于updateLiveSockets 似乎只更新套接字持有者的状态,你为什么要创建新的SocketHolder 对象及其容器?是因为你想要一个不可变的SocketHolder 对象吗?在这种特殊情况下,允许SocketHolder 的状态发生突变(显然以线程安全的方式)是有意义的。有什么反对的理由吗?
  • 是的,我这样做只是为了拥有一个不可变的 SocketHolder 对象。我做错了什么吗?
  • 没有错,将类设计为不可变通常是一个好习惯。在目前的情况下,我对你的代码库的其余部分知之甚少(例如,哪些代码正在使用SocketHolder 对象,它们是如何使用的,如果你可以改变状态,是否会有潜在的缺陷等),但它如果您允许更新 SocketHolder 的状态 - 例如使用 volatile boolean status,您似乎可以避免在 SocketManager 中复制地图和套接字列表。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-10-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多