【问题标题】:Lock-free container with flipping buffers带有翻转缓冲区的无锁容器
【发布时间】:2019-03-11 09:36:47
【问题描述】:

对于我的一个项目,它必须支持并发读取和写入,我需要一个能够缓冲项目的容器,直到消费者一次获取每个当前缓冲的项目。由于生产者应该能够生成数据,无论消费者是否读取当前缓冲区,我想出了一个自定义实现,在AtomicReference 的帮助下,将每个条目添加到支持ConcurrentLinkedQueue 直到执行翻转,这导致在存储具有空队列和元数据的新条目时要返回的当前条目以原子方式存储在 AtomicReference 中。

我想出了一个解决方案,例如

public class FlippingDataContainer<E> {

  private final AtomicReference<FlippingDataContainerEntry<E>> dataObj = new AtomicReference<>();

  public FlippingDataContainer() {
    dataObj.set(new FlippingDataContainerEntry<>(new ConcurrentLinkedQueue<>(), 0, 0, 0));
  }

  public FlippingDataContainerEntry<E> put(E value) {
    if (null != value) {
      while (true) {
        FlippingDataContainerEntry<E> data = dataObj.get();
        FlippingDataContainerEntry<E> updated = FlippingDataContainerEntry.from(data, value);
        if (dataObj.compareAndSet(data, updated)) {
          return merged;
        }
      }
    }
    return null;
  }

  public FlippingDataContainerEntry<E> flip() {
    FlippingDataContainerEntry<E> oldData;
    FlippingDataContainerEntry<E> newData = new FlippingDataContainerEntry<>(new ConcurrentLinkedQueue<>(), 0, 0, 0);
    while (true) {
      oldData = dataObj.get();
      if (dataObj.compareAndSet(oldData, newData)) {
        return oldData;
      }
    }
  }

  public boolean isEmptry() {
    return dataObj.get().getQueue().isEmpty();
  }
}

由于需要将当前值推送到后备队列,因此现在需要格外小心。 from(data, value) 方法的当前实现看起来像这样:

static <E> FlippingDataContainerEntry<E> from(FlippingDataContainerEntry<E> data, E value) {
  Queue<E> queue = new ConcurrentLinkedQueue<>(data.getQueue());
  queue.add(value);
  return new FlippingDataContainerEntry<>(queue,
      data.getKeyLength() + (value.getKeyAsBytes() != null ? value.getKeyAsBytes().length : 0),
      data.getValueLength() + (value.getValueAsBytes() != null ? value.getValueAsBytes().length : 0),
      data.getAuxiliaryLength() + (value.getAuxiliaryAsBytes() != null ? value.getAuxiliaryAsBytes().length : 0));
}

由于其他线程可能会在该线程执行更新之前更新值,因此我需要在每次写入尝试时复制实际队列,否则即使原子操作也会将条目添加到共享队列中参考资料无法更新。因此,简单地将值添加到共享队列可能会导致值条目被多次添加到队列中,而实际上它应该只出现一次。

复制整个队列是一项相当昂贵的任务,所以我只是设置当前队列而不是在 from(data, value) 方法中复制队列,而不是将 value 元素添加到执行块中的共享队列中发生更新:

public FlippingDataContainerEntry<E> put(E value) {
  if (null != value) {
    while (true) {
      FlippingDataContainerEntry<E> data = dataObj.get();
      FlippingDataContainerEntry<E> updated = FlippingDataContainerEntry.from(data, value);
      if (data.compareAndSet(data, updated)) {
        updated.getQueue().add(value);
        return updated;
      }
    }
  }
  return null;
}

from(data, value)内我现在只设置队列不直接添加value元素

static <E> FlippingDataContainerEntry<E> from(FlippingDataContainerEntry<E> data, E value) {
  return new FlippingDataContainerEntry<>(data.getQueue(),
      data.getKeyLength() + (value.getKeyAsBytes() != null ? value.getKeyAsBytes().length : 0),
      data.getValueLength() + (value.getValueAsBytes() != null ? value.getValueAsBytes().length : 0),
      data.getAuxiliaryLength() + (value.getAuxiliaryAsBytes() != null ? value.getAuxiliaryAsBytes().length : 0));
}

虽然与复制队列的代码相比,这允许运行测试快 10 倍以上,但它也经常使消费测试失败,因为现在可能在消费者线程翻转队列后立即将 value 元素添加到队列中排队并处理数据,因此似乎并非所有项目都已被消耗。

现在的实际问题是,是否可以避免备份队列的复制以获得性能提升,同时仍然允许使用无锁算法自动更新队列的内容,从而避免中途丢失一些条目?

【问题讨论】:

    标签: java multithreading lock-free


    【解决方案1】:

    首先,让我们明确一点——最好的解决方案是避免编写任何此类自定义类。也许像java.util.concurrent.LinkedTransferQueue 这样简单的东西也能正常工作,而且不容易出错。如果LinkedTransferQueue 不起作用,那么LMAX disruptor 或类似的东西呢?您是否查看过现有的解决方案?

    如果您仍然需要/想要自定义解决方案,那么我有一个稍微不同的方法的草图,可以避免复制:

    这个想法是让put 操作围绕一些原子变量旋转,试图设置它。如果一个线程设法设置它,那么它将获得对当前队列的独占访问权,这意味着它可以附加到它。追加后,它重置原子变量以允许其他线程追加。它基本上是spin-lock。这样,线程之间的争用发生在追加到队列之前,而不是之后。

    【讨论】:

      【解决方案2】:

      我需要在每次写入尝试时复制实际队列

      您的想法听起来像 RCU (https://en.wikipedia.org/wiki/Read-copy-update)。通过为您解决释放问题(我认为),Java 被垃圾收集使 RCU 变得更加容易。

      如果我从您的问题的快速浏览中正确理解,您的“读者”实际上想为自己“声明”容器的全部当前内容。这也使他们成为有效的编写者,但是他们可以构造一个空容器,并原子地交换顶级引用来指向它,而不是读取+复制。 (因此要求旧容器独占访问。)

      RCU 的一大好处是容器数据结构本身不必到处都是原子的;一旦你引用了它,其他人就不会修改它。


      当作者想要向非空容器中添加新内容时,唯一棘手的部分就出现了。然后复制现有容器并修改副本,并尝试将更新后的副本通过 CAS(比较交换,即compareAndSet())转换为共享顶级AtomicReference

      作家不能只是无条件地交换,因为它最终可能会得到一个非空容器并且无处可放。除非作者可以坚持完成一批工作并旋转等待读者清空队列...


      我在这里假设您的作者有一批工作要立即排队;否则 RCU 对作者来说可能太贵了。 抱歉,如果我错过了您问题中排除这一点的细节。我不经常使用Java,所以我只是快速写下这篇文章以防万一。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-09-26
        • 2012-12-03
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-11-19
        相关资源
        最近更新 更多