跳转到内容

Apache Kafka — 分布式流处理平台

待复核

Kafka 是一个持久化日志 + 发布订阅系统,2010 年 LinkedIn 开发并在 2011 年开源给 Apache 基金会。一个集群每秒能处理百万条消息,被 LinkedIn、Uber、Netflix、Twitter 等大型实时数据系统当成数据管道的中枢。

日常类比:就像超市的传送带——商家(producer)把商品放到传送带上,多个收银员(consumer)按顺序从传送带上取货。传送带自己记录每件商品的位置(offset),收银员只要告诉传送带”我取到了第 42 号商品”,下次就能从第 43 号继续取。

你写一个最简单的发消息:

Terminal window
bin/kafka-console-producer --topic events --bootstrap-server localhost:9092
> hello kafka

另一个进程订阅同一个 topic:

Terminal window
bin/kafka-console-consumer --topic events --from-beginning --bootstrap-server localhost:9092
hello kafka

发送方不需要知道接收方在哪、是谁、什么时候来——这就是发布订阅最基础的解耦能力。

不理解 Kafka 的设计哲学,下面这些事都没法解释:

  • 为什么 LinkedIn / Uber / Netflix / Twitter 这种”每秒百万事件”的公司,数据管道几乎都是 Kafka
  • 为什么 Kafka Streams + ksqlDB 让流处理可以写得像查询——过滤、聚合直接跑在事件流上
  • 为什么 Kafka Connect 生态有大量现成连接器(MySQL / Elasticsearch / S3 等),让”进出管道”少写胶水代码
  • 为什么在高吞吐日志型管道场景里,Kafka 比传统 broker 型队列(RabbitMQ / ActiveMQ)更常被选作中枢

简单说:这是过去 10 年实时数据基础设施的核心开源项目之一,尤其擅长”先落盘再多人消费”的管道。

Kafka 的核心模型可以拆成 三层

  1. Topic + Partition:一个 topic 是一类消息的分类(如 user-clicksorder-events)。一个 topic 切成多个 partition,每个 partition 是一个独立的有序日志,可以并行写入和读取——这是 Kafka 横向扩展的根本。

  2. Producer + Consumer Group:producer 把消息发到 topic,Kafka 决定写入哪个 partition(按 key hash 或轮询)。consumer 订阅 topic 时加入一个 consumer group,同 group 内的 consumer 自动分摊 partition——4 个 partition、2 个 consumer,每人吃 2 个;3 个 consumer,每人 1-2 个。

  3. Offset + 持久化:每条消息在 partition 内有一个递增编号叫 offset。Kafka 把消息持久化到磁盘(默认保留 7 天),consumer 自己记录”我消费到 offset 多少了”。所以消费者挂了重启也能从断点继续,不会丢消息。

简单说:topic 是邮筒,partition 是邮筒里的并行格子,consumer group 是分工取信的一组邮差,offset 是邮差的进度书签

案例 1:Docker 起一个单节点 KRaft Kafka

Section titled “案例 1:Docker 起一个单节点 KRaft Kafka”

最快上手可用 Bitnami / Apache 官方 quickstart;下面是一份最小可跑的单节点 KRaft 示意(生产请用官方文档补全安全与磁盘配置):

# docker-compose.yml(示意:单 broker + controller 合部署)
services:
kafka:
image: apache/kafka:3.7.0
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk
Terminal window
docker compose up -d
# 再用容器内 kafka-topics / console-producer 验证连通

注意:这里没有 ZooKeeper——3.3+ 可用 KRaft(broker 自管元数据);新集群应优先 KRaft

Terminal window
# 创建 topic(3 个 partition = 最多 3 个同组 consumer 并行)
bin/kafka-topics.sh --create --topic events --partitions 3 \
--bootstrap-server localhost:9092
# 发消息
bin/kafka-console-producer.sh --topic events --bootstrap-server localhost:9092
> {"user": "alice", "action": "click"}
# 收消息(从头回放)
bin/kafka-console-consumer.sh --topic events --from-beginning \
--bootstrap-server localhost:9092

