10.5 使用 Spring Cloud Stream 产生和消费消息

在前面的部分中,您学习了函数式编程范式以及它如何适应 Spring 生态系统,使用 Spring Cloud Function 和 Spring Cloud Stream。最后一节将指导您完成生产者和消费者的实现。

正如您将看到的,消费者与您在 Dispatcher Service 中编写的函数没有太大区别。另一方面,生产者略有不同,因为与函数和消费者不同,它们不是自然激活的。我将向您展示如何在 Order Service 中同时使用它们,同时实现 Polar Bookshop 系统订单流程的最后一部分。

10.5.1 实现事件消费者,以及幂等性问题

我们之前构建的 Dispatcher Service 应用程序在订单被调度时生成消息。Order Service 应在发生这种情况时被通知,以便它可以在数据库中更新订单状态。

首先,打开您的 Order Service 项目(order-service),并在 build.gradle 文件中添加对 Spring Cloud Stream 和测试绑定器的依赖。添加新依赖后,请记住刷新或重新导入 Gradle 依赖。

清单 10.17 添加 Spring Cloud Stream 和测试绑定器的依赖

dependencies {
 ...
 implementation 'org.springframework.cloud:spring-cloud-stream-binder-rabbit'
 testImplementation("org.springframework.cloud:spring-cloud-stream") {
 artifact {
 name = "spring-cloud-stream"
 extension = "jar"
 type = "test-jar"
 classifier = "test-binder"
 }
 }
}

接下来,我们需要建模 Order Service 要监听的事件。创建一个新的 com.polarbookshop.orderservice.order.event 包,并添加一个 OrderDispatchedMessage 类来保存已调度订单的标识符。

清单 10.18 表示订单被调度事件的 DTO

package com.polarbookshop.orderservice.order.event;

public record OrderDispatchedMessage (
 Long orderId
){}

现在,我们将使用函数式方法实现业务逻辑。创建一个 OrderFunctions 类(com.polarbookshop.orderservice.order.event 包),并实现一个函数来消费 Dispatcher Service 应用程序在订单被调度时生成的消息。该函数将是一个 Consumer,负责监听传入的消息并相应地更新数据库实体。Consumer 对象是有输入但没有输出的函数。为了保持函数简洁易读,我们将 OrderDispatchedMessage 对象的处理转移到 OrderService 类(我们将在一分钟后实现)。

清单 10.19 从 RabbitMQ 消费消息

package com.polarbookshop.orderservice.order.event;

import java.util.function.Consumer;
import com.polarbookshop.orderservice.order.domain.OrderService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.Flux;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class OrderFunctions {

 private static final Logger log = LoggerFactory.getLogger(OrderFunctions.class);

 @Bean
 public Consumer<Flux<OrderDispatchedMessage>> dispatchOrder(
 OrderService orderService) {
 return flux -> orderService.consumeOrderDispatchedEvent(flux)
 // 对于数据库中更新的每个订单,它记录一条消息
 .doOnNext(order -> log.info("The order with id {} is dispatched", order.id()))
 .subscribe();
 // 订阅响应式流以激活它。没有订阅者,就没有数据流过流
 }
}

Order Service 是一个响应式应用程序,因此 dispatchOrder 函数将作为响应式流(OrderDispatchedMessage 的 Flux)消费消息。响应式流仅在有订阅者对接收数据感兴趣时才会激活。因此,我们必须通过订阅来结束响应式流,否则将永远不会处理任何数据。在前面的示例中,订阅部分由框架透明地处理(例如,当使用响应式流通过 REST 端点返回数据或向备份服务发送数据时)。在这种情况下,我们必须使用 subscribe() 子句显式执行此操作。

接下来,让我们在 OrderService 类中实现 consumeOrderDispatchedMessageEvent() 方法,以便在订单被调度后更新数据库中的状态。

清单 10.20 实现将订单更新为已调度的逻辑

