【问题标题】:Multiple threads in java appending to a queuejava中的多个线程附加到队列
【发布时间】:2013-09-03 17:33:58
【问题描述】:

我有多个线程正在运行,需要附加到同一个队列。这个队列被分成多个变量,所以有效地,我正在调用一个函数,该函数在某个位置 i 处附加到每个变量。从http://www.youtube.com/watch?v=8sgDgXUUJ68&feature=c4-overview-vl&list=PLBB24CFB073F1048E 显示,我为每种方法添加了锁,如此处所示。 http://www.caveofprogramming.com/java/java-multiple-locks/ 。为数组中的 100,000 个对象创建锁似乎并不有效。在java中,假设我有一个管理大量对象队列的对象。如何正确同步 addToQueue,而不会通过同步方法而牺牲性能,只是我要附加到的 floatQueue 和 intQueue 中的位置?

class QueueManager {
    private float[] floatQueue;
    private int[] intQueue;

    private int currentQueueSize;
    private int maxQueueSize;

    QueueManager(int sizeOfQueue) {
        intQueue = new int[sizeOfQueue];
        floatQueue = new float[sizeOfQueue];

        currentQueueSize = 0;
        maxQueueSize = sizeOfQueue;
    }

    public boolean addToQueue(int a, float b) {
        if (currentQueueSize == maxQueueSize) {
            return false;
        }
        else{
            intQueue[currentQueueSize] = a;
            floatQueue[currentQueueSize] = b;

            currentQueueSize++;

            return true;
        }
    }
}

【问题讨论】:

  • 我想到的一种方法是将队列管理器的实例添加到每个线程,并将变量更改为静态易失性数组。有没有更好的方法来做到这一点?
  • 听起来你想要的是一个使用 offer 方法的 BlockingQueue?
  • 附带说明,不幸的是,java 中没有无锁队列实现
  • 有什么方法可以锁定特定位置吗?如果是这样,我是否必须为每个职位创建一个锁?
  • 您需要在一个原子操作中更新该位置的位置和内容,以使其线程安全。这显然需要某种锁定——我认为没有任何方法可以在不使用本机实现的情况下提高 BlockingQueue 实现的性能。

标签: java multithreading locks


【解决方案1】:

【讨论】:

  • 我希望它是非阻塞的。
【解决方案2】:

java.util.concurrent.ConcurrentLinkedQueue 是一个非阻塞、无锁的队列实现。它使用比较和交换指令而不是锁定。

【讨论】:

  • 速度怎么样?
  • 在我的系统上最多支持四个并行线程,将一百万个项目添加到同一个队列需要不到 2.5 秒。如果你想自己测试一下,我这里有基准测试程序:gist.github.com/kyledewey/6399477
【解决方案3】:

您最好使用Blocking Queue,而不是重新发明轮子。如果您这样做是为了练习,那么您可以使用两个锁,一个用于putting,一个用于getting。具体参考Blocking queue的源码。

附带说明一下,确保正确的并发性很棘手,因此不要在部署应用程序中使用队列,而应依赖 Java 的实现。

【讨论】:

    【解决方案4】:

    要保持并行数组“同步”,您需要同步。这是保持对数组及其关联索引的修改原子性的唯一方法。

    您可能会考虑使用不可变元组对象而不是并行数组,为您的intfloat 提供final 字段,然后将它们添加到BlockingQueue 实现中,例如ArrayBlockingQueue

    如果有很多流失(频繁创建和丢弃对象),synchronized 可能会比元组执行得更好。您需要对它们进行概要分析才能确定。

    【讨论】:

      【解决方案5】:

      我发现一个可能的解决方案是使用 AtomicIntegers,它类似于锁定,但运行的级别比同步的低得多。这是我的代码的更新副本,使用了 AtomicInteger。请注意,它仍然需要测试,但理论上应该足够了。

      import java.util.concurrent.atomic.AtomicInteger;
      
      class ThreadSafeQueueManager {
          private float[] floatQueue;
          private int[] intQueue;
      
          private AtomicInteger currentQueueSize;
          private int maxQueueSize;
      
          QueueManager(int sizeOfQueue) {
              intQueue = new int[sizeOfQueue];
              floatQueue = new float[sizeOfQueue];
      
              currentQueueSize = new AtomicInteger(0);
              maxQueueSize = sizeOfQueue;
          }
      
          public boolean addToQueue(int a, float b) {
              if (currentQueueSize.get() == maxQueueSize) {
                  return false;
              }
              else{
                  try {
                      // Subtract one so that the 0th position can be used
                      int position = currentQueueSize.incrementAndGet() - 1;
                      intQueue[position] = a;
                      floatQueue[positions] = b;
      
                      return true;
                  }
                  catch (ArrayIndexOutOfBoundsException e) { return false;}
              }
          }
      }
      

      另外,作为参考,值得一读

      https://www.ibm.com/developerworks/java/library/j-jtp11234/
      http://www.ibm.com/developerworks/library/j-jtp04186/

      【讨论】:

      • 我很抱歉这么说,但这很好地说明了为什么这种事情很难。假设代码已被充分更改以使其能够编译,那么在检查当前队列大小和将数据添加到队列之间仍然会在 addToQueue(int, float) 中发生数据竞争。
      • 我不太确定这是否属实。在这里stackoverflow.com/questions/10077937/… 发现了一个类似的问题,但我只希望 addToQueue 的每次调用只占据一个位置,因为 floatQueue 和 intQueue 中的那个位置取决于 incrementAndGet(),为什么会有问题?这些写入不需要是顺序的,唯一的要求是这些数组中的相同位置在每个工作负载中只写入一次。
      • 所以如果我用 try catch 包裹 else 块,那么我就安全了吗?
      • 假设addToQueue(...)中的初始测试改成if (currentQueueSize.get() == maxQueueSize)这样代码就可以编译了,那么可能会出现以下情况: 1.线程进入方法并且测试通过,但是在int position = currentQueueSize.incrementAndGet() - 1 被执行,它被抢占。 2.抢占线程也将通过测试,因为计数器还没有增加。该线程完成该方法。 3. 一段时间后,来自 (1) 的原始线程被重新安排。如果数组现在已满,则会出现ArrayIndexOutOfBounds
      • 我认为它会以某种解决方法的方式出现。您也应该将初始条件从== 更改为>=。这也有点难看,因为currentQueueSize 在解决方法的情况下不一定具有当前的队列大小。 A 可能是实现一个方法int incrementAndGetWithLimit(AtomicInteger val, int limit),它只会将计数器递增到一个点。您可以查看AtomicInteger#incrementAndGet() 以获取有关如何完成此操作的一些灵感docjar.com/html/api/java/util/concurrent/atomic/…
      猜你喜欢
      • 2014-07-27
      • 2014-06-12
      • 1970-01-01
      • 2014-11-20
      • 2018-03-09
      • 2013-06-18
      • 2010-10-28
      • 1970-01-01
      • 2013-01-14
      相关资源
      最近更新 更多