Reactor Sinks

🧩 一、什么是 Sinks(汇流器)

在 Reactor 中,Sinks 是一种用于向 Flux 或 Mono 发送数据的机制。它们允许你以非阻塞、响应式的方式将数据推送到响应式流中。

✅ 核心概念:

  • SinksFlux 或 Mono 的“生产端” ,而 Subscribers 是“消费端”。
  • 你可以通过 Flux.sink()Mono.sink() 创建一个 Sink。
  • Sink 提供了多种方法来发送数据,例如:
    • next():发送一个值
    • error():发送一个错误
    • complete():完成流
    • cancel():取消订阅(如果尚未完成)

🧱 二、Sinks 的类型

1. FluxSink

  • 用于创建一个 Flux 的 Sink。
  • 提供了以下方法:
    • next(T value)
    • error(Throwable error)
    • complete()
    • cancel()
    • requestFusion(int mode):控制背压(backpressure)策略

2. MonoSink

  • 用于创建一个 Mono 的 Sink。
  • 提供了以下方法:
    • success(T value)
    • error(Throwable error)
    • complete()
    • cancel()

🔄 三、Sinks 的使用场景

1. 手动控制数据流

  • 你可以通过 Sinks 手动控制何时发送数据、何时出错、何时完成。
  • 适用于需要异步处理外部事件驱动的场景。

2. 与外部系统集成

  • 例如,将外部 API 调用、数据库查询、消息队列等集成到响应式流中。
  • 你可以通过 Sinks 将这些外部操作的结果推送到 Flux 或 Mono 中。

3. 背压控制

  • Sinks 提供了 requestFusion() 方法,允许你控制背压策略。
  • 你可以选择:
    • EAGER:立即请求所有数据
    • ASYNC:异步请求
    • BOUNDED:有界请求

🧪 四、示例代码

示例 1:使用 Flux.sink()

Flux<String> flux = Flux.sink();
FluxSink<String> sink = flux.sink();sink.next("Hello");
sink.next("World");
sink.complete();

示例 2:使用 Mono.sink()

Mono<String> mono = Mono.sink();
MonoSink<String> sink = mono.sink();sink.success("Hello");
sink.complete();

示例 3:使用 Sinks 处理外部事件

Flux<String> flux = Flux.sink();
FluxSink<String> sink = flux.sink();// 模拟外部事件
externalEventSource.subscribe(value -> sink.next(value),error -> sink.error(error),() -> sink.complete()
);

⚠️ 五、注意事项

  1. 不要重复使用 Sinks

    • 每个 Sink 只能被订阅一次。如果多次订阅,可能会导致错误或不可预测的行为。
  2. 注意背压

    • 如果你没有正确处理背压,可能会导致内存溢出或性能问题。
    • 使用 requestFusion() 来控制背压策略。
  3. Sinks 是“生产端”

    • 它们不是用于消费数据的,而是用于向响应式流中“推”数据的。

✅ 六、总结

项目描述
SinksReactor 中用于向 Flux 或 Mono 发送数据的机制
FluxSink用于创建 Flux 的 Sink,支持发送多个值、错误、完成
MonoSink用于创建 Mono 的 Sink,支持发送一个值、错误、完成
使用场景手动控制数据流、集成外部系统、处理背压
注意事项不要重复使用、注意背压、Sinks 是“生产端”

在 Reactor 中,Sinks 是一种用于手动触发信号(如 onNextonErroronComplete)的机制,它允许你以独立的方式创建一个类似于 Publisher 的结构,从而可以处理多个订阅者(Subscriber)。Sinks 提供了多种类型,每种类型适用于不同的使用场景。以下是对你提到的 Sinks 类型的详细解释:


1. many().multicast()

  • 描述:这是一个多播(multicast)的 Sinks,它会将新推送的数据传输给其订阅者,并且尊重每个订阅者的背压(backpressure)。
  • 特点
    • 多播:一个 Sinks 可以有多个订阅者。
    • 背压处理:每个订阅者在订阅后只接收在该订阅者订阅之后推送的数据。
    • 适用场景:当你需要将数据分发给多个订阅者,并且每个订阅者只关心自己订阅后的新数据时,可以使用这种类型。

2. many().unicast()

  • 描述:这是一个单播(unicast)的 Sinks,它与 multicast() 类似,但有一个额外的特性:在第一个订阅者注册之前推送的数据会被缓冲。
  • 特点
    • 单播:一个 Sinks 只能有一个订阅者。
    • 缓冲机制:在第一个订阅者订阅之前,所有推送的数据都会被缓冲,直到订阅者到来。
    • 适用场景:当你需要将数据发送给一个订阅者,并且希望在订阅者到来之前保留所有数据时,可以使用这种类型。

3. many().replay()

  • 描述:这是一个重播(replay)的 Sinks,它会将指定历史大小的推送数据重播给新的订阅者,然后继续推送新数据。
  • 特点
    • 重播机制:可以为新的订阅者重播之前推送的数据。
    • 配置方式
  • limit(int size):指定重播的历史大小。
  • limit(Duration duration):基于时间窗口重播。
  • limit(int size, Duration duration):结合大小和时间窗口重播。
  • latest():只重播最后一个元素。
    • 适用场景:当你需要将历史数据发送给后来加入的订阅者时,可以使用这种类型。

4. one()

  • 描述:这是一个单元素的 Sinks,它只支持推送一个元素。
  • 特点
    • 单元素:只能推送一个元素。
    • 等效于 Mono:可以通过 asMono() 方法将其视为一个 Mono
    • 适用场景:当你只需要推送一个元素,并且不需要处理多个元素的场景时,可以使用这种类型。

5. empty()

  • 描述:这是一个空的 Sinks,它只推送一个终端信号(error 或 complete)。
  • 特点
    • 终端信号:只能推送一个错误或完成信号。
    • 等效于 Mono:可以通过 asMono() 方法将其视为一个 Mono
    • 适用场景:当你只需要推送一个终端信号(如错误或完成),并且不需要推送任何数据时,可以使用这种类型。

总结

Sinks 类型描述特点
many().multicast()多播 Sinks,将新推送的数据传输给订阅者多播、背压处理
many().unicast()单播 Sinks,缓冲在订阅者之前推送的数据单播、缓冲机制
many().replay()重播 Sinks,重播历史数据重播机制、可配置
one()单元素 Sinks,只推送一个元素单元素、等效于 Mono
empty()空 Sinks,只推送终端信号终端信号、等效于 Mono

通过选择合适的 Sinks 类型,你可以更灵活地控制数据的推送和处理,从而构建出更健壮的响应式系统。


📚 参考文档

  • Project Reactor Sinks 文档

如果你需要我帮你写一个完整的 Sinks 示例,或者想了解如何在 Spring WebFlux 中使用 Sinks,也可以告诉我!