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);