译者按

  • 核心就一句:Kafka 不是靠什么神仙算法,就是“追加写日志 + 线性 IO”,只要记住这一点,后面所有调优和架构选型都有谱了。
  • 国内中小厂如果日吞吐还在 GB 级以下,直接上集群多半是给自己找事——那套分布式复杂度、7×24 值守和磁盘规划,先把运维干掉半条命;这种量级其实用托管 Kafka 或普通 MQ 更香。
  • 如果非要 K3s/K8s 私有云自建,心里得有数:网络和存储才是爸爸,PVC 的延迟和跨节点带宽一旦拉胯,Kafka 副本追不上、ISR 缩水都是白给的事故。
  • 分层存储那套 S3/MinIO 冷热分离,在中小环境别抱太大希望,私有云对象存储的带宽和成本根本给不到它发挥的前提,老老实实按“日写入量 × 保留天数”容量买磁盘更实在。
  • 落地建议先做减法:集群规划、分区数、监控(ISR、消费 lag)这几样砸实了,再考虑 Cruise Control、KRaft 这些高级特性,否则大厂最佳实践到你这儿全变成事故预案。

Apache Kafka 架构入门:从日志到生态

本文为技术翻译稿,内容基于 Apache Kafka Committer Stanislav Kozlovski 的文章精简优化。

Apache Kafka 是目前最流行的开源分布式流处理平台,已成为实时数据流事实标准。本文从底层日志结构出发,讲清 Kafka 的核心机制、性能优化、容错设计,以及周边生态组件。

核心:日志(Log)

Kafka 中数据存储在 Topic 中,而 Topic 的基础是 Log——一个简单的有序数据结构,按顺序追加记录。

Log 具有两个关键特性:

  • 不可变性:已写入的记录不会被修改
  • O(1) 读写:只要从尾部写入、从头部或尾部读取,访问速度不会随着日志增大而变慢

选择 Log 作为核心结构的根本原因,是它针对机械硬盘(HDD)做了优化。HDD 最擅长线性读写,而 Log 的操作模式恰好就是线性读写。这使 Kafka 能以极低成本存储大量数据,同时保持高性能。

性能优化

一个调优良好的 Kafka 集群,瓶颈通常在网络层面,吞吐可达每秒数 GB。性能来自多个层面的优化:

磁盘持久化

Kafka 实际上把所有数据写入磁盘,不显式维护内存缓存。它通过协议批量聚合消息,减少网络开销;服务端一次性持久化一批消息(线性写),消费者一次性拉取大块数据(线性读)。

线性读写避免了磁盘寻道,而操作系统还会进一步优化:

  • 预读(read-ahead):提前加载大块数据到内存,后续读取无需触盘
  • 后写(write-behind):小逻辑写合并为大物理写。Kafka 不使用 fsync,写操作异步落盘

Pagecache

现代操作系统会用空闲内存缓存磁盘内容,即页缓存(Pagecache)。Kafka 在整个链路(生产者 → Broker → 消费者)中保持消息的二进制格式不变,因此可以利用 零拷贝(zero-copy) 技术,让操作系统将数据从页缓存直接复制到 socket,绕过 Kafka 进程。

不过零拷贝的实际收益有限:一是优化良好的集群 CPU 很少成为瓶颈,二是生产环境必须开启 SSL/TLS,加密过程会修改消息内容,从而禁用零拷贝。

基础概念

Broker 与副本

Kafka 分布式节点称为 Broker。每个 Topic 分为多个 分区(Partition),每个分区按复制因子保存 N 个 副本(Replica) 以实现高可用。每个副本就是一组日志文件,记录按 偏移量(Offset) 单调递增标记。

每个分区有且仅有一个 Leader 副本,负责处理读写。其余副本为 Follower,状态分为 in-sync(同步中)和 out-of-sync(已落后)。

写入

写入请求只能发往 Leader。客户端生产者(Producer)通过 acks 参数控制持久化级别:

  • acks=0:不等待任何确认,立即视为成功
  • acks=1:Leader 落盘后返回成功
  • acks=all(默认):所有 in-sync 副本都落盘后才返回成功

为避免只有一个 in-sync 副本时退化为 acks=1,可用 min.insync.replicas 指定最少同步副本数。

读取

消费者(Consumer)可以从任意副本读取,通常选择网络拓扑最近的副本。多个 Consumer 构成 消费者组(Consumer Group),通过 Broker 协调彼此,组内每个分区同一时刻只能被一个消费者读取,保证分区内有序消费。

