资讯中心

搭建MQ消息组件Kafka服务环境:Lottery分布式抽奖系统的消息解耦实战

📅 2026/9/27 23:17:18
搭建MQ消息组件Kafka服务环境:Lottery分布式抽奖系统的消息解耦实战
文档教程后端【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址https://gitcode.com/gh_mirrors/code/CodeGuide点击查看免费下载本篇技术指南聚焦于 Lottery 分布式抽奖系统基于 DDD 领域驱动设计的四层架构实践中 MQ 消息组件的落地从 Kafka 的安装配置、主题Topic创建到与 SpringBoot 整合完成消息的生产与消费再到将 MQ 串联进抽奖 → 发奖的解耦流程。读完本文你将掌握在本地与 Docker 环境中搭建 Kafka 服务、按业务语义配置消息主题、并通过 SpringBoot 验证消息收发的一整套可复现方案。一、为什么抽奖系统需要 KafkaMQ 解耦的引入背景在 Lottery 抽奖系统的整体流程中用户抽奖与奖品发放原本是一条串行的长链路——抽奖完成后立刻执行发奖流程过长会导致用户长时间等待。为此项目采用 MQ 消息组件将流程解耦抽奖完成即发送 MQ 消息异步驱动后续的发奖流程避免抽奖与发奖互相阻塞参见 Lottery 面试技能、简历、问题汇总 中的项目介绍。这也是 第16节使用MQ解耦抽奖发货流程 的核心工作。在本节第15节搭建MQ消息组件Kafka服务环境的开发日志中明确记录了两项任务搭建 Kafka 环境配置消息主题。注意MQ 消息的使用不非得局限于 Kafka也可以使用 RocketMq 等其他消息中间件SpringBoot 整合 Kafka验证消息的生产和消费。Apache Kafka 是一个分布式发布-订阅消息系统也是一个强大的队列可以处理大量的数据并允许将消息从一个端点传递到另一个端点。Kafka 适合离线和在线消息消费消息保留在磁盘上并在群集内复制以防止数据丢失它构建在 ZooKeeper 同步服务之上并与 Apache Storm 和 Spark 有着良好的集成常用于实时流式数据分析。从源码结构与工程演进看Lottery 之所以选择 Kafka 作为 MQ 载体正是看中了它面向高吞吐、高可靠的以下四大特性特性说明可靠性Kafka 是分布式、分区、复制且容错的消息在集群内复制可防止数据丢失可扩展性Kafka 消息传递系统可以轻松扩展无需停机耐用性Kafka 使用分布式提交日志消息尽可能快地保留在磁盘上因此是持久的性能Kafka 对于发布和订阅消息都具有高吞吐量即使存储了许多 TB 的消息也保持稳定的性能Kafka 非常快并保证零停机和零数据丢失。这些特性与抽奖系统瞬时峰值高、历史数据不一定多、但必须快速响应的业务特征详见 notes.md 中关于分库分表理由的阐述是匹配的。二、Kafka 安装与启动本地命令行方式Kafka 依赖 ZooKeeper 完成分布式协调因此标准的本地启动流程分为三步先启动 ZooKeeper再启动 Kafka Broker最后创建消息主题。以下命令基于 Kafka 安装目录执行相关命令同样记录在 第16节使用MQ解耦抽奖发货流程 中# 1. 启动 ZooKeeper后台运行 bin/zookeeper-server-start.sh -daemon config/zookeeper.properties # 2. 启动 Kafka Broker后台运行 bin/kafka-server-start.sh -daemon config/server.properties # 3. 创建消息主题Topic bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic lottery_invoice命令参数说明zookeeper-server-start.sh读取config/zookeeper.properties配置启动协调服务-daemon表示后台守护进程方式运行kafka-server-start.sh读取config/server.properties启动 Broker 节点其中server.properties内可配置broker.id、listeners、log.dirs等核心参数kafka-topics.sh --create用于创建主题关键参数包括--zookeeper localhost:2181指定 ZooKeeper 连接地址--replication-factor 1副本因子为 1本地单机环境--partitions 1分区数为 1本地验证环境--topic lottery_invoice主题名语义即发货单消息抽奖完成后向该主题发送发货单消息再由消费端异步处理发货流程。在抽奖系统正式落地时主题的分区数与副本数需要结合消息量与可用性要求调整分区数影响消费并行度副本数影响容错能力。三、Kafka 安装与启动Docker 容器方式除了本地命令行方式第03节部署环境 Kafka 提供了面向云服务器/容器环境的 Docker 部署方案适合在 Lottery 抽奖系统项目介绍 中提到的 Docker 运维实践环节使用。# 1. 拉取镜像 docker pull wurstmeister/kafka docker pull wurstmeister/zookeeper # 2. 启动 ZooKeeper 容器 docker run -d --name zookeeper -p 2181:2181 -t wurstmeister/zookeeper# 3. 启动 Kafka 容器需传入 ZooKeeper 连接地址等环境变量 docker run -d --name kafka \ -p 9092:9092 \ -e KAFKA_BROKER_ID0 \ -e KAFKA_ZOOKEEPER_CONNECTzookeeper:2181 \ -e KAFKA_ADVERTISED_LISTENERSPLAINTEXT://宿主机IP:9092 \ -e KAFKA_LISTENERSPLAINTEXT://0.0.0.0:9092 \ wurstmeister/kafka要点说明-p 2181:2181将容器内 ZooKeeper 端口映射到宿主机供 Kafka 与本地应用访问Kafka 容器必须通过KAFKA_ZOOKEEPER_CONNECT指定 ZooKeeper 地址KAFKA_ADVERTISED_LISTENERS是容器部署最容易踩坑的参数它告诉客户端去哪里连接 Kafka生产环境务必配置为宿主机可达的 IP 或域名否则本地程序消费端会因拿到容器内地址而连接失败容器启动后同样使用kafka-topics.sh在后台添加抽奖系统需要的 Topic 主题并在本地程序中验证消息的生产与消费参见 第03节部署环境 Kafka 的运维日志。四、SpringBoot 整合 Kafka消息的生产与消费骨架本节的第二个开发日志项是SpringBoot 整合 Kafka验证消息的生产和消费。在 Lottery 的 SpringBoot 工程中Kafka 集成通常由以下三部分构成下为标准 Spring Boot Kafka 的通用实现骨架可对照本节分支211023_xfg_mq_kafka的提交内容进行验证1. 引入依赖dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependency2. 配置连接与序列化在application.yml中配置 Kafka 连接地址与消息序列化方式spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer consumer: group-id: lottery-invoice-consumer key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializerbootstrap-serversBroker 地址列表必须与前面启动的 Kafka 端口本地 9092 / Docker 映射端口保持一致group-id消费者组 ID同组消费者共同分摊该主题的消息若消息体为 JSON 对象可将 value 序列化器替换为JsonSerializer/JsonDeserializer并配合spring.kafka.producer.properties指定目标类型包路径。3. 消息生产与消费骨架// 生产者通过 KafkaTemplate 发送发货单消息 Service public class InvoiceMessageProducer { private final KafkaTemplateString, Object kafkaTemplate; public InvoiceMessageProducer(KafkaTemplateString, Object kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } public void sendInvoice(String topic, String orderId, Object payload) { kafkaTemplate.send(topic, orderId, payload); } } // 消费者监听 lottery_invoice 主题异步触发发奖流程 Component public class InvoiceMessageConsumer { KafkaListener(topics lottery_invoice) public void onMessage(ConsumerRecordString, Object record) { // 解析发货单消息驱动后续异步发奖处理 } }生产端通过KafkaTemplate.send(topic, key, value)发送消息消费端通过KafkaListener(topics lottery_invoice)订阅主题在消息到达时触发发奖逻辑。这样即完成了本节SpringBoot 整合 Kafka验证消息的生产和消费的验证闭环。五、MQ 解耦抽奖发货流程主题落地与可靠性补偿Kafka 环境与主题就绪后MQ 真正融入抽奖业务是在 第16节使用MQ解耦抽奖发货流程 中完成的其核心工作包括在数据库表user_strategy_export中添加字段mq_state用于在 MQ 发送成功后更新库表状态如果 MQ 消息发送失败则通过定时任务补偿 MQ 消息可使用该分支下的 sql 更新自己的库表启动 Kafka 新增主题lottery_invoice用于发货单消息抽奖完成后发送发货单再异步处理发货流程——这就是 MQ 解耦流程的使用在ActivityProcessImpl#doDrawProcess活动抽奖流程编排中补全用户抽奖后发送 MQ 触达异步奖品发送的流程。从消息可靠性闭环看第18节扫描库表补偿发货单MQ消息 进一步补全了补偿机制通过分布式任务调度扫描抽奖发货单消息状态对未发送 MQ或发送失败的消息进行补偿发送处理保障全流程可靠性。实现上需要分库分表路由组件db-router-spring-boot-starter支持设置路由到的库和表以便循环扫描每个库下的多张表中的每条用户记录。至此抽奖系统形成了完整的消息链路用户抽奖 → 抽奖结果落库(mq_state) → 发送 MQ(lottery_invoice) → 异步消费 → 发奖流程 │ └── 发送失败 → 定时任务扫描补偿(第18节)六、小结本节完成了 Lottery 抽奖系统 MQ 消息组件的环境搭建本地命令行方式启动 ZooKeeper 与 Kafka、创建lottery_invoice发货单主题Docker 方式部署 Kafka 到云环境并验证了 SpringBoot 生产者/消费者的消息收发。后续章节在此基础上完成了 MQ 解耦发奖第16节与定时补偿第18节构成抽奖 → 消息 → 异步发奖 → 失败补偿的完整可靠性链路这也正是抽奖项目简历中解耦抽奖流程把抽奖和发奖用MQ消息串联起来这一条核心经验的落点。在面试中关于为什么使用 MQ、为什么选择 Kafka、消息丢失与幂等如何处理等追问都可以基于本节的搭建过程与后续章节的补偿设计给出完整回答参考 notes.md 面试问题汇总。赞分享文档教程后端【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址https://gitcode.com/gh_mirrors/code/CodeGuide点击查看免费下载相关推荐Lottery 分布式抽奖系统用简单工厂搭建发奖领域服务Lottery 分布式抽奖系统用简单工厂搭建发奖领域服务 本文以 Lottery 分布式抽奖系统 Part 2 第 07 节“简单工厂搭建发奖领域”为主线讲文档教程后端云服务器部署 Docker 实战为 Lottery 抽奖系统搭建容器环境与 Portainer 面板云服务器部署 Docker 实战为 Lottery 抽奖系统搭建容器环境与 Portainer 面板 在 Lottery 抽奖系统基于 DDD 四层架构的分文档教程后端EasyWeChat 3.x 消息体系全解析统一抽象的消息类型、服务端回复与客服消息实战EasyWeChat 3.x 消息体系全解析统一抽象的消息类型、服务端回复与客服消息实战 EasyWeChat 将微信开放平台 API 中形形色色的消息后端即时通讯创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

看完文章,想为自己的企业也做一次专业网站诊断?

尧图顾问免费为您评估现有网站,并给出建站/改版建议与报价方案。

免费获取方案