尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Springy-Store-Microservices事件驱动架构揭秘:Spring Cloud Stream 一套代码玩转 RabbitMQ 与 Kafka 双中间件

Springy-Store-Microservices事件驱动架构揭秘:Spring Cloud Stream 一套代码玩转 RabbitMQ 与 Kafka 双中间件 Springy-Store-Microservices事件驱动架构揭秘Spring Cloud Stream 一套代码玩转 RabbitMQ 与 Kafka 双中间件【免费下载链接】Springy-Store-MicroservicesSpringy Store is a conceptual simple μServices-based project using the latest cutting-edge technologies, to demonstrate how the Store services are created to be a cloud-native and 12-factor app agnostic. Those μServices are developed based on Spring Boot Cloud framework that implements cloud-native intuitive, design patterns, and best practices.项目地址: https://gitcode.com/gh_mirrors/sp/Springy-Store-MicroservicesSpringy-Store-Microservices是一个基于 Spring Boot 与 Spring Cloud 构建的云原生微服务教学项目用「写操作走消息、读操作走 HTTP」的**事件驱动架构EDA**演示如何用Spring Cloud Stream统一抽象消息层让同一套服务代码在RabbitMQ和Kafka双消息中间件之间一键切换完整覆盖重试、死信、分区等生产级实战要点是学习微服务消息驱动开发的绝佳实战样本。1️⃣ 事件驱动架构总览读写分离的微服务全景Springy-Store 模拟了一个在线商店系统由 4 个核心业务微服务组成store-service门店聚合服务对外提供复合商品 APIproduct-service商品服务、review-service评论服务、recommendation-service推荐服务其事件驱动设计的核心思想非常清晰操作类型通信方式说明读操作HTTPWebClient组合查询商品评论推荐带熔断与重试写操作异步消息事件发布 CREATE / DELETE 事件到主题Store 服务 ──发布事件──▶ products 主题 ──▶ product-service ──发布事件──▶ recommendations 主题 ──▶ recommendation-service ──发布事件──▶ reviews 主题 ──▶ review-service客户端只需向 store-service 发起一次创建复合商品请求store-service 就拆解出三类事件异步分发各服务独立消费落库MongoDB / MySQL彼此完全解耦——这就是典型的事件驱动架构在 Spring Cloud 生态中的落地。上图即为项目自带的Springy Store 事件驱动微服务架构全景图左侧 Edge Server 网关统一入口Store 微服务通过 Circuit Breaker 与 Read API 完成同步读三条绿色 Publish 链路分别把写事件发布到 Reviews / Products / Recommendations 三个主题由对应微服务异步消费——读走 HTTP、写走事件的 EDA 双通道设计一目了然。2️⃣ 统一事件模型一个 Event 类贯穿所有消息所有服务共享同一个不可变的事件对象定义在store-common/store-api/模块的event/Event.java中结构极为简洁Event { Type eventType // CREATE 或 DELETE K key // 业务键如商品 ID T data // 事件负载 LocalDateTime eventCreatedAt }这个设计有两个实战价值自描述事件消费端只看eventType就知道该建还是该删天然支持分区key字段商品 ID正是后续 Kafka / RabbitMQ 分区路由的分区键partition key。3️⃣ Spring Cloud Stream抽象层让中间件可插拔项目的精妙之处在于Java 代码完全不感知底层是 RabbitMQ 还是 Kafka。发布端——store-service 在store-services/store-service/的integration/StoreIntegration.java中通过EnableBinding绑定三个输出通道Output(output-products) // → products 主题 Output(output-recommendations) // → recommendations 主题 Output(output-reviews) // → reviews 主题发送消息只需一行伪代码messageSources.outputProducts().send(withPayload(new Event(CREATE, id, product)))消费端——各服务使用 Spring 内置的Sink接口用StreamListener(target Sink.INPUT)监听输入通道。以 product-service 的infra/MessageProcessor.java为例StreamListener(target Sink.INPUT) public void process(EventInteger, Product event) { switch (event.getEventType()) { case CREATE - productService.createProduct(event.getData()); case DELETE - productService.deleteProduct(event.getKey()); default - throw new EventProcessingException(...); } }MessageChannel之上是 Spring Cloud Stream 的 Binder 抽象层——换中间件只改配置不改代码这正是Binder 可插拔设计的精髓。4️⃣ 一份配置切换 RabbitMQ 与 KafkadefaultBinder 与 kafka 配置文件Binder 的选择集中在配置服务器仓库config/repo/product.yml通过 config-server 统一下发配置项作用spring.cloud.stream.defaultBinder: rabbit默认使用RabbitMQ Binderdefault.contentType: application/json消息 JSON 序列化bindings.input.destination: products输入通道绑定products主题bindings.input.group: productsGroup消费者组群发模式consumer.maxAttempts: 3最多尝试 3 次含首次backOffInitialInterval: 500/backOffMultiplier: 2.0指数退避重试rabbit.bindings.input.consumer.autoBindDlq: trueRabbitMQ 自动绑定死信队列kafka.bindings.input.consumer.enableDlq: trueKafka 侧启用 DLQ注意这里同时预留了kafka绑定配置段——RabbitMQ 与 Kafka 两套可靠性配置写在同一份 yml 里切换 Binder 只需通过 Spring 配置文件激活kafkaprofile该 profile 定义在config/repo/application.yml中Binder 会自动换成 Kafka Binder。5️⃣ 三种运行形态docker-compose 一键切换消息中间件项目根目录提供三份 docker-compose 文件对应三种运行形态这是双中间件实战最直观的入口文件消息形态启动命令docker-compose.ymlRabbitMQ无分区docker-compose -p ssm up -ddocker-compose-partitions.ymlRabbitMQ 2 分区每服务双实例docker-compose -p ssm -f docker-compose-partitions.yml up -ddocker-compose-kafka.ymlKafka Zookeeper2 分区双实例docker-compose -p ssm -f docker-compose-kafka.yml up -d以 Kafka 形态为例compose 文件拉起wurstmeister/kafka与 Zookeeper 容器并通过SPRING_PROFILES_ACTIVEdocker,streaming_partitioned,streaming_instance_0,kafka激活分区与 Kafka 配置RabbitMQ 形态则使用rabbitmq:3-management镜像可在http://localhost:5672guest/guest管理界面直接查看主题、分区、DLQ 与消息负载。分区Partitioning如何生效激活streaming_partitionedprofile 后生产者按商品 ID 路由output-products.producer: partition-key-expression: payload.key # 用商品ID做分区键 partition-count: 2 # 两个分区配合streaming_instance_0 / streaming_instance_1两个 profile每个服务跑两个实例分别消费不同分区——同一个商品的所有事件始终落在同一分区、被同一实例处理保证顺序性同时双实例水平扩展消费能力。Kafka 依赖原生 partitionRabbitMQ 则借助instanceIndex映射到分区队列抽象层抹平了差异。6️⃣ 消息可靠性实战重试、指数退避与死信队列生产级事件驱动系统的三大可靠性武器本项目全部配置齐活重试 指数退避maxAttempts: 3首次失败 500ms 后重试退避乘数 2.0避免瞬时故障丢消息死信队列DLQRabbitMQ 侧autoBindDlq: truerepublishToDlq: true三次尝试仍失败的消息转投products-dlqKafka 侧enableDlq: true等效兜底测试验证store-services/store-service/下的MessagingTests.java使用 Spring Cloud Stream 的MessageCollector测试 Binder断言发布到各通道的Event内容与期望完全一致——中间件无关测试也无关。7️⃣ Zipkin 追踪事件链路看得见的事件流事件驱动系统的痛点是消息流转不可见。Springy-Store 引入Zipkin完成分布式追踪上图为 Zipkin 的 Dependencies 依赖关系图可以清晰看到 gateway → store → product/review/recommendation 的调用边以及各服务与 rabbitmq、MongoDB、MySQL 之间的消息与存储链路——异步事件链路上的每一跳都被串成了完整的 Trace让 EDA 系统具备了与同步调用同等的可观测性。8️⃣ 快速上手三步跑通双中间件第一步克隆并安装共享模块git clone https://gitcode.com/gh_mirrors/sp/Springy-Store-Microservices.git cd Springy-Store-Microservices ./setup.sh第二步构建并测试所有微服务./mvnw clean verify -Ddockerfile.skip第三步选择消息形态启动# 形态一RabbitMQ 基础版 docker-compose -p ssm up -d # 形态二RabbitMQ 分区双实例 docker-compose -p ssm -f docker-compose-partitions.yml up -d # 形态三Kafka Zookeeper 分区双实例 docker-compose -p ssm -f docker-compose-kafka.yml up -d启动后可通过 Swaggerhttps://localhost:8443/swagger-ui.html向 store-service 创建复合商品再登录 RabbitMQ 管理界面或借助 Zipkin观察事件在主题间的流转。9️⃣ 关键文件导航统一事件模型store-common/store-api/src/main/java/com/siriusxi/ms/store/api/event/Event.java事件发布端Store 聚合服务store-services/store-service/src/main/java/com/siriusxi/ms/store/pcs/integration/StoreIntegration.java事件消费端商品服务store-services/product-service/src/main/java/com/siriusxi/ms/store/ps/infra/MessageProcessor.java消息绑定配置RabbitMQ/Kafka 双 Binderconfig/repo/product.yml、config/repo/store.yml、config/repo/application.yml三形态编排文件docker-compose.yml、docker-compose-partitions.yml、docker-compose-kafka.yml消息发布测试store-services/store-service/src/test/java/com/siriusxi/ms/store/pcs/MessagingTests.java架构图docs/diagram/app_ms_landscape.png、docs/diagram/Zipkin.png10️⃣ 总结这套实战教会我们什么Springy-Store-Microservices 用最精简的商店场景完整演示了事件驱动架构在 Spring Cloud 生态中的最佳实践Spring Cloud Stream 是双中间件的万能转接头——代码只面向MessageChannelRabbitMQ 与 Kafka 靠defaultBinder与 profile 一键切换统一不可变事件模型让 CREATE/DELETE 语义清晰且key天然成为分区键重试、指数退避、死信队列三件套是消息可靠性的底线配置分区 多实例展示了 EDA 系统水平扩展的正确姿势Zipkin 管理端点补齐了异步系统的可观测性短板。如果你想深入理解一套代码、双消息中间件、事件驱动架构是如何落地的这个项目从Event.java到 docker-compose 每一层都值得逐行品读。【免费下载链接】Springy-Store-MicroservicesSpringy Store is a conceptual simple μServices-based project using the latest cutting-edge technologies, to demonstrate how the Store services are created to be a cloud-native and 12-factor app agnostic. Those μServices are developed based on Spring Boot Cloud framework that implements cloud-native intuitive, design patterns, and best practices.项目地址: https://gitcode.com/gh_mirrors/sp/Springy-Store-Microservices创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表