)
示例工程教程【免费下载链接】java-design-patternsDesign patterns implemented in Java项目地址https://gitcode.com/GitHub_Trending/ja/java-design-patterns点击查看免费下载本文基于 java-design-patterns 仓库中的microservices-messaging模块深入讲解 Microservices Messaging微服务消息设计模式服务如何通过消息代理Message Broker异步交换消息实现解耦、可扩展与容错。读完本文你将掌握该模式的核心概念、在 Java 中的完整实现含Message、KafkaMessageProducer、KafkaMessageConsumer与订单服务示例并能在本地通过 Docker Compose Kafka 一键运行演示程序、查看运行日志同时了解消费组、序列化、offset 管理等底层细节。模式概览又名与意图Microservices Messaging 模式也被称为Asynchronous Messaging异步消息Event-Driven Communication事件驱动通信Message-Oriented Middleware / MOM面向消息的中间件其核心意图是让微服务之间通过消息传递进行异步通信服务不再直接互相调用而是经由消息代理所管理的话题Topic交换消息从而获得更好的解耦能力、伸缩性与容错能力。Wikipedia 对面向消息中间件的定义与本模式一脉相承MOM 是支持分布式系统之间收发消息的软件或硬件基础设施它允许应用模块分布在异构平台上并降低跨操作系统与网络协议开发应用的复杂度。真实世界示例电商订单处理设想一个电商平台客户下单后**Order Service订单服务**向消息代理发布一条 Order Created 消息。多个服务同时监听该消息Inventory Service更新库存水位Payment Service处理支付Notification Service发送确认邮件。每个服务都独立运行按自己的节奏处理消息、互不阻塞。如果 Payment Service 暂时宕机消息代理会持有这条消息直到它恢复从而确保数据不丢失。用一句话概括该模式让服务通过消息代理异步通信彼此无需等待即可独立工作。下图展示了整个消息流的整体流程Java 程序化示例订单处理系统本模块microservices-messaging演示了服务如何通过消息代理通信而无需直接耦合。示例是一个订单处理系统OrderService作为生产者向 Kafka 发布订单事件InventoryService、PaymentService、NotificationService作为消费者分别响应。消息载体Message 类Message类表示服务间交换的数据由唯一 ID、内容与时间戳三部分组成Getter public class Message { private final String id; private final String content; private final LocalDateTime timestamp; public Message(String content) { this.id UUID.randomUUID().toString(); this.content content; this.timestamp LocalDateTime.now(); } // JSON constructor for deserialization JsonCreator public Message( JsonProperty(id) String id, JsonProperty(content) String content, JsonProperty(timestamp) LocalDateTime timestamp) { this.id id; this.content content; this.timestamp timestamp; } Override public String toString() { /* ... */ } }与 README 中的简化版本不同Message.java 还提供了带JsonCreator/JsonProperty注解的 JSON 反序列化构造器这为后续通过 Jackson 在 Kafka 线路上传输对象提供了基础。id由UUID.randomUUID()生成可作为消息幂等判重的天然标识。消息代理KafkaMessageProducer 与 KafkaMessageConsumerREADME 用内存中的ConcurrentHashMap演示了MessageBroker的路由思想subscribe按 topic 注册处理器publish向订阅者分发。而本仓库的实际实现则使用Apache Kafka 作为生产级消息代理对应两个核心类生产者 KafkaMessageProducer.javapublic class KafkaMessageProducer implements AutoCloseable { private final ProducerString, String producer; private final ObjectMapper objectMapper; public KafkaMessageProducer(String bootstrapServers) { this(createDefaultProducer(bootstrapServers)); } private static ProducerString, String createDefaultProducer(String bootstrapServers) { Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.RETRIES_CONFIG, 3); return new KafkaProducer(props); } public void publish(String topic, Message message) { try { String json objectMapper.writeValueAsString(message); ProducerRecordString, String record new ProducerRecord(topic, message.getId(), json); producer.send(record, (metadata, exception) - { /* 日志记录 partition 与 offset */ }); } catch (Exception e) { LOGGER.error(Error serializing message: {}, e.getMessage(), e); } } Override public void close() { producer.flush(); producer.close(); } }值得注意的生产者配置细节ACKS_CONFIG all要求所有副本确认写入牺牲部分吞吐换取更强的一致性保障RETRIES_CONFIG 3发送失败最多重试 3 次提升可靠性key 使用消息 IDvalue 使用 Jackson 序列化后的 JSON注册了JavaTimeModule以正确序列化LocalDateTime时间戳发送回调会记录目标 topic 的partition与offset方便追踪每条消息在 Kafka 中的落盘位置。消费者 KafkaMessageConsumer.javapublic class KafkaMessageConsumer implements AutoCloseable, Runnable { private final ConsumerString, String consumer; private final ObjectMapper objectMapper; private final String topic; private final java.util.function.ConsumerMessage messageHandler; private final AtomicBoolean running new AtomicBoolean(true); public KafkaMessageConsumer( String bootstrapServers, String groupId, String topic, java.util.function.ConsumerMessage messageHandler) { this(createDefaultConsumer(bootstrapServers, groupId), topic, messageHandler); } private static ConsumerString, String createDefaultConsumer(String bootstrapServers, String groupId) { Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true); return new KafkaConsumer(props); } Override public void run() { consumer.subscribe(Collections.singletonList(topic)); while (running.get()) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); records.forEach(record - { Message message objectMapper.readValue(record.value(), Message.class); messageHandler.accept(message); }); } consumer.close(); } public void stop() { running.set(false); } }关键消费者配置解读GROUP_ID_CONFIG每个服务使用独立的消费组inventory-group、payment-group、notification-group使同一 topic 的消息能被多个服务各自完整消费——这是多个服务同时监听同一事件的核心机制AUTO_OFFSET_RESET_CONFIG earliest从最早 offset 开始消费服务重启后不会丢消息ENABLE_AUTO_COMMIT_CONFIG true开启自动提交 offset简化代码在真实生产场景中通常需要改为手动提交以配合幂等处理消费者实现Runnable可放入线程池独立运行stop()通过AtomicBoolean安全地终止轮询循环。生产者角色OrderServiceOrderService.java 是消息生产者发布三类订单消息到order-topicpublic class OrderService { private static final String ORDER_TOPIC order-topic; private final KafkaMessageProducer producer; public OrderService(KafkaMessageProducer producer) { this.producer producer; } public void createOrder(String orderId) { LOGGER.info(Creating order: {}, orderId); Message message new Message(Order Created: orderId); producer.publish(ORDER_TOPIC, message); } public void updateOrder(String orderId) { Message message new Message(Order Updated: orderId); producer.publish(ORDER_TOPIC, message); } public void cancelOrder(String orderId) { Message message new Message(Order Cancelled: orderId); producer.publish(ORDER_TOPIC, message); } }OrderService只依赖KafkaMessageProducer抽象完全不感知下游有哪些消费者——这正是生产端与消费端解耦的直接体现。消费者角色三个独立服务InventoryService库存服务响应Order Created时updateInventory预留库存响应Order Cancelled时restoreInventory释放库存其余消息仅记录 debug 日志public void handleMessage(Message message) { if (message.getContent().contains(Order Created)) { updateInventory(message); } else if (message.getContent().contains(Order Cancelled)) { restoreInventory(message); } else { LOGGER.debug(No inventory action needed for: {}, message.getContent()); } }PaymentService支付服务响应Order Created时processPayment扣款响应Order Cancelled时refundPayment退款public void handleMessage(Message message) { if (message.getContent().contains(Order Created)) { processPayment(message); } else if (message.getContent().contains(Order Cancelled)) { refundPayment(message); } }NotificationService通知服务对Order Created/Order Updated/Order Cancelled分别发送订单确认、更新通知与取消通知。三个服务的实现文件分别为 InventoryService.java、PaymentService.java 与 NotificationService.java。它们的类注释均明确说明各自运行在独立消费组中从而能够独立处理消息、横向扩容同组内增加实例、并在重启后从上次提交的 offset 继续消费。主程序组装与运行App.java 完成组装创建生产者与三个消费者将它们提交到 3 线程的ExecutorService中并发运行依次演示创建、更新、取消订单最后优雅关闭public class App { private static final String BOOTSTRAP_SERVERS localhost:9092; static long sleepMs 2000; // 测试中可改为 0 加速 public static void main(String[] args) throws InterruptedException { KafkaMessageProducer producer new KafkaMessageProducer(BOOTSTRAP_SERVERS); KafkaMessageConsumer inventoryConsumer new KafkaMessageConsumer( BOOTSTRAP_SERVERS, inventory-group, order-topic, new InventoryService()::handleMessage); KafkaMessageConsumer paymentConsumer new KafkaMessageConsumer( BOOTSTRAP_SERVERS, payment-group, order-topic, new PaymentService()::handleMessage); KafkaMessageConsumer notificationConsumer new KafkaMessageConsumer( BOOTSTRAP_SERVERS, notification-group, order-topic, new NotificationService()::handleMessage); run(producer, inventoryConsumer, paymentConsumer, notificationConsumer); } }主程序运行时会看到类似如下日志README 中的输出示例Published order message: ORDER-123 Inventory Service received: Order Created: ORDER-123 Updating inventory... Payment Service received: Order Created: ORDER-123 Processing payment...实际接入 Kafka 后日志会更加丰富消费者先输出Consumer subscribed to topic: order-topic生产者回调会输出Published message to topic order-topic [partition0, offset0]每个服务还会按业务分支输出更新库存 / 处理支付 / 发送订单确认等明细。下面的时序图展示了App中生产者与各消费者之间的完整交互过程如何运行示例程序本模块的运行依赖一个可用的 Kafka 实例默认连接localhost:9092并提供三种方式方式一自动化脚本推荐在microservices-messaging模块目录下直接运行辅助脚本脚本会自动检测端口 9092 上是否已有 Kafka若没有且本机装有 Docker则自动通过 Docker Compose 拉起 Kafka 容器并等待其就绪随后编译并启动应用WindowsPowerShellpowershell -ExecutionPolicy Bypass -File .\run-app.ps1Linux / macOS./run-app.sh从 run-app.sh 的源码可以看到其执行逻辑先探测localhost:9092是否可达nc或/dev/tcp不可达且存在docker命令时执行docker compose up -d最多重试 20 次、每次间隔 2 秒等待 Kafka 就绪最后执行../mvnw compile exec:java -Dexec.mainClasscom.iluwatar.messaging.App。方式二Docker Compose 手动启动如果希望手动控制容器生命周期可先启动 Kafka 容器再运行应用# 在 microservices-messaging 目录下启动 Kafka 容器端口 9092 docker compose up -d # 编译并运行应用 ../mvnw compile exec:java -Dexec.mainClasscom.iluwatar.messaging.App # 结束演示后停止并移除容器 docker compose downdocker-compose.yml 使用confluentinc/cp-kafka:7.5.0镜像并在KRaft 模式KAFKA_PROCESS_ROLES: broker,controller无需 Zookeeper下运行关键环境变量包括环境变量值作用KAFKA_NODE_ID1Kafka 节点 IDKAFKA_ADVERTISED_LISTENERSPLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092容器内外可达地址宿主机通过localhost:9092访问KAFKA_LISTENERSPLAINTEXT://kafka:29092,CONTROLLER://kafka:29093,PLAINTEXT_HOST://0.0.0.0:9092三种监听器业务、控制器、宿主机入口KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR1单节点环境下 offset topic 副本因子KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS0消费组初始再平衡延迟设为 0加速演示启动KAFKA_PROCESS_ROLESbroker,controller单进程同时承担 broker 与 controller 角色KRaftCLUSTER_IDMkU3OEVBNTcwNTJENDM2QkKRaft 集群 ID方式三本地手工启动 Kafka如果不想用 Docker可参考 App.java 类注释中给出的手工启动步骤# 启动 Zookeeper bin/zookeeper-server-start.sh config/zookeeper.properties # 启动 Kafka bin/kafka-server-start.sh config/server.properties # 创建主题 bin/kafka-topics.sh --create --topic order-topic --bootstrap-server localhost:9092注意本模块的 Kafka 客户端版本为kafka-clients 3.9.2见 pom.xml需保证本地 Kafka 版本与之兼容docker compose up -d后应用第一次运行时会自动创建order-topicKafka 默认auto.create.topics.enabletrue。测试与验证模块提供了完整的单元测试src/test/java/com/iluwatar/messaging无需真实 Kafka 即可验证核心逻辑AppTest.java将App.sleepMs置为0加速运行并借助 Kafka 官方的MockProducer/MockConsumer构造无 broker 依赖的完整调用链验证App.run()不会抛出异常MessageTest、OrderServiceTest、InventoryServiceTest、PaymentServiceTest、NotificationServiceTest、KafkaMessageProducerTest、KafkaMessageConsumerTest分别验证消息结构、生产者发布与消费者接收/反序列化流程。依赖方面pom.xml 引入了kafka-clients、jackson-databind含jackson-datatype-jsr310支持时间序列化、SLF4J Logback 日志以及 JUnit 5 Mockito 测试框架。何时使用该模式服务之间需要非阻塞通信时系统需要组件间松散耦合时事件驱动架构中多个服务需要对同一事件作出响应时需要通过消息缓冲应对流量尖峰时分布式系统中部分服务可能临时不可用时。真实世界的应用场景使用 Apache Kafka、RabbitMQ 或 ActiveMQ 进行服务通信的 Java 应用电商平台的订单处理与库存管理金融服务的交易处理与通知推送IoT 系统的传感器数据处理与事件响应。收益与权衡收益服务松散耦合可独立开发与部署消息缓冲提升了服务临时不可用时的系统韧性支持发布/订阅、请求/响应等多种通信模式通过并行处理消息增强可扩展性天然契合事件驱动架构。权衡引入消息代理基础设施增加额外复杂度消息代理自身需要高可用部署采用最终一致性而非即时一致性异步流的调试比同步调用更困难需要处理消息重复投递并保证消费者幂等。相关设计模式本模块与仓库中多个模式密切相关以下链接已转换为仓库根目录相对路径Saga 模式使用消息编排分布式事务CQRS 模式常借助消息分离读写操作事件溯源Event Sourcing将状态变更以消息形式存储API Gateway 模式以同步请求补足消息通信的不足。若希望进一步理解消息驱动系统的演进可同时阅读仓库中的 event-driven-architecture、event-aggregator 与 publish-subscribe 等模块它们从不同角度展现了事件与消息在分布式系统中的应用。赞分享示例工程教程【免费下载链接】java-design-patternsDesign patterns implemented in Java项目地址https://gitcode.com/GitHub_Trending/ja/java-design-patterns点击查看免费下载相关推荐微服务架构模式实战详解服务拆分、通信与容错设计——基于 agents 仓库 microservices-patterns 技能微服务架构模式实战详解服务拆分、通信与容错设计——基于 agents 仓库 microservices patterns 技能 本指南以开源仓库 agentsAI 插件AI 技能开发工具GoogleCloudPlatform/microservices-demo微服务通信模式对比分析GoogleCloudPlatform/microservices demo微服务通信模式对比分析 概述 在现代微服务架构中通信模式的选择直接影响系统的性能后端微服务电商云原生XCP-ng自动化部署使用脚本打造高效虚拟化 infrastructureXCP ng自动化部署使用脚本打造高效虚拟化 infrastructure XCP ng作为一款强大的开源虚拟化平台为企业和个人用户提供了稳定可靠的虚拟化解上一篇NVIDIA Profile Inspector完全指南免费解锁显卡200隐藏设置游戏性能飙升30%下一篇NVIDIA Profile Inspector完全指南解锁显卡隐藏性能的终极工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考