【问题标题】:Java Concurrency in put/get in collections集合中 put/get 中的 Java 并发
【发布时间】:2011-11-29 15:19:28
【问题描述】:

我通过在不同线程上运行的交互式会话 + 私有源 (InputStream) 连接到外部服务。在交互式会话中,我发送传出消息并接收包含不同字段的对象的同步响应,其中一个是 ID 和确认成功或失败的“状态”。同时,我在此 ID 的私人订阅源上收到消息,并带有进一步的“状态”更新。我目前将有关每个 ID 的状态信息存储在 ConcurrentHashMap 中。我必须在这些对象上保持正确的事件序列,但我目前遇到了竞争条件,有时我会在接收和处理交互式会话上的同步响应之前处理和更新私人提要上的对象,因此让我ID 的状态已过时且不正确。

理想情况下,我希望有某种类型的带有 PutIfKeyExistOrWait (w timeout) 方法的集合,它只会在密钥存在或等待时更新值,我可以在处理私有提要上的对象时使用它。

有谁知道是否有合适的集合可用,或者可以建议我的问题的替代解决方案吗?谢谢。

【问题讨论】:

  • PutIfKeyExistOrWait 如果稍后(当 ID 已经存在)您收到 2 个异步通知,但由于竞争条件而以错误的顺序处理它们,则将无济于事。听起来唯一的防弹方法是传入消息是否具有序列号。如果他们不这样做,我想下一个最好的办法是在您收到每条消息时附加时间戳,然后根据时间戳对状态进行排序。
  • @EliAcherkan 谢谢。私人订阅源上的异步消息是在单个线程中接收的,因此我可以“保证”这些消息的正确顺序。只有同步响应与异步消息会导致问题。时间戳也不起作用。消息通常在同一毫秒内收到,并且无法保证在私人提要上的消息之前收到同步响应,因此存在问题。
  • 对不起,我能想到的唯一选择是让交互式会话有一个 synchronized 块,用于检查密钥是否已由私人提要插入到地图中,并插入/更新相应的状态。
  • 你的问题我还是不清楚。您向该服务发送一条消息,然后您会从该服务获得响应并从提要中获得“更新”。 ID 是否像消息号一样唯一? (这意味着您从服务中获得一条消息,并从具有此 ID 的提要中获得一条消息。)您正在等待处理,直到您从服务和提要中获得一对相关消息?您能否从给定 ID 的提要中获取额外消息?
  • @toto2 ID 就像一个订单 ID,我需要始终保持订单的准确状态。我在交互式会话中获得了一个同步响应,并且(通常)在私人提要上获得了 ID 的多个状态更新。竞争条件发生在初始交互式响应与同时接收的私有源更新之间。私人提要上的任何后续更新(几秒钟或几分钟后)都不是问题。

标签: java collections concurrency


【解决方案1】:

您可以尝试将处理这种情况的逻辑封装到地图的值中,如下所示:

  • 如果 feed 线程是第一个为特定 id 添加值的线程,则该值被认为是不完整的,线程会一直等待直到完成
  • 如果交互式会话线程不是第一个添加值的,它会将不完整的值标记为完整
  • 不完整的值在从地图中获取时被视为不存在

这个解决方案是基于putIfAbsent()的原子性。

public class StatusMap {
    private Map<Long, StatusHolder> map = new ConcurrentHashMap<Long, StatusHolder>();

    public Status getStatus(long id) {
        StatusHolder holder = map.get(id);
        if (holder == null || holder.isIncomplete()) {
            return null;
        } else {
            return holder.getStatus();
        }
    }

    public void newStatusFromInteractiveSession(long id, Status status) {
        StatusHolder holder = StatusHolder.newComplete(status);
        if ((holder = map.putIfAbsent(id, holder)) != null) {
            holder.makeComplete(status); // Holder already exists, complete it
        } 
    }

