
简介基于Java语言实现的Kafka消息队列系统设计源码面向具备一定分布式基础的后端开发者用于解决大数据传输与实时数据处理场景下的消息解耦、高吞吐与可靠性问题。资源共42个文件压缩包约77.3MB涵盖27个Java源文件生产者、消费者及Kafka集群交互核心逻辑、6个XML配置文件数据库与服务器等关键参数、YAML与properties属性配置、SQL脚本初始化消息数据、集群搭建Markdown文档以及LICENSE、.gitignore等并附带Kafka与Zookeeper相关安装包便于本地快速构建环境。已有316人学习下载。通过源码与配置可系统掌握Kafka生产消费模型、集群配置、消息持久化与容错机制适合希望在大数据和实时计算方向提升工程能力的开发者研读参考。1. 基于 Java 语言的 Kafka 消息队列系统设计这套源码到底拆了什么先说结论如果你正在为一个高吞吐场景找可落地的 Java 榜样这份基于 Java 语言的 Kafka 消息队列系统设计源码省掉的不只是写 Demo 的时间——它把集群搭建、Spring Boot 集成、生产者消费者、数据库初始化全串在一条线上。我拆完第一反应是这不是给人看概念的这是给要上线的人垫底的。里面带的 kafka_2.13-2.7.0 和 apache-zookeeper-3.7.0-bin 两个归档包直接把服务端版本锁死不会出现网上教程那种版本错位导致的环境玄学。资源适合两类人一类是从零搭 Kafka 开发环境、想知道配置每一步在做什么的入门者另一类是已经跑过 Spring Boot 项目、但被消息丢失和重复消费磨得没脾气的开发。接下来我按自己复现这套工程的顺序从集群拉起写到调参排坑尽量让每个命令都解释清楚为什么这么写。2. 先把环境搭起来Zookeeper 3.7.0 与 Kafka 2.7.0 的集群配置和启动2.1 为什么这套设计绕不开 Zookeeper很多新人在看到 Kafka 2.7 时会有个疑问已经有 KRaft 模式了为什么项目里还给 Zookeeper 单独打个包因为 Kafka 2.7.0 这个版本里 KRaft 还不是稳定默认选项生产环境最常见、社区资料最完善的还是 Zookeeper 做协调节点的模式。Zookeeper 负责 broker 的注册、分区 leader 选举、消费组元数据维护而 Kafka 只专注消息的存储与分发。这套源码包把 Zookeeper 3.7.0 一起打进来说明作者的意图很清楚先用最成熟的拓扑把链路跑通。我一般建议不要在这个环节自作聪明升级到 KRaft除非你连 Zookeeper 的坑都没踩过。毕竟集群搭建.md 里的步骤大概率是按 ZK 模式写的你换模式等于给自己加戏。2.2 在 Linux 上把两个服务依次拉起来先把包解压建议目录层级分开不要直接放在同一层。解压命令如下# 解压两个归档包 tar -zxvf apache-zookeeper-3.7.0-bin.tar.gz tar -zxvf kafka_2.13-2.7.0.tgz # 进入 Zookeeper 配置目录把样例配置复制成正式配置 cd apache-zookeeper-3.7.0-bin/conf cp zoo_sample.cfg zoo.cfg这里的核心是把 ZooKeeper 的数据目录指向一个有足够磁盘空间的路径不要放在 /tmp否则重启一次丢掉全部元数据这个血泪经验后面再展开。进入 zoo.cfg 后至少要确认两个参数dataDir/opt/zk-data和admin.serverPort8081。后者容易被忽略Zookeeper 3.7.0 默认起一个 8080 的管理端口和本地 Spring Boot 项目端口撞车是常有的事。启动 Zookeeper 我习惯用前台方式观察日志确认没有异常再切后台cd apache-zookeeper-3.7.0-bin bin/zkServer.sh start-foreground看到绑定 2181 端口的日志后说明协调服务已经就绪。接着改 Kafka 的 server.properties这一步是整个集群能不能被人连上的关键。关键参数如下# kafka_2.13-2.7.0/config/server.properties 中必改项 broker.id0 log.dirs/opt/kafka-logs zookeeper.connectlocalhost:2181 advertised.listenersPLAINTEXT://localhost:9092advertised.listeners是最容易翻车的参数。很多教程只写 listeners不写 advertised结果本机连接正常一旦换台机器或换 IP 就连不上。这个参数是告诉客户端去连接哪个地址你写 localhost 就只适合单机测试真实集群要写具体 IP 或机器名。确认无误后再启动 Kafkacd kafka_2.13-2.7.0 bin/kafka-server-start.sh -daemon config/server.properties后面加-daemon表示后台运行但第一次启动我更推荐不加能直接看到有没有报错。2.3 先用主题命令验证环境健康度环境起没起来不要急着写代码先用自带 CLI 探底。创建一个专用主题bin/kafka-topics.sh --create \ --topic order-topic \ --partitions 3 \ --replication-factor 1 \ --bootstrap-server localhost:9092--partitions 3是我在这个阶段常用的分区数不是越大越好。分区数是 Kafka 并行度的上限单机测试时给太多分区反而让日志管理复杂。--replication-factor 1在单 broker 情况下只能写 1写 3 会直接报错。如果你用的 Windows注意命令要换成 bin\windows\kafka-topics.bat这是源码包里自带但很多教程不提的隐藏路径。这一步成功后说明 Zookeeper 和 broker 之间的连接已经通了可以进入 Spring Boot 集成环节。Windows 安装 Kafka 的读者尤其要记得先起 Zookeeper 再起 Kafka两个窗口不要关防火墙要放行 2181 和 9092 端口否则后面所有 Java 代码都会卡在连接超时。3. Spring Boot 集成要点pom、YAML 与生产者消费者核心源码3.1 配置链路pom.xml 中的依赖版本和 XML/YAML 边界这套源码里有两个 Maven 工程study-spring-kafka 和 kafka-study。前者是主工程后面其实是用来做隔离验证的简化版。打开 study-spring-kafka/pom.xml核心依赖是 Spring Kafka版本要和本机 Kafka 客户端匹配。我用的是 2.8.x 对应 Kafka 2.7 客户端这段依赖是项目能跑起来的底座dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId version2.8.6/version /dependency dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version2.7.0/version /dependency版本匹配是这套工程里最需要警惕的部分。Spring Kafka 2.8.6 的内部实现依赖 kafka-clients 2.7.0你要硬换成一个 3.x 的客户端编译能过但运行期可能冒出奇怪的反序列化异常。源码包里同时有 XML 配置和 YAML 配置我看下来 XML 主要是给老项目迁移用的模板新工程直接走 YAML 更简洁。application.yml 里核心配置如下spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: all retries: 3 batch-size: 16384 consumer: group-id: study-group enable-auto-commit: false auto-offset-reset: earliest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer这个 YAML 里的几个重点我要单独说。acks: all是确保消息写入所有副本后才返回成功单 broker 环境看不出差别但这是生产可用的最低门槛。enable-auto-commit: false是后续解决重复消费问题的前提后面避坑章会详细讲。auto-offset-reset: earliest表示消费组没有历史偏移量时从最早开始读调试阶段非常有用生产环境一般改成 latest否则上线第一天会把积压的老消息全部拉一遍。3.2 Java 生产者核心逻辑与 KafkaTemplate 使用看源码里生产者实现核心类就一个 KafkaTemplate它是 Spring 对原生 Producer 的封装。我会在业务代码里单独抽出 MessageProducer 类避免把发送逻辑散落在 Service 里Service public class MessageProducer { private static final Logger log LoggerFactory.getLogger(MessageProducer.class); private final KafkaTemplateString, String kafkaTemplate; public MessageProducer(KafkaTemplateString, String kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } public void sendMessage(String topic, String key, String message) { ListenableFutureSendResultString, String future kafkaTemplate.send(topic, key, message); future.addCallback( result - log.info(消息发送成功, topic{}, partition{}, offset{}, result.getRecordMetadata().topic(), result.getRecordMetadata().partition(), result.getRecordMetadata().offset()), ex - log.error(消息发送失败, key{}, message{}, key, message, ex) ); } }这段代码的逻辑说明send方法本身是异步的它立刻返回一个 ListenableFuture真正把数据写到 broker 是在后台线程完成的。如果你只调 send 不拿回调错误会在不打印任何日志的情况下被吞掉这是消息丢失的第一大隐藏原因。我在回调里故意把 partition 和 offset 打出来是为了快速判断消息被路由到哪个分区。参数方面key 如果为 nullKafka 会用轮询方式选分区如果指定 key则按 key 的哈希值选分区相同 key 的消息一定进同一个分区这保证了局部顺序性。3.3 Java 消费者核心逻辑与监听容器工厂消费者部分源码用的注解式监听器是最省事的方案。一个 KafkaListener 注解挂在方法上Spring 容器就会为它创建监听线程Component public class MessageConsumer { private static final Logger log LoggerFactory.getLogger(MessageConsumer.class); KafkaListener(topics order-topic, groupId study-group) public void onMessage(ConsumerRecordString, String record) { log.info(收到消息, key{}, value{}, partition{}, offset{}, record.key(), record.value(), record.partition(), record.offset()); // 在这里写真正业务逻辑入库、调用接口、写文件等 } }这里有个细节很多人不留意KafkaListener 方法里如果抛出异常Spring 默认行为是让消息处理失败并可能触发重试逻辑。如果你没有配置异常处理器异常会打印堆栈但不会阻塞消费线程。我一般会配一个专门的 ErrorHandler把失败消息记录到一张失败表而不是让 Kafka 无限重投。关于并发消费KafkaListener 注解支持 concurrency 属性它控制的是监听器容器的线程数。默认情况下一个分区的消息只会被一个消费者线程处理你即使把 concurrency 调到 10如果没有对应数量的分区多出的线程也只会空转。源码里具体数字写在 application.yml 的 spring.kafka.listener.concurrency 位置看到这个属性时不要盲目调大它必须和主题分区数保持一定关系否则白白浪费线程资源。4. 避坑指南重复消费、NoBrokersAvailable 与消息超限的排查顺序4.1 重复消费自动提交偏移量埋下的定时炸弹现象消费者明明把消息处理完了数据库里记录了这条数据重启服务或消费者组发生 rebalance 后同一条消息被再次消费业务数据出现重复。我见过的是订单回调接口同一笔订单被通知了好几次。原因Spring Kafka 默认 enable-auto-commit 是 true客户端每 5 秒自动提交一次偏移量。如果消息处理时间超过了这个间隔或者处理完成到提交之间的窗口期发生崩溃broker 里记录的偏移量还停留在上一次提交位置下一次 poll 就会把已处理过的消息重新拉回来。解决在 application.yml 里设置enable-auto-commit: false然后在监听容器工厂里指定 AckMode.MANUAL代码里处理完业务逻辑后手动提交偏移量。注意是处理完之后提交不是刚收到消息就提交。我在代码里会这样处理KafkaListener(topics order-topic, groupId study-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { try { // 业务逻辑幂等写入 orderService.save(record.value()); ack.acknowledge(); } catch (Exception e) { log.error(消息处理失败, 跳过手动提交, e); } }这个方案并不能彻底解决重复消费它只是把提交时机延后到业务成功之后。真正彻底解决重复需要消费端做幂等例如数据库主键唯一约束或 Redis setnx。源码里 test.sql 那条建表语句如果设置了唯一索引就是在配合这条策略。4.2 连接报错 NoBrokersAvailableadvertised.listeners 配置不一致现象本地 Kafka 生产消息时抛 NoBrokersAvailable但明明 kafka-server-start.sh 已经跑起来用 kafka-topics.sh 也能列出主题唯独 Java 客户端连不上。Windows 安装 Kafka 时这个报错尤其多。原因broker 实际监听地址和客户端解析地址不一致。只配置了 listenersPLAINTEXT://0.0.0.0:9092 或默认的 localhost而客户端从另一个机器访问时broker 返回的 advertised address 是 localhost客户端自然连不上。解决把 server.properties 里的 advertised.listeners 显式写成局域网 IP例如advertised.listenersPLAINTEXT://192.168.1.10:9092。如果是同一台机器调试写成 localhost 没问题但注意这是唯一允许的情况。改完配置必须重启 broker这个参数在运行时不会热更新。4.3 消息延迟高消费者线程和拉取参数不匹配现象生产者发消息很快但消费者侧积压持续上涨监控里看到消息延迟高到几十秒甚至分钟级。这时候先从自身找问题别急着怪网络。原因我排查过的案例里八成是消费者并发度过低或 fetch.max.poll.records 太小。默认情况下一个消费者只拉一条批次的记录如果单条消息处理逻辑很重吞吐就会被拖垮。另一种原因是 max.poll.interval.ms 设置过短处理太慢导致消费者被认为失联触发 rebalance进一步放大延迟。解决优先增加监听器 concurrency让它等于主题分区数的值然后把 fetch.max.poll.records 从默认的 500 往下调到一个合理值比如 100减少单次拉取后的处理压力。参数之间是联动的不要把 fetch.max.poll.records 调大同时又把 max.poll.interval.ms 调小那是自杀式配置。源码里 kafka-study 工程给了一个偏保守的参数组合适合用来对照。4.4 消息超限RecordTooLargeException 的默认 1M 限制现象发送一条稍大的 JSON 消息生产者报 RecordTooLargeException消费者端也可能因为拉取上限过小而报同样错误。原因broker 端 message.max.bytes 默认 1MB生产者 max.request.size 默认 1MB消费者端 max.partition.fetch.bytes 默认 1MB。任何一个环节不放开都会拒绝大消息。这是 Kafka 默认值的一个经典边界网上搜 kafka 接收1m 指的就是这个问题。解决需要同时改 broker、生产者、消费者三处配置只改一处没有任何作用。下一章我会给出完整参数表这里先记住一个原则单条消息上限不是单一参数决定的它是整条链路的共同约束。5. 调优实践把单条消息上限推到 1M 以上并压低消息延迟5.1 1MB 限制的完整配置链和参数表如果你需要传递超过默认 1MB 的消息比如大图片转 base64 或一批批量数据需要按下面这张表逐项核对。我实际测过100KB 消息不加任何配置没问题到了 1.2MB 就会触发报错边界卡得很死。配置位置参数名默认值建议值说明Broker 全局message.max.bytes100001210485760broker 接受单条消息上限Topic 级max.message.bytes100001210485760覆盖全局配置只作用于指定主题生产者max.request.size104857610485760单次请求最大字节消费者fetch.max.bytes5242880052428800一次 fetch 能拉取的所有分区总和上限消费者max.partition.fetch.bytes104857610485760单个分区单次拉取上限要注意的是 5MB 和 10MB 这种值不要拍脑袋写放得太大意味着单条消息会占用更多内存GC 压力和网络传输风险都会上升。我一般建议按业务里最大消息的 1.5 倍设置比如最大消息 3MB就设 5242880 左右。改完 broker 配置要重启生产者消费者则重新启动应用即可。在推送大消息时源码中生产者代码可以加一个手动指定 properties 的方法而不是全部依赖 YAMLMapString, Object props new HashMap(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 10485760); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 5);这段参数说明MAX_REQUEST_SIZE_CONFIG 直接对应 5.1 表格的生产者上限不设它就会吃默认亏。LINGER_MS_CONFIG 设为 5 毫秒是延迟和吞吐之间的折中下面单独展开。5.2 消息延迟高的另一面linger.ms、batch.size 与并发度的取舍很多人困惑为什么同一个主题别人跑起来延迟低自己跑起来像爬。我刚接触 Kafka 时也以为调大 batch 就能提吞吐结果延迟反而变高。这里有个核心概念linger.ms 控制生产者在发送前愿意等待多久来攒批。默认值是 0也就是来一条发一条延迟最低但吞吐也最低。你把 linger.ms 调到 20 毫秒吞吐会显著上升但极端情况下单条消息可能多等 20 毫秒才发送。我在调这种场景时一般是延迟敏感型业务保持 linger.ms0 或 5吞吐优先型设 50 毫秒以上。batch.size 是攒批的字节阈值默认 16KB如果单条消息 10KB那一个批次只能装一条批次效果打折扣。要注意 batch.size 设置得过大不会提高吞吐只会浪费内存。消费者端同理fetch.min.bytes 默认 1 字节表示有数据就拉如果改成 1024broker 会等攒够 1KB 再返回这会积累更多消息但牺牲响应速度。源码 kafka-study 工程里的配置已经给出一个偏稳定的组合建议先按它的来再根据监控调整。5.3 给 Kafka 配上可视化管理界面实际操作中命令行工具能查 LAG 和主题但看趋势和排查问题效率太低。有 UI 界面的工具能直接省掉大量重复劳动。我看到源码包里没有自带 UI 组件如果需要的话常见的开源方案可以选 Kafka UI 这类 Web 项目或者 Offset Explorer 这种桌面客户端。它们功能区别不大核心都要填 bootstrap-server 地址也就是 localhost:9092。有个细节要注意新版 UI 工具会要求填 Schema Registry 地址不用的场景直接留空就行别被必填项吓到。以 Kafka UI 为例配置好后能直接看到主题分区、消费者组、每个分区的 LAG 曲线。我这里提 LAG 是因为它是排查重复消费和延迟高的第一指标比盯日志直观得多。UI 工具不能改 broker 的 message.max.bytes 这类参数改配置还是得到 server.properties但用来观察是够用的。6. 用 test.sql 和消费者组做一次端到端验收验证链路真正通6.1 先把数据库表建起来配合消费者写库这套源码里带了一个 test.sql它在整个工程里的角色很明确给消费端提供一个落库目标让整个链路不止停留在日志输出层面。我通常会先执行这个脚本在本地 MySQL 建一张消息记录表字段至少包含消息内容、分区号、偏移量、创建时间。建表语句大致如下CREATE TABLE IF NOT EXISTS message_record ( id BIGINT AUTO_INCREMENT PRIMARY KEY, topic_name VARCHAR(64) NOT NULL, partition_no INT NOT NULL, offset_no BIGINT NOT NULL, message_text TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_offset (topic_name, partition_no, offset_no) );这个唯一键 uk_offset 是故意设计的它保证同一个主题、同一个分区、同一个偏移量的消息不会因为重复消费插入两次从消费端做幂等。测试时先用 kafka-console-producer 发送几条消息bin/kafka-console-producer.sh \ --broker-list localhost:9092 \ --topic order-topic输入任意字符串回车然后去看数据库表能看到对应记录就行。这一步验证的是消费者代码真的在跑而不是只在 IDE 里编译通过。6.2 用消费者组命令检查积压量和提交位置端到端验证的关键一步是看消费者组的偏移量状态。命令如下bin/kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --group study-group \ --describe输出里最值得关注的两列是 CURRENT-OFFSET 和 LOG-END-OFFSET。前者是消费组已经提交到的位置后者是主题日志末尾位置。两者相减就是 LAG即待消费消息数。如果 LAG 持续为 0证明消费者已经跟上生产速度如果 LAG 不断增长说明消费者处理能力不足回到第 5 章去调参数。这个描述命令是排查消息延迟高问题时我的首选工具比猜配置靠谱得多。6.3 把验证变成固定动作在那之前我经常因为只做了一轮消息收发测试就认为系统没问题结果在生产环境第一个晚上就碰到重复消费。从那以后我每次接手新的 Kafka 工程都会强制走一遍从集群启动到生产者发送、消费者落库、查看 LAG 为 0 的完整流程确认无误后再把源码里的配置往测试环境迁。这套验证流程不复杂但能挡掉至少一半的线上翻车。希望整套源码和这篇拆解能帮到你。本文还有配套的精品资源点击获取