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