8.1.2 Project Reactor:使用 Mono 和 Flux 的响应式流

响应式 Spring 基于 Project Reactor,这是一个在 JVM 上构建异步、非阻塞应用程序的框架。Reactor 是 Reactive Streams 规范的实现,它旨在提供"具有非阻塞背压的异步流处理标准"(www.reactive-streams.org)。

从概念上讲,响应式流类似于 Java Stream API,我们使用它们来构建数据管道。关键区别之一是 Java 流是基于拉取的:消费者以命令式和同步方式处理数据。相反,响应式流是基于推送的:当新数据可用时,生产者通知消费者,因此处理是异步发生的。

响应式流根据生产者/消费者范式工作。生产者被称为发布者。它们生产可能最终可用的数据。Reactor 提供了两个实现 Producer 接口的核心 API,用于类型为 的对象,它们用于组合异步、可观察的数据流:Mono 和 Flux

  • Mono——表示单个异步值或空结果(0..1)
  • Flux——表示零个或多个项目的异步序列(0..N)

在 Java 流中,你将处理像 Optional 或 Collection 这样的对象。在响应式流中,你将有 Mono 或 Flux。响应式流的可能结果是空结果、值或错误。所有这些都被视为数据。当发布者返回所有数据时,我们说响应式流已成功完成。

消费者被称为订阅者,因为它们订阅发布者,并在新数据可用时收到通知。作为订阅的一部分,消费者还可以通过通知发布者它们一次只能处理一定数量的数据来定义背压。这是一个强大的功能,使消费者能够控制接收多少数据,防止它们被淹没并变得无响应。只有在有订阅者时,响应式流才会被激活。

你可以构建响应式流,组合来自不同来源的数据,并使用 Reactor 丰富的运算符集对其进行操作。在 Java 流中,你可以使用流畅的 API 通过 map、flatMap 或 filter 等运算符处理数据,每个运算符都构建一个新的 Stream 对象,保持前一步不可变。类似地,你可以使用流畅的 API 和运算符构建响应式流来处理异步接收的数据。

除了 Java 流可用的标准运算符之外,你可以使用更强大的运算符来应用背压、处理错误并提高应用程序弹性。例如,你将看到如何使用 retryWhen() 和 timeout() 运算符使 Order Service 和 Catalog Service 之间的交互更加健壮。运算符可以对发布者执行操作并返回新的发布者,而无需修改原始发布者,因此你可以轻松构建函数式和不可变的数据流。

Project Reactor 是 Spring 响应式技术栈的基础,它允许你使用 Mono 和 Flux 实现业务逻辑。在下一节中,你将了解更多关于使用 Spring 构建响应式应用程序的选项。

results matching ""

    No results matching ""