@Service
public class OrderService {
 // ...
 public Flux<Order> consumeOrderDispatchedEvent(
 // 接受 OrderDispatchedMessage 对象的响应式流作为输入
 Flux<OrderDispatchedMessage> flux) {
 return flux
 .flatMap(message -> orderRepository.findById(message.orderId()))
 // 对于流中发出的每个对象,它从数据库中读取相关订单
 .map(this::buildDispatchedOrder)
 // 使用 "dispatched" 状态更新订单
 .flatMap(orderRepository::save);
 // 将更新后的订单保存在数据库中
 }

 private Order buildDispatchedOrder(Order existingOrder) {
 return new Order(
 existingOrder.id(),
 existingOrder.bookIsbn(),
 existingOrder.bookName(),
 existingOrder.bookPrice(),
 existingOrder.quantity(),
 OrderStatus.DISPATCHED,
 // 给定一个订单,它返回一个具有 "dispatched" 状态的新记录
 existingOrder.createdDate(),
 existingOrder.lastModifiedDate(),
 existingOrder.version()
 );
 }
}

消费者由到达队列的消息触发。RabbitMQ 提供至少一次投递保证,因此您需要注意可能的重复项。我们实现的代码将特定订单的状态更新为 DISPATCHED,这是一个可以多次执行并产生相同结果的操作。由于该操作是幂等的,因此代码对重复项具有弹性。进一步的优化将是检查状态,如果订单已经调度,则跳过更新操作。

最后,我们需要在 application.yml 文件中配置 Spring Cloud Stream,以便 dispatchOrder-in-0 绑定(从 dispatchOrder 函数名推断)映射到 RabbitMQ 中的 order-dispatched 交换机。另外,请记住将 dispatchOrder 定义为 Spring Cloud Function 应该管理的函数,以及与 RabbitMQ 的集成。

清单 10.21 配置 Cloud Stream 绑定和 RabbitMQ 集成

spring:
 cloud:
 function:
 definition: dispatchOrder
 # Spring Cloud Function 管理的函数定义
 stream:
 bindings:
 dispatchOrder-in-0:
 # 输入绑定
 destination: order-dispatched
 # 绑定器绑定到的代理上的实际名称(RabbitMQ 中的交换机)
 group: ${spring.application.name}
 # 对目标感兴趣的消费者组(与应用程序名称相同)
 rabbitmq:
 host: localhost
 port: 5672
 username: user
 # 配置与 RabbitMQ 的集成
 password: password
 connection-timeout: 5s

如您所见,它的工作方式与 Dispatcher Service 中的函数相同。Order Service 中的消费者将属于 order-service 消费者组,Spring Cloud Stream 将在它们和 RabbitMQ 中的 order-dispatched.order-service 队列之间定义一个消息通道。

接下来,我们将通过定义一个负责触发整个过程的供应商来完成订单流程。

10.5.2 实现事件生产者,以及原子性问题

供应商是消息源。它们在事件发生时生成消息。在 Order Service 中,供应商应在订单被接受时通知相关方(在本例中为 Dispatcher Service)。与函数和消费者不同,供应商需要被激活。它们仅在被调用时才起作用。

Spring Cloud Stream 提供了几种定义供应商的方法,并覆盖不同的场景。在我们的例子中,事件源不是消息代理,而是 REST 端点。当用户向 Order Service 发送 POST 请求以购买书籍时,我们希望发布一个信号,表明订单是否已被接受。

让我们首先将该事件建模为 DTO。它将与我们在 Dispatcher Service 中使用的 OrderAcceptedMessage record 相同。在 Order Service 项目的 com.polarbookshop.orderservice.order.event 包中添加该 record。

清单 10.22 表示订单被接受事件的 DTO

package com.polarbookshop.orderservice.order.event;

public record OrderAcceptedMessage (
 Long orderId
){}