    public void newStatusFromFeed(long id, Status status) {
        StatusHolder incomplete = StatusHolder.newIncomplete();
        StatusHolder holder = null;
        if ((holder = map.putIfAbsent(id, incomplete)) == null) {
            holder = incomplete; // New holder added, wait for its completion
            holder.waitForCompletion();
        }
        holder.updateStatus(status);
    }
}

public class StatusHolder {
    private volatile Status status;
    private volatile boolean incomplete;
    private Object lock = new Object();

    private StatusHolder(Status status, boolean incomplete) { ... }

    public static StatusHolder newComplete(Status status) {
        return new StatusHolder(status, false);
    }

    public static StatusHolder newIncomplete() {
        return new StatusHolder(null, true);
    }

    public boolean isIncomplete() { return incomplete; }

    public void makeComplete(Status status) {
        synchronized (lock) {
            this.status = status;
            incomplete = false;
            lock.notifyAll();
        }
    }

    public void waitForCompletion() {
        synchronized (lock) {
            while (incomplete) lock.wait();
        }
    }
    ...
}

【讨论】:

  • 谢谢。 +1 雄心勃勃的答案。您的建议与我目前正在做的一种解决方法是一致的,我尝试根据来源和状态来锻炼序列。该映射涉及一些相当棘手的场景,例如在某些情况下,私人提要没有更新,所以我想看看是否有我忽略的更清洁、更简单的解决方案。
【解决方案2】:

您已经有一些 ConcurrentHashMap iDAndStatus 来存储 ID 和最新状态。但是,我只会让处理服务的线程在该映射中创建一个新条目。

当消息从提要到达时,如果iDAndStatus 中已经存在该ID,它只是修改状态。如果密钥不存在,只需将 ID/状态更新临时存储在其他数据结构中,pendingFeedUpdates

每次在 iDAndStatus 中创建新条目时,请检查 pendingFeedUpdates 以查看是否存在针对新 ID 的一些更新。

我不确定pendingFeedUpdates 使用什么同步数据结构:您需要按 ID 检索,但每个 ID 可能有很多消息,并且您希望保持消息的顺序。也许是一个同步的 HashMap,将每个 ID 与某种类型的同步有序队列相关联?

【讨论】:

  • 谢谢。我将采用您的建议的一个版本。当我在交互式会话上处理同步消息时,我会检查我是否已经从私人提要中获得了相同 ID 的状态更新。在这种情况下,我只是忽略同步响应并且不更新状态。我将所有消息存储在一个单独的集合中以保留订单跟踪,但我可以忍受该列表中的轻微混乱。感谢大家提供出色的 cmets 和答案。
【解决方案3】:

我建议您查看 Collections.getSynchronized 集合:http://docs.oracle.com/javase/1.4.2/docs/api/java/util/Collections.html#synchronizedList%28java.util.List %29

这可能会解决您的问题,另一个选项取决于调用的方式,使方法是同步方法,允许线程安全执行并确保事务的原子性。见http://docs.oracle.com/javase/tutorial/essential/concurrency/syncmeth.html

第三个选项是根据您要实现的目标,采用乐观或悲观的方法在应用程序中实施并发管理控制。这是 3 个中最复杂的一个,但如果与之前的选项结合使用,您将获得更大的控制权。

这真的取决于你的具体实现。

【讨论】:

  • 谢谢。我不确定普通的同步或锁如何解决我的问题。我仍然会遇到竞争情况,其中私人提要上的消息可以在同步响应之前获得锁定。我应该指出,在我在交互式会话上发送请求之前我不知道 ID,否则我可能会在收到并处理同步响应之前锁定与 ID 对应的密钥。
  • 您可以尝试乐观锁定,而不是锁定记录,方法是检测冲突并分别处理这些冲突agiledata.org/essays/concurrencyControl.html#OptimisticLocking 在增量行版本中,我个人偏好有几种方法可以检测并发冲突。因此,如果行版本不等于您的数据库中的内容,则让交互式发回 ID 和行版本,这将允许您单独处理它。
猜你喜欢
  • 1970-01-01
  • 2012-03-02
  • 2013-11-15
  • 1970-01-01
  • 2021-07-10
  • 2013-03-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多