--from-beginning 是杀手锏:消息已落盘,新 consumer 仍可从最早 offset 重放。

两个终端都用 group analytics

Terminal window
# 终端 1
bin/kafka-console-consumer.sh --topic events --group analytics \
--bootstrap-server localhost:9092
# 终端 2
bin/kafka-console-consumer.sh --topic events --group analytics \
--bootstrap-server localhost:9092

Kafka 把 3 个 partition 分给两个 consumer(一人 2、一人 1)。再加第三个,正好一人一个。同 group 分摊 partition = 流处理水平扩展的基本实现。

  1. Partition 数预估错:partition 数是 topic 创建时定的,事后只能加不能减。设少了无法横向扩 consumer(consumer 数 ≤ partition 数);设多了 broker 元数据爆炸(每个 partition 占内存和文件句柄)。经验值:每 broker 不超过 2000-4000 partition。

  2. Consumer Rebalancing 期间消息处理停顿:consumer 加入或退出 group 时会触发 rebalance,所有 consumer 暂停消费几秒到几十秒。Kafka 2.4 引入 cooperative-sticky 分配策略缓解了这个问题——只重分配必要的 partition,不全量打散。

  3. Exactly-once 语义复杂:默认是 at-least-once(可能重复)。要做到 exactly-once 必须开启 idempotent producer + transactional consumer,配置一堆参数(enable.idempotencetransactional.idisolation.level=read_committed),且只在 Kafka 内部端到端有效——出了 Kafka 还是要业务层去重。

  4. 监控指标多到爆炸:consumer lag(落后多少消息)/ throughput(吞吐)/ disk usage / GC 时间 / network IO ——任何一个炸了集群都可能挂。生产环境必须上 Confluent Control Center 或 LinkedIn 开源的 Burrow,光看 broker 自己的 JMX 不够。

适用

  • 大型实时数据管道(用户行为、日志聚合、CDC 数据库变更捕获)
  • 微服务间事件驱动架构(订单事件 → 库存 + 物流 + 通知 多个下游)
  • 流处理(Kafka Streams / Flink / Spark Streaming 的输入源)
  • 需要消息持久化和回放的场景(消息保留 N 天,新业务上线可以从头消费)

不适用

  • 低延迟点对点通信(毫秒级 RPC 用 gRPC 更合适)
  • 小规模消息队列(几千 QPS 用 RabbitMQ / Redis Stream 更轻)
  • 严格 FIFO 跨多 partition(Kafka 只保证单 partition 内有序)
  • 复杂路由规则(topic exchange、header routing 这种用 RabbitMQ)
  • 2010 年:LinkedIn 工程师 Jay Kreps 主导开发,名字来源于作家 Franz Kafka——“我喜欢一个写作系统命名的项目,所以叫了 Kafka”。
  • 2011 年:开源给 Apache 基金会,进入孵化器。
  • 2014 年:Jay Kreps 等核心团队从 LinkedIn 离职创立 Confluent,主导 Kafka 商业化。
  • 2017 年:Kafka Streams 发布,第一个内嵌在 Kafka 里的流处理库(不需要 Spark/Flink)。
  • 2020 年:KIP-500 推进移除 ZooKeeper 依赖(KRaft 自管元数据)。
  • 2023–2024:KRaft 生产可用(约 3.3+);3.5 起 ZooKeeper 模式标记 deprecated;3.7 仍可跑 ZK 但新集群应选 KRaft;Kafka 4.0 移除 ZK。Tiered Storage 等让冷数据可下沉对象存储。