消费者组的消费进度(offset)持久化在名为 __consumer_offsets 的 Topic 中,该 Topic 分区的 Leader Broker 充当 Group Coordinator,负责管理组成员和存活状态。

Kafka 胜过传统消息队列的核心原因:消息不会在消费后被删除。生产者和消费者完全解耦,消费者慢不会阻塞生产者。

容错与共识

Controller

每个 Kafka 集群有一个活跃 Controller,负责管理所有需要全局唯一决策的元数据操作,如创建/删除 Topic、分区副本分配等。最重要的是,它负责每个分区的 Leader 选举。

从 ZooKeeper 到 KRaft

早期 Kafka 依赖 ZooKeeper 进行 Controller 选举和元数据存储。Broker 启动时竞争注册 /controller zNode,先到者成为 Controller。

最近几年 Kafka 逐步脱离 ZooKeeper,转向自研共识协议 KRaft(Kafka Raft)。KRaft 是 Raft 的变体,将集群元数据表达为日志中的有序事件流。集群中由 N 个控制器(通常 3 个)组成仲裁,通过 Raft 选举出活跃 Controller。所有 Broker 都异步复制 __cluster_metadata 主题来更新本地元数据。

KRaft 支持组合模式和隔离模式部署。生产级支持从 Kafka 3.3 开始,ZooKeeper 预计在 Kafka 4.0 完全移除。

分层存储(Tiered Storage)

传统架构中,Broker 将数据全部保存在本地磁盘。当单 Broker 存储近 10TB 时,会带来几个严重问题:

  1. 异常恢复慢:非正常关闭后需要重建本地索引文件,10TB 磁盘可能需要数小时甚至数天
  2. 历史读取消耗 IOPS:HDD 的 IOPS 通常只有 120 左右,历史数据读取会让消费端与生产端争抢 IOPS,性能急剧下降
  3. 故障复制放大:磁盘故障后需要从零复制 10TB 数据,期间产生大量跨 Broker 历史读
  4. 分区重分配成本高:副本迁移需要全量复制数据

分层存储通过将大部分数据存放到远程对象存储(如 S3/MinIO)来解决上述问题。Broker 的热数据留在本地,冷数据自动分层到对象存储;Leader 负责分层,Follower 也可直接从对象存储读取历史数据。该特性已进入 Early Access,测试显示在存在历史消费者时,生产者性能提升 43%。

辅助生态

分区重平衡与 Cruise Control

Kafka 集群运行一段时间后,容易出现热点或不均衡。为此社区开发了 Cruise Control(LinkedIn 开源),它从 Kafka Topic 读取所有 Broker 指标,在内存中构建集群模型,通过启发式装箱算法计算优化方案,再利用 Kafka 底层 API 增量执行分区重分配。

Cruise Control 支持多种可配置的 Goal(目标),持续监控指标并自动触发重平衡,同时提供 API 简化 Broker 的增删操作。

Kafka Connect

Kafka Connect 是 Apache 开源项目的一部分,用于将 Kafka 与其他系统进行插件化集成。它分为:

  • Source Connector:从外部系统导入数据到 Kafka
  • Sink Connector:从 Kafka 导出数据到外部系统

Connect 运行时支持单机模式和分布式模式。分布式模式下,Worker 节点利用 Kafka 内部 Topic 存储配置、状态和偏移量,同时复用消费者组协议处理故障和任务分配。社区提供了大量成熟 Connector(如 Elasticsearch、PostgreSQL、MySQL 等),用户只需 REST API 配置即可。

Kafka Streams

Kafka Streams 是 Apache 开源客户端库,提供实时流处理 API(连接流、转换、丰富数据)。它不是运行在 Broker 上,而是作为普通 Kafka 客户端嵌入你的应用中,无需单独部署集群。它支持精确一次(exactly-once)处理语义。

现状与展望

Kafka 虽然诞生于 2011 年,但社区依然活跃。目前趋势是在 Kafka API 之上竞争底层实现:

  • Confluent Kora:云原生 Kafka 引擎
  • RedPanda:C++ 重写 Kafka
  • WarpStream:重度依赖 S3,避免复制和 Broker 有状态

总之,Kafka 是成熟、广泛采用的流数据平台,开源社区健康,历经十余年依然保持强劲创新力。

原文来源:High Scalability