【问题标题】:concurrency object for writer-takes-precedence-over-readerwriter-takes-precedence-over-reader 的并发对象
【发布时间】:2009-08-07 01:52:52
【问题描述】:

我正在寻找可以帮助以下用例的并发对象:

  • 线程/实体:1 个发布者(唯一),0 个读者
  • 发布者频繁/不规律地更新数据结构,需要快速更新数据结构,并且延迟时间最短
  • 每个读取器都具有对数据结构的读取访问权限(通过某种不允许写入的方式,或者因为读取器隐含承诺不会更改数据)
  • 每个读者都愿意反复尝试访问数据结构,只要它能够检测到发布者何时到来并更改它,因为它知道它最终会获得足够的时间来阅读它需要的内容。

有什么建议吗?我可以使用ReentrantReadWriteLock,但有点担心阻止发布者。我宁愿出版商能够破坏读者阅读的机会,也不愿让读者能够阻止出版商。

发布者线程:

 PublisherSignal ps = new PublisherSignal();
 publishToAllReaders(ps.getReaderSignal());

   ...

 while (inLoop())
 {
      ps.beginEdit();
      data.setSomething(someComputation());
      data.setSomethingElse(someOtherComputation());
      ps.endEdit();

      doOtherStuff();
 }

读者话题:

 PublisherSignal.Reader rs = acquireSignalFromPublisher();

      ...

 while (inLoop())
 {
      readDataWhenWeGetAChance();

      doOtherStuff();
 }

      ...

 public readDataWhenWeGetAChance()
 {
      while (true)
      {
           rs.beginRead();
           useData(data.getSomething(), data.getSomethingElse());
           if (rs.endRead())
           {
               // we get here if the publisher hasn't done a beginEdit()
               // during our read.
               break;
           }

           // darn, we have to try again. 
           // might as well yield thread if appropriate
           rs.waitToRead();
      }
 }

编辑:在更高的层次上,我想做的是让发布者每秒更改数千次数据,然后让读者以更慢的速度显示最新更新(5 -10 次/秒)。我会使用 ConcurrentLinkedQueue 来发布已发生更新的事实,除了 (a) 同一项目上可能有数百个更新,我想合并这些更新,因为似乎必须重复复制大量数据 就像浪费是一个性能问题,并且(b)有多个读者似乎排除了一个队列......我想我可以有一个主代理读者并让它通知每个真正的读者。

【问题讨论】:

    标签: java concurrency


    【解决方案1】:

    为什么不使用BlockingQueue

    您的发布者可以独立于正在阅读的内容写入此队列。读者(类似地)可以从队列中取出东西,而不必担心阻塞作者。线程安全由队列处理,因此 2 个线程可以写入/读取,无需进一步同步等。

    来自链接的文档:

     class Producer implements Runnable {
       private final BlockingQueue queue;
       Producer(BlockingQueue q) { queue = q; }
       public void run() {
         try {
           while(true) { queue.put(produce()); }
         } catch (InterruptedException ex) { ... handle ...}
       }
       Object produce() { ... }
     }
    
     class Consumer implements Runnable {
       private final BlockingQueue queue;
       Consumer(BlockingQueue q) { queue = q; }
       public void run() {
         try {
           while(true) { consume(queue.take()); }
         } catch (InterruptedException ex) { ... handle ...}
       }
       void consume(Object x) { ... }
     }
    

    【讨论】:

    • 我会让一个消费者获取数据,合并它,然后将其传递给多个下游消费者(可能通过其他队列?)
    【解决方案2】:

    嗯...我想我的绊脚石是共享数据结构本身...我一直在使用类似的东西

     public class LotsOfData
     {
          int fee;
          int fi;
          int fo;
          int fum;
    
          long[] other = new long[123];
    
          /* other fields too */
     }
    

    发布者经常更新数据,但一次只更新一个字段。

    听起来也许我应该找到一种有利于使用生产者-消费者队列的方式来序列化更新:

     public class LotsOfData
     {
          enum Field { FEE, FI, FO, FUM };
          Map<Field, Integer> feeFiFoFum = new EnumMap<Field, Integer>();
    
          long[] other = new long[123];
    
          /* other fields too */
     }
    

    然后将更改项发布到队列中,例如 feeFiFoFum 字段的 (FEE, 23) 和 other 数组的 (33, 1234567L)。 (出于性能原因,几乎可以肯定 Bean 类型的反射已被淘汰。)

    不过,我似乎被简单地让出版商写任何它想要的东西,并且知道(最终)有时间让读者进入并获得一致的集合数据,如果它有一个标志,它可以用来判断数据是否已被修改。

    更新: 有趣的是,我尝试了这种方法,使用一个 ConcurrentLinkedQueue 的 Mutation 对象(仅存储 1 次更改所需的状态),用于类似于上面第一个 LotusOfData 的类(4 个 int 字段和一个27 个 long 的数组),一个生产者在大约 10000 个批次之间使用 Thread.sleep(1) 产生总共 1000 万个突变,一个消费者每 100 毫秒检查一次队列,并消耗存在的任何突变。我以多种方式进行了测试:

    • 测试框架内的空操作(仅循环 1000 次,调用 Thread.sleep(1),并检查是否使用空对象):在运行 jre6u13 的 3GHz Pentium 4 上为 1.95 秒。
    • 测试操作 1 -- 仅在生产者端创建和应用突变:4.3 秒
    • 测试操作 2 -- 在生产者端创建和应用突变,在创建时将每个突变放入队列:12 秒

    所以创建每个突变对象平均需要 230 纳秒,平均 770 纳秒将每个突变对象入队/出列到生产者的队列中并在消费者中拉出(执行原始类型的突变的时间似乎可以忽略不计,与对象创建和队列操作相比,应该是这样)。我猜还不错,它为我提供了一些指标来估计这种方法的性能成本。

    【讨论】:

      猜你喜欢
      • 2017-06-10
      • 2020-05-17
      • 1970-01-01
      • 1970-01-01
      • 2017-04-25
      • 2012-06-02
      • 2022-12-02
      • 2017-07-04
      相关资源
      最近更新 更多