Kafka 高吞吐来自顺序追加、批处理和页缓存,不是“磁盘永远比内存快”。事务则把同一 Kafka 流程中的多条输出和消费位置原子提交;它不会自动让外部数据库写入 Exactly-Once。
Table of contents
Open Table of contents
日志怎样定位
每个分区由多个 segment(段)组成。日志文件保存记录,稀疏 offset 索引把目标 offset 定位到附近文件位置,时间索引支持按时间查找。旧段按保留时间或容量清理;日志压缩则按 key 保留较新值,两者解决不同问题。
分区数过多会增加文件、恢复和控制面成本。消息过大降低批处理效率并放大网络和内存压力。容量规划应同时计算写入字节率、保留期、副本因子、压缩率和磁盘水位。
Exactly-Once 的限定范围
事务生产者用稳定 transactional.id 初始化事务,写输出记录,并把输入消费位置一同提交。下游设置 isolation.level=read_committed 后只读取已提交数据。旧生产者实例会被 fencing,避免两个实例使用同一身份并发写。
这能覆盖“Kafka 输入 → 处理 → Kafka 输出”。若处理还写 MySQL,Kafka 无法单方面原子提交 MySQL 事务;可把 offset 与结果写入同一数据库,或使用 Outbox、幂等和对账。
Kafka Streams 与 RocksDB
Kafka Streams 把拓扑拆成与分区对应的任务。无状态 map/filter 不需本地状态;聚合、窗口和 join 使用状态库,常见持久实现是 RocksDB。Changelog Topic 记录状态变化,使任务迁移后可恢复。
状态库不是唯一真相:恢复能力取决于 changelog 的持久性、消费位置和拓扑版本。上线前测试重启恢复时间、迟到事件、窗口保留、状态目录磁盘不足和重新分区。官方 Kafka 设计文档也明确说明对外部系统的 Exactly-Once 需要外部系统配合。
下一步
继续阅读08-18 RocketMQ 概念、客户端与消息类型,比较面向业务消息的另一套模型。