【发布时间】: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