友情支持

如果您觉得这个笔记对您有所帮助,看在D瓜哥码这么多字的辛苦上,请友情支持一下,D瓜哥感激不尽,😜

支付宝

微信

有些打赏的朋友希望可以加个好友,欢迎关注D 瓜哥的微信公众号,这样就可以通过公众号的回复直接给我发信息。

wx jikerizhi

公众号的微信号是: jikerizhi因为众所周知的原因,有时图片加载不出来。 如果图片加载不出来可以直接通过搜索微信号来查找我的公众号。

70. Flow

早在 2013年,一些知名的有影响力的网络公司提出 Reactive Streams 提案,旨在标准版软件组件之间的异步数据交换。

为了减少重复和不兼容性,Java 9 引入了 java.util.concurrent.Flow 类,统一并规范了 Reactive Streams 的接口。但是, Flow 只定义了接口,并没有给出具体实现。下面是一个简单实现:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
package com.diguage.truman.concurrent;

import org.junit.jupiter.api.Test;

import java.util.concurrent.SubmissionPublisher;
import java.util.concurrent.TimeUnit;

import static java.util.concurrent.Flow.Subscriber;
import static java.util.concurrent.Flow.Subscription;

public class FlowTest {

  @Test
  public void test() throws InterruptedException {
    SubmissionPublisher<Integer> publisher = new SubmissionPublisher<>();
    publisher.subscribe(new PrintSubscriber());

    System.out.println("Submitting items...");
    for (int i = 0; i < 10; i++) {
      publisher.submit(i);
    }

    TimeUnit.SECONDS.sleep(1);
    publisher.close();
  }

  public static class PrintSubscriber implements Subscriber<Integer> {
    private Subscription subscription;

    @Override
    public void onSubscribe(Subscription subscription) {
      this.subscription = subscription;
      subscription.request(1);
    }

    @Override
    public void onNext(Integer item) {
      System.out.println("Received item: " + item);
      subscription.request(1);
    }

    @Override
    public void onError(Throwable throwable) {
      System.out.println("Error occurred: " + throwable.getMessage());
    }

    @Override
    public void onComplete() {
      System.out.println("PrintSubscriber is complete");
    }
  }
}

70.1. 参考资料