【问题标题】:Coordinating multiple concurrent queues协调多个并发队列
【发布时间】:2011-01-06 16:47:31
【问题描述】:

我目前有一个使用BlockingQueue 作为数据存储的并发队列实现。我现在需要引入第二种具有更高优先级的对象,将我引向原始队列的饥饿/优先级队列。因此,我们正在处理由多个线程生成的 A 类型和 B 类型的对象。任何类型 B 的对象都应该在类型 A 的对象之前处理,但除了必须保持 FIFO 顺序之外。因此,如果插入 { 1A, 1B, 2A, 3A, 2B },则顺序应为 {1B, 2B, 1A, 2A, 3A}

我尝试了单个PriorityBlockingQueue 将类型 B 推到前面,但我无法保持 FIFO 要求(相同类型的项目之间没有自然顺序)。

我的下一个想法是使用两个并发队列。在协调两个队列之间的访问时,我正在寻找常见的问题或注意事项。理想情况下,我想做这样的事情:

   public void add(A a)
   {
       aQueue.add(a);
   }
   public void add(B b)
   {
       bQueue.add(b);
   }

   private void consume() 
   {
       if(!bQueue.isEmpty())
          process(bQueue.poll());
       else if(!aQueue.isEmpty())
          process(aQueue.poll());
   }

如果两个队列都是ConcurrentLinkedQueue(或在此处插入更合适的结构),我是否需要任何同步或锁定?注意我有很多生产者,但只有一个消费者(单线程ThreadPoolExecutor)。

编辑:如果 B 在 isEmpty() 检查后进入,则可以处理 A 并在下一次 consume() 调用时处理它。

【问题讨论】:

  • 在处理 A 对象之前处理 B 对象有多重要?在您当前的代码中,如果在调用 isEmpty() 之后插入了 B 对象,您会错过它。有两个不同的答案,取决于你必须做什么......
  • @Jonathan 不重要。我打算添加它并以某种方式将其排除在外。问题已更新。

标签: java concurrency queue


【解决方案1】:

我不确定我的情况是否正确,但我认为应该可以使用单个队列来解决这个问题。

您说您的对象(在队列中)应该按自然顺序和类型进行比较。如果没有自然顺序,只需有一个序列生成器(即AtomicLong),它将为您的对象提供唯一的、始终递增的队列 ID。从AtomicLong 获取数据应该不会花时间,除非您处于纳秒级的世界。

所以你的Comparator.compare 应该是这样的:

1) 检查对象类型。如果不同(A VS B),返回 1/-1。否则,请参见下文
2)检查身份证。保证不一样。

如果您无法更改对象(A 和 B),您仍然可以将它们包装到另一个包含该 ID 的对象中。

【讨论】:

  • 我唯一关心的包装和使用增量是通过这里推送的对象的绝对数量。
  • 这是一个权衡 :) 你得到了简单性并失去了一些 GC 时间。这完全取决于队列的强度、您的内存限制等。做一个测量并决定......我会选择简单,但我不知道完整的上下文。
  • 我要试一试。我绝对信任 java 并发开发人员,而不是我放在一起的任何东西。现有的实现已经在生产中锤炼多年,我不想彻底改变它。
  • 我也喜欢在不可避免地添加新类型时我可以处理它们(在 A 之前但在 B 之后执行 C - ack!)
【解决方案2】:

根据您的要求,您的方法应该可以正常工作。

我将对您的代码进行以下更改。这将阻塞您的 getNextItem() ,直到其中一个队列为您返回一个对象。

private Object block = new Object();

public void add(A a)
{
    synchronized( block )
    {
        aQueue.add( a );
        block.notifyAll();
    }
}

public void add(B b)
{
    synchronized( block )
    {
        bQueue.add( b );
        block.notifyAll();
    }
}

private Object consume() 
{
    Object value = null
    synchroinzed( block )
    {
        while ( return == null )
        {
            value = bQueue.poll();
            if ( value == null ) value = aQueue.poll();
            if ( value == null ) block.wait();
        }
     }

   return value;
}

【讨论】:

  • 实施也不错。在我的情况下,我还有其他避免阻塞的原因,但如果我的生产者和消费者的耦合度较低,这看起来不错。
【解决方案3】:

您需要同步一个 getQueueItem() 方法。

A 和 B 实现 QueueItem 接口。

   public void add(A a)
   {
       aQueue.add(a);
   }
   public void add(B b)
   {
       bQueue.add(b);
   }

   private void consume() 
   {
       process(getNextItem());
   }

    private QueueItem getNextItem()
    {
       synchronized(bQueue) {
          if(!bQueue.isEmpty()) return bQueue.poll();
          return aQueue.poll();
       }
    }

【讨论】:

  • 感谢您的回答。我将尝试@mindas 建议的单个队列,但我最终可能会回到两个队列。
猜你喜欢
  • 1970-01-01
  • 2017-05-18
  • 2020-04-21
  • 2015-10-19
  • 2017-06-27
  • 2012-03-17
  • 1970-01-01
  • 2023-04-09
  • 1970-01-01
相关资源
最近更新 更多