import java.util.concurrent.Flow;
import java.util.concurrent.Flow.Publisher;
import java.util.concurrent.Flow.Subscriber;
public class ReactiveExample {
public static void main(String[] args) {
// 创建一个发布者,发布一系列的数字
Publisher<Integer> publisher = new Publisher<>() {
@Override
public void subscribe(Flow.Subscriber<? super Integer> subscriber) {
subscriber.onSubscribe(new Flow.Subscription() {
private final int[] numbers = {1, 2, 3, 4, 5};
private int index = 0;
@Override
public void request(long n) {
for (int i = 0; i < n && index < numbers.length; i++) {
subscriber.onNext(numbers[index++]);
}
if (index == numbers.length) {
subscriber.onComplete();
}
}
@Override
public void cancel() {
// 取消订阅
}
});
}
};
// 创建一个订阅者,订阅发布者的数据
Subscriber<Integer> subscriber = new Subscriber<>() {
@Override
public void onSubscribe(Flow.Subscription subscription) {
System.out.println("Subscribed to publisher");
subscription.request(Long.MAX_VALUE); // 请求尽可能多的数据项
}
@Override
public void onNext(Integer item) {
System.out.println("Received: " + item);
}
@Override
public void one rror(Throwable throwable) {
throwable.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("Completed");
}
};
// 订阅发布者
publisher.subscribe(subscriber);
}
}
标签:Subscriber,Programming,void,Flow,subscriber,Reactive,Override,public
From: https://www.cnblogs.com/jiayuan2006/p/18294758