import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.Flow.*; public class ConduitExample { public static void main(String[] args) throws InterruptedException { ExecutorService executorService = Executors.newFixedThreadPool(2); SubmissionPublisher<String> publisher = new SubmissionPublisher<>(executorService, 2); MySubscriber<Integer> subscriber = new MySubscriber<>(); publisher.subscribe(subscriber); publisher.submit("Hello"); publisher.submit("World"); publisher.submit("!"); executorService.awaitTermination(1, TimeUnit.SECONDS); publisher.close(); executorService.shutdown(); } static class MySubscriber<T> implements Subscriber<T> { private Subscription subscription; @Override public void onSubscribe(Subscription subscription) { this.subscription = subscription; subscription.request(1); } @Override public void onNext(T item) { System.out.println("Received: " + item); subscription.request(1); } @Override public void onError(Throwable throwable) { throwable.printStackTrace(); } @Override public void onComplete() { System.out.println("Done"); } } }


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