如何基于Java 9 Flow实现仅单个订阅者消费数据的发布方式?
Absolutely, you can build exactly this kind of work queue-style pub/sub model with Java 9's Flow API! The default behavior of SubmissionPublisher is to broadcast each message to all subscribers, but with a custom Processor or modified publisher implementation, you can enforce that each message is consumed by only one subscriber (the first available one that can handle it).
Core Concept
The key is to add a middle layer (a Flow.Processor) that acts as both a subscriber to your main publisher and a publisher to your worker subscribers. This layer will hold incoming messages in a queue, and when a subscriber signals it's ready to receive more data (via request(n)), the processor will dispatch the next available message to that subscriber—guaranteeing no message gets sent to more than one consumer.
Example Implementation
Here's a simplified, thread-safe version of an exclusive-dispatch processor that handles backpressure correctly:
import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.Flow.*; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; public class ExclusiveWorkQueueProcessor implements Processor<String, String> { private final ConcurrentLinkedQueue<String> messageQueue = new ConcurrentLinkedQueue<>(); private final ConcurrentLinkedQueue<SubscriptionState> subscribers = new ConcurrentLinkedQueue<>(); private Subscription upstreamSubscription; private final AtomicBoolean isDispatching = new AtomicBoolean(false); @Override public void subscribe(Subscriber<? super String> subscriber) { SubscriptionState state = new SubscriptionState(subscriber); subscribers.add(state); subscriber.onSubscribe(state); tryDispatch(); } @Override public void onSubscribe(Subscription subscription) { this.upstreamSubscription = subscription; subscription.request(1); // Request first message from upstream publisher } @Override public void onNext(String item) { messageQueue.add(item); tryDispatch(); upstreamSubscription.request(1); // Request next message (adjust based on your backpressure needs) } @Override public void onError(Throwable throwable) { subscribers.forEach(state -> state.subscriber.onError(throwable)); subscribers.clear(); } @Override public void onComplete() { subscribers.forEach(state -> state.subscriber.onComplete()); subscribers.clear(); } private void tryDispatch() { if (!isDispatching.compareAndSet(false, true)) { return; // Skip if we're already in the middle of dispatching } try { while (!messageQueue.isEmpty() && !subscribers.isEmpty()) { SubscriptionState currentState = subscribers.peek(); if (currentState == null) break; if (currentState.availablePermits.get() > 0) { String message = messageQueue.poll(); if (message == null) break; currentState.subscriber.onNext(message); currentState.availablePermits.decrementAndGet(); // Rotate subscribers for fairness (remove this if you want strict first-come-first-served) subscribers.poll(); subscribers.add(currentState); } else { // Move to next subscriber if this one has no pending requests subscribers.poll(); subscribers.add(currentState); } } } finally { isDispatching.set(false); } } private class SubscriptionState implements Subscription { private final Subscriber<? super String> subscriber; private final AtomicInteger availablePermits = new AtomicInteger(0); public SubscriptionState(Subscriber<? super String> subscriber) { this.subscriber = subscriber; } @Override public void request(long n) { if (n <= 0) { subscriber.onError(new IllegalArgumentException("Request must be a positive number")); return; } availablePermits.addAndGet((int) n); tryDispatch(); } @Override public void cancel() { subscribers.remove(this); } } }
How to Wire It Up
You can pair this processor with a SubmissionPublisher (your main task publisher) and multiple worker subscribers like this:
public class WorkQueueDemo { public static void main(String[] args) throws InterruptedException { SubmissionPublisher<String> taskPublisher = new SubmissionPublisher<>(); ExclusiveWorkQueueProcessor workQueue = new ExclusiveWorkQueueProcessor(); // Connect the main publisher to the work queue processor taskPublisher.subscribe(workQueue); // Spawn 3 worker subscribers for (int i = 1; i <= 3; i++) { int workerId = i; workQueue.subscribe(new Subscriber<>() { private Subscription subscription; @Override public void onSubscribe(Subscription subscription) { this.subscription = subscription; subscription.request(1); // Signal we're ready for our first task } @Override public void onNext(String task) { System.out.printf("Worker %d processed task: %s%n", workerId, task); // Request next task after finishing the current one subscription.request(1); } @Override public void onError(Throwable throwable) { throwable.printStackTrace(); } @Override public void onComplete() { System.out.printf("Worker %d has finished all tasks%n", workerId); } }); } // Publish 10 sample tasks for (int i = 1; i <= 10; i++) { taskPublisher.submit("Task #" + i); Thread.sleep(100); // Simulate delay between task submissions } taskPublisher.close(); workQueue.onComplete(); Thread.sleep(1000); // Wait for all workers to finish processing } }
Key Details to Note
- Backpressure Compliance: The
SubscriptionStatetracks how many tasks each worker has requested, ensuring we never send more tasks than a worker can handle. - Fairness: The example rotates subscribers after each dispatch to prevent any single worker from hogging all tasks—you can remove this rotation logic if you strictly want the first available worker to pick up every task.
- Thread Safety: We use concurrent queues and atomic variables to safely handle multi-threaded access, which is critical for Flow components designed to run in parallel environments.
This setup gives you exactly the behavior you're after: a single publisher feeding tasks into a work queue, where each task is processed by exactly one subscriber.
内容的提问来源于stack exchange,提问作者Roy Kornviets