15 年从 LinkedIn 内部工具到全球数据管道标配。

  1. 持久化日志是个好抽象——把”消息队列”和”分布式日志”统一了,副作用变成事实记录
  2. Partition + Consumer Group 是横向扩展的范式——任何流处理系统都在重新发明这个模型
  3. 零拷贝 + 顺序写磁盘——Kafka 性能秘诀不是内存而是顺序 IO 比随机内存还快
  4. 生态比单点功能重要——Kafka Connect 的大量连接器让它常成为”系统之间的胶水层”
  • redis —— list / pub-sub / stream 也能做队列,但持久化与跨机扩展通常弱于 Kafka
  • flink —— 有状态流计算,生产里常订阅 Kafka topic
  • spark —— Spark Structured Streaming 常见数据源之一是 Kafka
  • rabbitmq —— 传统 broker 型消息队列,路由灵活,高吞吐日志管道场景常让位 Kafka
  • pulsar —— 另一套云原生日志/消息系统,常与 Kafka 对照选型
  • zookeeper —— 旧版 Kafka 元数据依赖;KRaft 后新集群不再需要
  • appwrite —— Appwrite — 自己能装一遍的开源 Firebase
  • asynq —— Asynq — Go 版 Sidekiq,把后台任务丢进 Redis 慢慢跑
  • bullmq —— BullMQ — Node.js 上的 Redis 任务队列
  • celery —— Celery — Python 把慢任务搬到后台干的工头
  • centrifugo —— Centrifugo — Go 写的开源实时消息服务器
  • debezium —— Debezium — 把数据库的”刚刚改了”变成消息流
  • docker —— Docker — 容器化平台
  • druid —— Apache Druid — 流批一体的实时分析数据库
  • ejabberd —— ejabberd — Erlang 写的电信级 XMPP/MQTT 多协议服务器
  • elasticsearch —— Elasticsearch — 分布式搜索引擎
  • emqx —— EMQX — 单集群千万连接的 MQTT 物联网消息总线
  • encore —— Encore — 类型安全 Go/TS 后端框架,基础设施即代码
  • grpc-go —— gRPC-Go — Google RPC 框架的官方 Go 实现
  • inngest —— Inngest — 让 async 函数自动从断点恢复的工作流引擎
  • memgraph —— Memgraph — 内存图数据库
  • mosquitto —— Mosquitto — C 写的轻量 MQTT 消息中转站
  • nats —— NATS — 极简云原生消息系统
  • nats-server —— NATS Server — 极简云原生消息总线
  • nginx —— nginx — 高性能 Web 服务器
  • nsq —— NSQ — Go 写的去中心化消息队列
  • openhab —— openHAB — Java OSGi 家庭自动化框架
  • orleans —— Orleans — 让分布式服务写起来像单机对象
  • pg-boss-readme —— pg-boss — 只用 Postgres 就能跑的任务队列
  • pino —— pino — 日志不该阻塞热路径
  • pinot —— Apache Pinot — LinkedIn 起家的实时 OLAP
  • postfix —— Postfix — 把 sendmail 拆成一群最小权限的小工
  • prom-client —— prom-client — Node 服务暴露监控指标的事实标准 SDK
  • pulsar —— Apache Pulsar — 云原生消息队列
  • quarkus —— Quarkus — 让 Java 启动比 Node 还快的云原生框架
  • rabbitmq-server —— RabbitMQ — 用 Erlang 写的多协议消息总线
  • redis —— Redis — 内存键值数据库
  • redpanda —— Redpanda — Kafka 兼容的 C++ 实现
  • rethinkdb —— RethinkDB — 让数据库自己把更新推给客户端的先驱
  • risingwave —— RisingWave — Postgres 兼容的流式数据库,用物化视图替代 Flink + KV 组合
  • sidekiq —— Sidekiq — Ruby 后台任务的事实标准
  • temporal —— Temporal — 持久化工作流引擎
  • thrift —— Thrift — 写一份 IDL 自动生成 28 种语言的 RPC 代码
  • unstorage —— unstorage — 让 KV 存储不绑死运行时的统一抽象层
  • zeppelin —— Apache Zeppelin — JVM 多语言笔记本
  • zookeeper —— Apache ZooKeeper — 给一群机器装一个共同的小脑