我们可以使用 StreamBridge 对象将 REST 层与应用程序的流部分桥接起来,该对象允许我们以命令方式将数据发送到特定目标。让我们分解一下这个新功能。首先,我们可以实现一个方法,该方法接受 Order 对象作为输入,验证它是否被接受,构建 OrderAcceptedMessage 对象,并使用 StreamBridge 将其发送到 RabbitMQ 目标。

打开 OrderService 类,自动装配 StreamBridge 对象,并定义新的 publishOrderAcceptedEvent 方法。

清单 10.23 实现向目标发布事件的逻辑

package com.polarbookshop.orderservice.order.domain;

import com.polarbookshop.orderservice.book.BookClient;
import com.polarbookshop.orderservice.order.event.OrderAcceptedMessage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.stereotype.Service;
// ...

@Service
public class OrderService {
 private static final Logger log = LoggerFactory.getLogger(OrderService.class);

 private final BookClient bookClient;
 private final OrderRepository orderRepository;
 private final StreamBridge streamBridge;

 public OrderService(BookClient bookClient, StreamBridge streamBridge, 
 OrderRepository orderRepository) {
 this.bookClient = bookClient;
 this.orderRepository = orderRepository;
 this.streamBridge = streamBridge;
 }

 // ...

 private void publishOrderAcceptedEvent(Order order) {
 if (!order.status().equals(OrderStatus.ACCEPTED)) {
 return;
 // 如果订单未被接受,它不执行任何操作
 }

 var orderAcceptedMessage = new OrderAcceptedMessage(order.id());
 // 构建消息以通知订单已被接受

 log.info("Sending order accepted event with id: {}", order.id());

 var result = streamBridge.send("acceptOrder-out-0", orderAcceptedMessage);
 // 显式地向 acceptOrder-out-0 绑定发送消息

 log.info("Result of sending data for order with id {}: {}", order.id(), result);
 }
}

由于数据源是 REST 端点,因此没有我们可以注册到 Spring Cloud Function 的 Supplier Bean,因此没有触发器让框架创建与 RabbitMQ 的必要绑定。然而,在清单 10.23 中,StreamBridge 用于向 acceptOrder-out-0 绑定发送数据。它从哪里来?没有 acceptOrder 函数!

在启动时,Spring Cloud Stream 将注意到 StreamBridge 想要通过 acceptOrder-out-0 绑定发布消息,并将自动创建一个。类似于从函数创建的绑定,我们可以配置 RabbitMQ 中的目标名称。打开 application.yml 文件并按如下所示配置绑定。

清单 10.24 配置 Cloud Stream 输出绑定

spring:
 cloud:
 function:
 definition: dispatchOrder
 stream:
 bindings:
 dispatchOrder-in-0:
 destination: order-dispatched
 group: ${spring.application.name}
 acceptOrder-out-0:
 # 由 StreamBridge 创建和管理的输出绑定
 destination: order-accepted
 # 绑定器绑定到的代理上的实际名称(RabbitMQ 中的交换机)

现在剩下的就是在提交的订单被接受时调用该方法。这是一个关键点,也是 saga 模式的特征之一,saga 模式是微服务架构中分布式事务的流行替代方案。为了确保系统中的一致性,将订单持久化到数据库和发送有关它的消息必须原子地完成。要么两个操作都成功,要么两个操作都必须失败。确保原子性的一种简单而有效的方法是将两个操作包装在本地事务中。为此,我们可以依赖内置的 Spring 事务管理功能。

