public class LowLatencyQueue { private final RingBuffer<Event> ringBuffer; public LowLatencyQueue(int bufferSize) { EventFactory<Event> eventFactory = Event::new; Executor executor = Executors.newCachedThreadPool(); WaitStrategy waitStrategy = new BusySpinWaitStrategy(); ringBuffer = RingBuffer.createSingleProducer(eventFactory, bufferSize, waitStrategy, executor); } public void produce(Event event) { long sequence = ringBuffer.next(); try { Event newEvent = ringBuffer.get(sequence); newEvent.set(event); } finally { ringBuffer.publish(sequence); } } public void consume(EventHandler<Event> eventHandler) { BatchEventProcessor<Event> eventProcessor = new BatchEventProcessor<>(ringBuffer, ringBuffer.newBarrier(), eventHandler); ringBuffer.addGatingSequences(eventProcessor.getSequence()); } } public class LowLatencyQueue { private final ConcurrentLinkedQueue<Event> queue; public LowLatencyQueue() { queue = new ConcurrentLinkedQueue<>(); } public void produce(Event event) { queue.add(event); } public void consume(Consumer<Event> consumer) { while (!queue.isEmpty()) { Event event = queue.poll(); consumer.accept(event); } } } int bufferSize = 1024; LowLatencyQueue queue = new LowLatencyQueue(bufferSize); EventHandler<Event> eventHandler = (event, sequence, endOfBatch) -> { }; queue.consume(eventHandler); LowLatencyQueue queue = new LowLatencyQueue(); Consumer<Event> consumer = event -> { }; queue.consume(consumer);


上一篇:
下一篇:
切换中文