开篇词 为什么要学习 Kafka?

消息引擎系统 ABC
设计消息引擎的关键点:
- 序列化:CSV、XML、JSON、Protocol Buffer、Thrift。Kafka 默认使用纯二进制字节序列
- 传输模型:同时支持点对点模型和发布/订阅模型

设计消息引擎的关键点:

Kafka 是一个开源分布式事件流平台。最初由 LinkedIn 开发,现在是 Apache 顶级项目。


Kafka 源码四大模块:
Kafka 是 Apache 的开源项目。Kafka 既可以作为一个消息队列中间件,也可以作为一个分布式流处理平台。
Kafka 用于构建实时数据管道和流应用。它具有水平可伸缩性,容错性,快速快速性。
Kafka 的存储子系统是其性能与可靠性的基石。作为一个分布式消息中间件与流处理平台,Kafka 需要在每秒数百万级消息写入的高并发场景下,仍然保证消息持久化、低延迟读取以及水平扩展能力。
Kafka 的存储设计完全有别于传统 MQ 的「消息落库」思路:它将消息以**追加写(append-only)**的方式写入分区日志文件,并借助操作系统的 Page Cache、顺序磁盘 I/O、稀疏索引以及零拷贝(sendfile)等技术,使得「磁盘存储 + 高吞吐」成为可能。同时,通过日志段(Log Segment)切分、基于保留策略的过期清理和日志压实(Log Compaction),Kafka 能够在海量数据场景下长期稳定运行。
Kafka 是一个分布式的、可水平扩展的、基于发布/订阅模式的、支持容错的消息系统。
Kafka 使用 Zookeeper 来维护集群成员的信息。每个 Broker 都有一个唯一标识符,这个标识符可以在配置文件里指定,也可以自动生成。在 Broker 启动的时候,它通过创建临时节点把自己的 ID 注册到 Zookeeper。Kafka 组件订阅 Zookeeper 的 /broker/ids 路径,当有 Broker 加入集群或退出集群时,这些组件就可以获得通知。
可靠传输是任何消息中间件的核心命题。在分布式系统中,网络抖动、节点宕机、磁盘故障、进程崩溃是常态,如何在这样不可靠的环境下保证消息「不丢失、不重复、有序」地从一个系统可靠地传递到另一个系统,是 Kafka 设计与运维的核心目标。
Kafka 通过「分区多副本 + ISR 同步机制 + 生产者确认(acks) + 消费者位移提交」的组合,在不同层面提供了可调节的可靠性保证。它并不追求绝对的「Exactly-Once」(这是分布式系统的理论难题),而是提供清晰的语义边界:默认 At-Least-Once,配合幂等生产者(Idempotent Producer)和事务(Transaction)可实现 Exactly-Once。
消息引擎获取消息有两种模式:

不管是把 Kafka 作为消息队列系统、还是数据存储平台,总是需要一个可以向 Kafka 写入数据的生产者和一个可以从 Kafka 读取数据的消费者,或者是一个兼具两种角色的应用程序。
使用 Kafka 的场景很多,诉求也各有不同,主要有:是否允许丢失消息?是否接受重复消息?是否有严格的延迟和吞吐量要求?
不同的场景对于 Kafka 生产者 API 的使用和配置会有直接的影响。
Kafka Producer 发送的数据对象叫做 ProducerRecord ,它有 4 个关键参数:
Kafka 是 Apache 的开源项目。Kafka 既可以作为一个消息队列中间件,也可以作为一个分布式流处理平台。
Kafka 用于构建实时数据管道和流应用。它具有水平可伸缩性,容错性,快速快速性。
Kafka、消息队列、Topic、Partition、BrokerKafka、Producer、分区策略、消息发送、幂等生产者Kafka、Consumer、消费组、Offset、重平衡Kafka、集群管理、ZooKeeper、Controller 选举、副本机制Kafka、可靠传输、ACK 机制、ISR、Exactly-OnceKafka、存储机制、日志段、PageCache、零拷贝Kafka、流式处理、Kafka Streams、事件流、窗口计算Kafka、运维、集群监控、性能调优、故障排查数据流是无边界数据集的抽象表示。无边界意味着无限和持续增长。无边界数据集之所以是无限的,是因为随着时间的推移,新的记录会不断加入进来。