Kafka 入门
1. 一句话简介
Apache Kafka 是一个分布式流处理平台,由 LinkedIn 开源、后捐献给 Apache 基金会。它把消息以仅追加(append-only)日志的方式写入磁盘上的分区(Partition),并基于 Topic + Partition 的分区模型组织数据,通过消费者组(Consumer Group)并行消费,从而在分布式环境下提供超高的吞吐能力与持久化保障。Kafka 解决的核问题是:在大量系统与服务之间进行高性能、可回溯、可持久化的异步消息与事件流传输,让生产者与消费者彻底解耦、互不感知。
以本 demo(demo-mq-kafka)为例,生产端通过 KafkaTemplate.send(topic, message) 发送消息,Broker 将消息落到 test Topic(3 个分区),消费端通过 @KafkaListener 注解监听该 Topic,并配合手动提交 offset(关闭自动提交)。这样一个「生产者—Broker—消费者」模型,就是 Kafka 最基础也最典型的使用形态。
2. 什么时候使用
✅ 适用场景
- 高吞吐、大流量的数据管道:Kafka 单机即可支撑每秒百万级消息写入,适合日志采集、埋点上报、监控指标等海量数据的汇聚与转发。
- 事件驱动架构与系统解耦:多个服务通过 Topic 订阅同一批事件,生产者无需关心消费者是否存在、何时消费。
- 需要消息回溯 / 重复消费:消息持久化到磁盘且按 offset 管理,允许消费者回退到任意历史位置重放,适合离线分析、重放修复等场景。
- 多消费者并行消费:Consumer Group 天然支持一个 Topic 被多个消费者实例分摊消费,利于水平扩容提升处理能力,正如 demo 中配置了
concurrency并发消费。 - 严格要求消息不丢失的业务:关闭自动提交、采用手动 offset 提交(
MANUAL_IMMEDIATE)可保证消息处理完才提交,提高可靠性,正如本 demo 的MessageHandler所演示。
❌ 不适用 / 需谨慎
- 业务系统内部的路由型消息:Kafka 以 Topic 发布-订阅为主,缺少 RabbitMQ 那种灵活的交换机路由模型,复杂的消息路由、请求-应答场景不占优势。
- 低延迟到毫秒级的要求:Kafka 重吞吐、采用拉取(Pull)模式,端到端延迟相对较高,不适合对实时性要求极苛刻的应用。
- 极低频、极简单的单机小应用:部署和运维 Kafka 集群成本较高(需协调器、分区副本管理),如果只是传几万条消息,引入它是明显的过度设计。
- 需要保证消息严格全局有序且分区数受限:Kafka 只在单个分区内保证有序,跨分区的全局顺序无法保证,对强全局顺序有诉求的需谨慎。
- 组件少、无大数据/Hadoop/流计算生态的中小团队:Kafka 的最大价值常与日志、流计算、大数据链路绑定,若没有这些生态背景,运维收益比偏低。
3. 常见业务场景
日志与埋点采集:把各业务服务、网关产生的访问日志、行为埋点统一写入 Kafka 的相应 Topic,再由日志分析系统或 Flink/Spark 二次消费,实现集中存储、离线统计与实时告警,Kafka 的高吞吐和持久化回溯在此发挥最大价值。
流数据处理管道:作为数据仓库和大数据计算框架之间的「数据动脉」,Kafka 将实时产生的数据流按主题组织,供多个下游(实时计算、数据湖、搜索引擎索引)并行消费,天然契合事件流分析的范式。
系统间异步解耦:订单、支付、库存等上游系统只负责把「订单已创建」等事件写入 Kafka,下游(通知、积分、审计)各自订阅对应 Topic 独立消费,上游不再等待下游同步返回,正如本 demo 中生产者只 send 到 test Topic、消费者在 MessageHandler 中异步接收,两端彻底解耦。
削峰填谷:秒杀、抢购等瞬时请求洪峰到来时,先把请求消息写入 Kafka 缓冲,后端按自身吞吐能力从容消费处理,避免突增流量直接压垮数据库或下游服务。
分布式系统最终一致性:跨服务的数据变更以事件形式按序落盘并消费,配合手动 offset 提交(如 demo 中的 acknowledge())保证消息不丢,结合补偿机制在多个服务间达成最终一致。
4. 同类技术对比
| 维度 | Apache Kafka | RabbitMQ | RocketMQ | ActiveMQ |
|---|---|---|---|---|
| 设计定位 | 分布式流处理平台 | 通用消息代理 | 分布式消息中间件 | 通用 JMS 消息代理 |
| 吞吐量 | 极高(百万级 TPS) | 中等(万级 TPS) | 高(数十万级 TPS) | 中等 |
| 消息模型 | 发布-订阅(Topic+Partition) | 灵活路由(Direct/Fanout/Topic/Headers 交换机) | 发布-订阅 + 延迟消息、事务 | JMS Queue/Topic |
| 消息持久化与回溯 | 默认落盘,支持按 offset 回溯与重放 | 消费后即删除(TTL/死信需额外配) | 持久化且支持按时间回溯 | 持久化到数据库/文件 |
| 交付语义 | 拉模式(Pull),支持 Exactly-Once | 推模式(Push),天然支持 ack/事务 | 拉模式,支持事务 | 推模式,支持事务 |
| 语言/生态 | Java,绑定大数据、流计算生态 | 多语言(Erlang),HTTP/AMQP | Java,阿里系生态 | Java,JMS 规范 |
| 分布式能力 | 天然分布式、分区副本、水平扩容 | 集群化但适度,非大数据规模设计 | 分布式、多主多从 | 集群能力一般 |
| 运维复杂度 | 较高(需协调器、分区/副本管理) | 较低(自带 Web 管理界面) | 中 | 低 |
| 适用规模与场景 | 大数据量、事件流、日志、海量吞吐 | 业务消息、任务队列、灵活路由 | 大流量业务消息、电商场景 | 传统 JMS 业务系统 |
选型建议
- 数据量大、要求超高通量,且链路涉及日志采集、事件驱动、流计算或大数据分析 → 优先选 Apache Kafka,这也是大数据生态的既定事实标准。
- 业务内部消息、需要灵活路由、团队希望低运维成本、对功能上手快 → 选 RabbitMQ。
- 日活规模大、业务消息可靠性要求高,且已深度绑定阿里云计算栈 → 选 RocketMQ,其延迟消息、事务消息等特性贴合电商业务。
- 老旧的 JMS 规范应用或历史迁移负担重 → 选 ActiveMQ,新项目不建议从零引入。
一句话:要「速度与数据量」选 Kafka,要「灵巧路由与易运维」选 RabbitMQ,电商强一致业务可考虑 RocketMQ。