Kafka 是一个分布式流处理平台,广泛应用于实时数据处理、日志收集、消息队列等场景。它的高吞吐量、低延迟和可扩展性使其成为现代云端数据架构中的重要组成部分。然而,许多关于 Kafka 的教程都会直接抛出复杂的架构图和 Kafka 的底层实现,如零拷贝、分布式日志、ISR 等等,这些内容对于初学者来说可能过于复杂,容易让人迷失在 Kafka 的内部机制中。我写过不同语言版本的 Kafka 协议实现和客户端实现,所以我将尝试客户端和 Kafka 协议的角度来解释 Kafka 的基本概念和工作原理。这样对使用 Kafka 的开发者来说,理解 Kafka 的工作原理会更加直观和容易。
Kafka 的基本概念
Kafka 通常作为一个消息队列系统使用,但它与传统的消息队列有一些不同。Kafka 服务器节点称为 Broker,一个 Kafka 集群由多个 Broker 组成,它的核心概念包括:
- Topic:消息的类别或主题,生产者将消息发送到特定的 Topic,消费者从特定的 Topic 中消费消息。
- Partition:每个 Topic 可以被划分为多个 Partition,每个 Partition 是一个有序的、不可变的消息序列。Partition 提供了并行处理的能力。
- Offset:每条消息在 Partition 中的唯一标识,消费者通过 Offset
- Producer:生产者负责将消息发送到 Kafka 的 Topic 中。生产者可以选择将消息发送到特定的 Partition,也可以让 Kafka 根据某种策略(如轮询或基于 Key 的哈希)自动选择 Partition。
- Consumer:消费者负责从 Kafka 的 Topic 中拉取消息进行处理。消费者可以以组的形式工作,Kafka 会确保同一组内的消费者不会重复消费同
先简单介绍一下以上这些基本概念,方便接下来从客户端的角度理解 Kafka 的工作原理。
首先,Kafka 的 Topic 是逻辑上的概念,它将消息按照主题进行分类。你可以将 Topic 理解成一个分布式的消息队列。由于 Kafka 的高吞吐量和可扩展性,每个 Topic 可以被划分为多个 Partition 会被分布在不同的 Kafka broker 节点上,这样可以实现消息的并行发送和消费。既然是消息队列,那么消息就应该是有序的,但 Kafka 把一个 Topic 分成多个 Partition,不同于传统的消息队列,Kafka 并不保证整个 Topic 内的消息是有序的。它只保证每个 Partition 内的消息是有序的,但不同 Partition 之间的消息顺序是不保证的。你可以理解为,一个 Topic 是一个逻辑上的消息队列,而它的每个 Partition 是一个物理上的消息队列。
它与传统的消息队列的另一个不同点是,Kafka 的消息是不可变的。也就是说,一旦消息被写入到 Partition 中,它就不能被修改或删除。传统的消息队列通常会在消息被消费后将其删除,而 Kafka 则是通过设置消息的保留策略来决定消息的存储时间。Kafka 的这种设计使得它可以支持高吞吐量和持久化存储,同时也为消费者提供了灵活的消费方式。消费者可以通过 Offset 偏移量来定位和消费消息。Offset 是每条消息在 Partition 中的唯一标识,它是一个递增的整数值。消费者可以选择从特定的 Offset 开始消费消息。
Kafka 协议的核心交互机制
Kafka 的客户端有两种角色,生产者和消费者。其中生产者向 Kafka 的 Topic 发送消息,消费者从 Kafka 的 Topic 中读取消息。客户端与服务器之间的交互是通过 Kafka 协议在 TCP 连接上进行的。Kafka 协议是一个二进制协议,它定义了客户端和服务器之间的请求和响应格式。客户端通过发送请求来与 Kafka 服务器进行通信,服务器则返回相应的响应。
在深入具体命令前,我们需要了解 Kafka 客户端与 Broker 通信的两大基本规则:
- Request / Response 模式:通信由客户端(Producer/Consumer)主动发起,Broker 被动响应
- 元数据驱动(Metadata Driven):客户端在发送消息或消费数据前,必须先知道「去哪里找谁」
由于连接是全双工的,客户端和 Broker 可以同时发送请求和响应,但在 Kafka 协议中,客户端始终是主动发起请求的一方,而 Broker 则是被动响应的一方。客户端发送请求时,需要指定请求的类型(如 Produce、Fetch 等)以及相关的参数(如 Topic、Partition、消息内容等)。Broker 接收到请求后,会根据请求类型进行处理,并返回相应的响应给客户端。客户端需要根据响应的内容进行相应的处理,如确认消息是否发送成功、获取消息内容等。
生产端是如何工作的
要启动一个客户端,无论是生产者还是消费者,首先都需要知道 Kafka 集群的 Broker 节点信息。因此,一般我们都需要给客户端提供一个或多个 Broker 的地址(IP + Port),以便客户端能够与 Kafka 集群建立连接。这就是 bootstrap.servers 配置项的作用,它告诉客户端 Kafka 集群的入口点。
客户端在启动时通过向 bootstrap.servers 里的 broker 地址列表中的任意一个 Broker 发送元数据请求(Metadata Request),以获取集群中所有的 Broker 节点信息,以及所有 Topic 和 Partition 的信息,包括每个 Partition 的 Leader 节点和副本节点的信息。客户端根据这些元数据信息来决定将消息发送到哪个 Partition,以及从哪个 Partition 拉取消息。然后,客户端可以根据实际需要,向 Broker 节点建立连接。具体的连接策略由客户端库的设计决定。
建立好 TCP 连接后,客户端就可以开始发送请求了。但如果 Broker 配置了 SASL 或者需要认证,客户端在发送请求前还需要进行身份验证。Kafka 协议支持多种认证机制,如 SASL/PLAIN、SASL/SCRAM、SASL/GSSAPI 等。客户端在发送请求时,需要在请求头中包含认证信息,以便 Broker 验证客户端的身份。
之后,客户端就可以开始发送 Produce 请求,将消息发送到指定的 Topic 和 Partition 中。如果观察 Produce 请求的结构,你可以发现,一个 Produce 请求中可以包含多条消息,甚至是不同的 Topic 和 Partition 的消息。再看具体的 RecordBatch (早期叫 MessageSet) 结构,你会发现它是一个批量的消息集合,包含了多条消息。Kafka 正是通过这种批量发送消息的方式,提高了 Producer 的效率和吞吐量。
要注意的是,Produce 请求中包含一个选项叫 required.acks,它决定了 Broker 在返回响应前需要等待多少个副本节点确认消息的写入。这个选项有三个可能的值:0, 1, -1(或 all)。当 required.acks 为 1 时,Broker 会等待 Leader 节点确认消息写入成功后才返回响应,这种方式在保证一定可靠性的同时,吞吐量也较高。当 required.acks 为 -1 时,Broker 会等待所有副本节点确认消息写入成功后才返回响应,这种方式的可靠性最高,但吞吐量最低。特别需要注意的是:当 required.acks 为 0 时,Broker 不会返回响应。这也是 Kafka 的请求 / 响应模式中的一个特殊没有响应的情况。这种方式的吞吐量最高,但消息可能会丢失,因为 Broker 不会对这个请求给出任何响应,即使是写入失败,Broker 也可能不会主动关闭连接,这导致客户端无法确认消息是否成功写入。
这个 Required Acks 是 Kafka Produce 请求中非常重要的一个选项,它直接影响到消息的可靠性和吞吐量。但在 brod (一个 Erlang 语言的 Kafka 客户端库)中提供了一个 produce_no_ack 的函数。它这个 no ack 的意思与 Kafka 协议中的 required.acks 的意思并不同。它的意思其实是 produce and forget,即发送消息后不等待响应,在 Erlang 里相当于一个异步的函数调用,类似于 gen_server:cast/2。通过调用 produce_no_ack 实际发送 Produce 请求的时候,brod 客户端仍然会在请求中设置 required.acks 为 brod client 设定的值,可能并不是 0。这种在概念上容易混淆的地方很容易成为隐藏的坑,开发者在使用时需要注意。
Kafka 官方客户端库的 Producer 有个配置 max.in.flight.requests.per.connection 默认值为 5,意味着客户端可以在同一个 TCP 连接上同时发送最多 5 个请求,而不需要等待前一个请求的响应。这种设计可以提高客户端的吞吐量和性能,但也需要注意请求的顺序和响应的处理。所以我们可以判断,这个配置是给 required.acks 不为 0 的情况设计的,因为如果 required.acks 为 0,客户端发送请求后不会等待响应,也没有 in flight 的概念。
注意到, max.in.flight.requests.per.connection 这个配置是在客户端的,而不是在 Broker 端的。也就是说,Kafka 将这个控制权交给了客户端,让客户端根据自己的需求来决定是否允许同时发送多个请求,出了问题也是客户端自己负责处理。实际上在 Kafka 的设计中,大多是这样的,Broker 端尽量保持简单,而将复杂的逻辑交给客户端去处理。
消费端是如何工作的
消费者与 Broker 的连接过程与生产者类似,这里就不画图了。真正复杂的地方在于消费者组成一个 Consumer Group 后,如何与 Broker 协调分配 Partition,以及如何管理 Offset。这些就不画图了。
消费端与 Broker 建立连接,获取 Metadata 数据,通过 SASL 认证后,发送的是 Fetch 请求来获取订阅的数据。Fetch 请求中包含了需要消费的 Topic 和 Partition,以及从哪个 Offset 开始消费。这个 Offset 可以是消费者上次消费的最后一个 Offset,也可以是消费者指定的一个新的 Offset。消费者可以选择从最新的 Offset 开始消费,也可以选择从最早的 Offset 开始消费,这取决于消费者的需求和配置。消费者可以选择自己保存这个 Offset 值,也可以让 Kafka 帮助管理 Offset。Kafka 提供了一个特殊的 Topic,叫做 __consumer_offsets,用于存储消费者的 Offset 信息。消费者可以通过 Commit Offset 请求将自己的 Offset 提交到这个 Topic 中,以便在下次消费时能够通过 Fetch Offset 请求取加对应 Partition 上次提交的 Offset 然后从上次消费的位置继续消费。
每个 Fetch 请求中可以包含多个 Topic 和 Partition 的信息,消费者可以一次性从多个 Partition 拉取消息。这也是 Kafka 高吞吐量的一个重要原因。消费者在拉取消息时,可以设置一个最大拉取的字节数(max_bytes)和最大拉取的消息数(max_messages),以控制每次 Fetch 请求返回的数据量。消费者在处理完拉取到的消息后,可以选择立即提交 Offset,也可以选择延迟提交 Offset,这取决于消费者的业务逻辑和需求。
此外,Kafka 还提供了 Consumer Group 的概念,允许多个消费者组成一个组,共同消费一个 Topic 的消息。Kafka 会确保同一组内的消费者不会消费同一个 Partition,从而避免重复消费同一条消息,从而实现负载均衡和高可用性。
为了实现这个功能,Kafka 提供了一个协调者(Coordinator)机制,负责管理消费者组的成员和分配 Partition。消费者在加入消费者组时,会向协调者发送 Join Group 请求,协调者会根据当前组内的成员,选出一个 Group Leader 让其进行 Partition 分配,并返回给每个消费者一个分配结果。消费者根据分配结果,开始消费自己被分配的 Partition。当有新的消费者加入组或者有消费者离开组时,协调者会重新进行分配,确保每个 Partition 都有消费者在消费。这也会导致在消费者组内的消费者之间进行重新平衡(Rebalance),在这个过程中,消费者可能会暂时停止消费,直到重新分配完成。Kafka 提供的这个机制已经经过了充分的优化和测试,能够在大多数情况下保证消费者组的稳定性和高可用性,当然在一些极端情况下或者有特殊的需求时,开发者也可以选择自己实现消费者组的管理逻辑。
关于 Partition
Partition 是 Kafka 的核心概念之一,它是 Topic 的数据划分的基本单位,是一个有序的、不可变的消息序列。每个 Partition 都有一个唯一的标识符(Partition ID),并且每个 Partition 都有一个 Leader 节点和多个副本节点。
不像其它一些分布式的服务,每个节点都能处理请求,Kafka 的读写请求是由 Leader 节点负责处理的,如果将这种请求发送到副本或其它节点则会收到错误的响应信息。在 Leader 节点发生故障或其它原因导致不可用时,副本节点可以被选举为新的 Leader 节点,从而保证数据的高可用性和可靠性。这就意味着,客户端在发送请求时,需要知道每个 Partition 的 Leader 节点信息,以便将请求发送到正确的 Broker 节点。客户端可以通过 Metadata 请求获取这些信息。对于 Leader 发生变化的情况,客户端需要及时更新自己的元数据信息,以便将请求发送到新的 Leader 节点。
正如上文所说,一个 Topic 只是逻辑上的概念,用于业务层面划分消息,而 Partition 才是物理上的概念,用于 Kafka 内部的存储和分布式处理。一个 Topic 被分为多个 Partition 后,Kafka 可以将这些 Partition 分布在不同的 Broker 节点上,从而实现消息的并行处理和负载均衡。客户端在生产和消费消息时,也可以将数据发送到不同的 Partition,从而提高吞吐量和性能。客户端可以根据自己的需求选择将消息发送到特定的 Partition,也可以让根据某种策略(如轮询或基于 Key 的哈希)自动选择 Partition。这个选择都是在客户端进行的,Kafka 协议并没有规定客户端必须使用哪种策略。例如,在 Kafka 官方客户端库中就提供了多种分区策略,开发者可以根据自己的需求选择合适的策略。
但 Partition 的数量并不是越多越好,过多的 Partition 会增加 Broker 的负载和管理复杂度,也会影响消费者组的重新平衡时间,还会影响生产者端的吞吐。同时还要考虑到,一个消费者可以同时消费多个 Partition,但一个 Partition 只能被一个消费者消费。这就意味着,如果一个 Consumer Group 里的消费者数量大于 Partition 的数量,那么有些消费者就会处于空闲状态,无法消费消息。同时为了能够尽量平均分配 Partition 给消费者,Partition 的数量在设计时也需要考虑到消费者的数量和负载情况。我曾经见过有人将 Partition 的数量设置为 1999 个,这是一个质数,无法被整除,这会导致在消费者组内的消费者无法平均分配 Partition,从而影响消费的效率和性能。正确的做法是需要评估目标吞吐量和单个生产者和单个消费者的处理能力。
其它 API 请求类型
除了以上提到的 Produce、Fetch、Metadata、Commit Offset、Join Group 等请求类型,Kafka 协议还定义了许多其他的请求类型,如用于管理 Topic、Partition、Consumer Group,管理 Cluster 等的请求。以及后续版本新增加的各种功能的请求类型,包括与 Transaction 相关的请求类型,管理 KRaft 模式的请求类型等。
这些请求类型都是 Kafka 协议的一部分,客户端可以根据自己的需求选择使用。例如可以开发一个用于管理 Kafka 集群的工具,或者开发一个用于监控 Kafka 集群状态的工具,这些都需要使用到 Kafka 协议中的请求类型。
其实,早期的 Kafka 协议中定义的请求类型并不多,只有 17 种请求类型,但随着 Kafka 的发展和功能的增加,Kafka 协议中的请求类型也在不断增加。到目前为止,Kafka 协议中已经定义了超过 90 种 API 请求类型,每种请求类型都有其特定的用途和参数。有趣的是,Kafka 在协议设计上从 0.11 版本开始就提供了 API 版本号的概念,这意味着同一种请求类型可以有多个版本,每个版本可能会有不同的参数和行为。这让 Kafka Broker 可以向后兼容旧版本的客户端(最低 0.9),同时也允许客户端使用新版本的请求类型来利用新的功能。
完整的 API 请求可以查看 Kafka Protocol 文档中的 API Keys 部分。 具体的请求类型和协议细节可以参考 Kafka 官方文档 Kafka Protocol 部分。了解这些请求类型有助于深入理解 Kafka 的工作原理和客户端与 Broker 之间的交互机制。
错误处理
值得一提的是 Kafka 协议中定义了丰富的错误码(Error Code),用于表示请求处理过程中可能出现的各种错误情况。每个请求类型在响应中都会包含一个错误码字段,用于指示请求是否成功以及失败的原因。其中每个错误码都有一个表示是否可以重试的标志位,客户端可以根据这个标志位来决定是否需要重试请求。Kafka 协议中的错误码设计得非常详细和丰富,涵盖了各种可能的错误情况,如网络错误、认证错误、权限错误、请求参数错误、Broker 状态错误等。客户端库在设计的时候需要考虑到这些错误码,并根据不同的错误码进行相应的处理。例如,对于一些可以重试的错误码,客户端可以选择在库内部进行重试,而对于一些不可重试的错误码,客户端则需要抛出异常或者返回错误信息给调用方。调用方需要根据错误码进行相应的处理,如记录日志、重试请求、抛出异常等。
总结
从上文的介绍可以看出,Kafka 的设计理念是将复杂的逻辑尽量放在客户端去处理,包含错误处理,而 Broker 端尽量保持简单,但同时也提供了丰富的 API 请求类型供客户端使用其它非核心的功能。从协议的消息结构设计上看,Kafka 鼓励客户端使用批量的方式发送和接收消息,以提高吞吐量和性能。另外,协议的设计也充分考虑了向后兼容性和扩展性,使得客户端可以根据自己的需求选择使用不同版本的请求类型,从而实现更灵活的功能和更高的性能。
除了从 Kafka Broker 的存储和分布式角度理解 Kafka 的工作原理外,我们还可以从客户端的角度去理解,这不仅更加直观,而且也非常重要。即使这些信息通常用在库的设计和实现上,但对于使用者来说,了解这些基本的概念和工作原理,那些客户端的库就不再是黑箱了,而是可以更好地理解和使用它们,从而在实际应用中更好地利用 Kafka 的特性和优势。
希望这些内容能够帮助有需要的人,并在实际应用中更好地理解和使用 Kafka。