注意 saga 模式在 Chris Richardson 的《Microservices Patterns》(Manning,2018;https://livebook.manning.com/book/microservices-patterns/chapter-4)第 4 章中有广泛描述。如果您对设计跨多个应用程序的业务事务感兴趣,我建议您查看它。

在 OrderService 类中,修改 submitOrder() 方法以调用 publishOrderAcceptedEvent 方法,并使用 @Transactional 注解它。

清单 10.25 使用数据库和事件代理定义 saga 事务

@Service
public class OrderService {
 // ...

 @Transactional
 // 在本地事务中执行该方法
 public Mono<Order> submitOrder(String isbn, int quantity) {
 return bookClient.getBookByIsbn(isbn)
 .map(book -> buildAcceptedOrder(book, quantity))
 .defaultIfEmpty(buildRejectedOrder(isbn, quantity))
 .flatMap(orderRepository::save)
 // 将订单保存在数据库中
 .doOnNext(this::publishOrderAcceptedEvent);
 // 如果订单被接受,则发布事件
 }

 private void publishOrderAcceptedEvent(Order order) {
 if (!order.status().equals(OrderStatus.ACCEPTED)) {
 return;
 }
 var orderAcceptedMessage = new OrderAcceptedMessage(order.id());
 log.info("Sending order accepted event with id: {}", order.id());
 var result = streamBridge.send("acceptOrder-out-0", orderAcceptedMessage);
 log.info("Result of sending data for order with id {}: {}", order.id(), result);
 }
}

Spring Boot 预配置了事务管理功能,可以处理涉及关系数据库的事务操作(正如您在第 5 章中学到的那样)。但是,为消息生产者建立的与 RabbitMQ 的通道默认情况下不是事务性的。要使事件发布操作加入现有事务,我们需要在 application.yml 文件中启用 RabbitMQ 对消息生产者的事务支持。

清单 10.26 配置输出绑定为事务性

spring:
 cloud:
 function:
 definition: dispatchOrder
 stream:
 bindings:
 dispatchOrder-in-0:
 destination: order-dispatched
 group: ${spring.application.name}
 acceptOrder-out-0:
 destination: order-accepted
 rabbit:
 # Spring Cloud Stream 绑定的 RabbitMQ 特定配置
 bindings:
 acceptOrder-out-0:
 producer:
 transacted: true
 # 使 acceptOrder-out-0 绑定具有事务性

现在您可以为供应商和消费者编写新的集成测试,就像我们为 Dispatcher Service 中的函数所做的那样。我将自动测试留给您,因为您现在拥有所需的工具。如果您需要灵感,请查看本书附带的源代码(Chapter10/10-end/order-service)。

您还需要在现有的 OrderServiceApplicationTests 类中导入测试绑定器的配置(@Import(TestChannelBinderConfiguration.class)),以使其正常工作。

我们已经很好地完成了事件驱动模型、函数和消息传递系统的旅程。在结束之前,让我们看看实际的订单流程。首先,启动 RabbitMQ、PostgreSQL(docker-compose up -d polar-rabbitmq polar-postgres)和 Dispatcher Service(./gradlew bootRun)。然后运行 Catalog Service 和 Order Service(./gradlew bootRun 或在构建镜像后从 Docker Compose 运行)。

一旦所有这些服务都启动并运行,向目录中添加一本新书:

$ http POST :9001/books author="Jon Snow" \
 title="All I don't know about the Arctic" \
 isbn="1234567897" \
 price=9.90 \
 publisher="Polarsophia"

然后订购三本该书:

$ http POST :9002/orders isbn=1234567897 quantity=3

如果您订购一本存在的书,订单将被接受,Order Service 将发布 OrderAcceptedEvent 消息。订阅同一事件的 Dispatcher Service 将处理该订单并发布 OrderDispatchedEvent 消息。Order Service 将收到通知并更新数据库中的订单状态。

提示 您可以通过检查 Order Service 和 Dispatcher Service 的应用程序日志来跟踪消息流。

现在是关键时刻。从 Order Service 获取订单:

$ http :9002/orders

状态应为 DISPATCHED

{
 "bookIsbn": "1234567897",
 "bookName": "All I don't know about the Arctic - Jon Snow",
 "bookPrice": 9.9,
 "createdDate": "2022-06-06T19:40:33.426610Z",
 "id": 1,
 "lastModifiedDate": "2022-06-06T19:40:33.866588Z",
 "quantity": 3,
 "status": "DISPATCHED",
 "version": 2
}

确实如此。干得好!测试完系统后,停止所有应用程序(Ctrl-C)和 Docker 容器(docker-compose down)。

这就完成了 Polar Bookshop 系统业务逻辑的主要实现。下一章将介绍使用 Spring Security、OAuth 2.1 和 OpenID Connect 保护云原生应用程序的安全。


Polar Labs

请随时应用您在前面章节中学到的知识,并为 Dispatcher Service 应用程序的部署做好准备。

  1. 向 Dispatcher Service 添加 Spring Cloud Config Client,使其从 Config Service 获取配置数据。
  2. 配置 Cloud Native Buildpacks 集成,容器化应用程序,并定义部署管道的提交阶段。
  3. 编写 Deployment 和 Service 清单,以便将 Dispatcher Service 部署到 Kubernetes 集群。
  4. 配置 Tilt 以将 Dispatcher Service 的部署自动化到使用 minikube 初始化的本地 Kubernetes 集群。

然后更新 Docker Compose 规范和 Kubernetes 清单,以为 Order Service 配置 RabbitMQ 集成。

您可以参考本书附带源代码中的 Chapter10/10-end 文件夹以查看最终结果(https://github.com/ThomasVitale/cloud-native-spring-in-action)。从 Chapter10/10-end/polar-deployment/kubernetes/platform/development 文件夹中可用的清单部署备份服务,使用 kubectl apply -f services


本章小结

  • 事件驱动架构是通过产生和消费事件进行交互的分布式系统。
  • 事件是系统中发生的相关事件。
  • 在 pub/sub 模型中,生产者发布事件,这些事件被发送给所有订阅者进行消费。
  • RabbitMQ 和 Kafka 等事件处理平台负责从生产者收集事件,路由并将它们分发给感兴趣的消费者。
  • 在 AMQP 协议中,生产者将消息发送到代理中的交换机,交换机根据特定的路由算法将它们转发到队列。
  • 在 AMQP 协议中,消费者从代理中的队列接收消息。
  • 在 AMQP 协议中,消息是由键/值属性和二进制有效负载组成的数据结构。
  • RabbitMQ 是一个基于 AMQP 协议的消息代理,您可以用它来实现基于 pub/sub 模型的事件驱动架构。
  • RabbitMQ 提供高可用性、弹性和数据复制。
  • Spring Cloud Function 使您能够使用标准 Java FunctionSupplierConsumer 接口实现业务逻辑。
  • Spring Cloud Function 包装您的函数并提供几个令人兴奋的功能,如透明的类型转换和函数组合。
  • 在 Spring Cloud Function 上下文中实现的函数可以以不同的方式公开和与外部系统集成。
  • 函数可以公开为 REST 端点,打包并部署到 FaaS 平台作为无服务器应用程序(Knative、AWS Lambda、Azure Function、Google Cloud Functions),或者可以绑定到消息通道。
  • Spring Cloud Stream 构建在 Spring Cloud Function 之上,为您提供了将函数与 RabbitMQ 或 Kafka 等外部消息系统集成所需的所有基础设施。
  • 一旦您实现函数,就不必对代码进行任何更改。您只需添加对 Spring Cloud Stream 的依赖并对其进行配置以适应您的需求。
  • 在 Spring Cloud Stream 中,目标绑定器提供与外部消息系统的集成。
  • 在 Spring Cloud Stream 中,目标绑定(输入和输出)将应用程序中的生产者和消费者与 RabbitMQ 等消息代理中的交换机和队列桥接起来。
  • 函数和消费者在新消息到达时自动激活。
  • 供应商需要显式激活,例如通过显式地向目标绑定发送消息。

results matching ""

    No results matching ""