05Kafka
||约 9 分钟 · 2467 字|评论
请注意,本文编写于 830 天前,最后修改于 39 天前,其中某些信息可能已经过时。
Kafka 复习
一、采用 Kafka 作为消息队列的原因
- 兼顾实时与离线分析:出于后期打造实时分析平台和现在的离线数据分析需求,需要兼顾两者,故采用。
- 削峰填谷与解耦:数据采集速率和数据处理速率不同,Kafka 可动态提高消费能力,避免数据积压。
- 统一数据源:应用间解耦,当后期出现不同的应用需要原始数据时,Kafka 可以作为统一数据源。
二、Kafka 的基础架构分析
1. 生产者发送数据流程
- 生产者发送数据主要包含两个线程:main 线程和 sender 线程。main 线程中,消息通过封装后,调用
sender方法发送给 broker。 - main 线程中,消息依次通过拦截器、序列化器、分区器后进入缓冲,缓冲大小默认 32M。每个分区的数据进入对应的双端队列,通过
batch.size和linger.ms参数控制数据发送给 sender 的时机。 - sender 线程中,会创建网络客户端与 broker 建立联系。默认每个 broker 只能缓存 5 个请求,向 broker 发送数据后,根据配置参数
acks,broker 会汇报数据的落盘情况。
2. Broker 的初始化流程
- broker 初始化时,如果使用 ZooKeeper 作为注册中心,先向 ZooKeeper 注册的节点会当选为集群的 controller。
- controller 负责副本 leader 的选举、broker 的上下线等。当数据的 leader 出现故障,controller 根据 ZooKeeper 的选举机制,负责选举新的 broker 作为副本 leader,同时更新 ZooKeeper 节点的信息,其他的存活 follower 从 ZooKeeper 读取最新的动态。
- AR 队列(所有副本的总称):
- ISR(In-Sync Replicas):同步的副本节点,包括 leader 节点。
- OSR(Out-of-Sync Replicas):与 leader 失去联系的副本节点,不包括 leader 节点。
3. 消费者初始化流程
- 消费者的初始化流程主要包括消费者组中消费计划的制定。
- 消费者建立时,会选择其中的一个消费者作为 leader,制定分区消费策略并汇报给 coordinator,再由 coordinator 分发计划给其他的消费者。
4. 消费者的消费流程
消费时,每个消费者按照分配的消费策略,通过网络客户端对指定分区的数据发送请求,并将从 broker 抓取的数据放在队列中,再从队列中拉取。
Q: 如何保证消费的顺序?
Kafka 只能保证单分区内的数据有序性,具体顺序需求需要结合场景自定义实现。
三、生产者的分区器分类
-
默认的分区分配策略:Sticky Partitioner(黏性分区器)
- 是否指定分区?指定了,直接发往指定的分区。
- 消息是否拥有 key?拥有,对 key 算法取模,决定发往的分区。
- 否则使用黏性分区。
-
纯粹的黏性分区:Uniform Sticky Partitioner
- 不管消息是否有 key,直接使用黏性分区策略。
-
Round Robin 轮询分区策略
- 消息对 topic 的分区进行轮询分配。
四、消费者的分区分配策略
默认策略:Range + CooperativeSticky
- Range:对每个 topic,按
分区数量 / 消费者数决定每个消费者将要消费的平均分区数。当消费的 topic 数量较多时,可能出现数据倾斜。 - Round Robin:所有的 topic 的分区和 consumer 进行排序(按 hashcode),轮询分配,但没有考虑到 rebalance 发生时带来的重分区开销。
- 黏性消费(Sticky):和 Range 策略大致相同,但当消费组内出现了 rebalance,尽量遵守两个规则:
- 规则一:分区分配尽量均衡;
- 规则二:尽量复用 rebalance 前的消费策略。
- 如果两个规则冲突,以第一条规则为准。
五、Kafka 集群规模估算
从吞吐的要求来看,考虑磁盘的吞吐:
- 写入吞吐/s = 写入速率 * 副本数
- 读取吞吐/s = 读取速率 * (副本数 - 1 + 消费者组数量)
示例计算:
- 实际的数据速率高峰 =
20M/s,期望能支撑100M/s的写入速率 - 写吞吐 =
200M/s - 读吞吐 =
200M/s - 一般单台磁盘的 I/O 速率
100+M/s→\rightarrow→ 至少 2 台,建议配置 3 台。
监控工具:
- Kafka Eagle
六、数据量的估算
- 场景:100 万日活,每人每天平均产生 100 条数据,每条数据 0.5-2k,平均 1k。
- 每天条数 = 100万×100条=1亿条100万 \times 100条 = 1亿条100万×100条=1亿条
- 每天数据大小 = 1亿条×1k=100G1亿条 \times 1k = 100G1亿条×1k=100G(每 1 万日活对应 1G 数据,对应 100 万条)
- 平均条数 = 1亿条/(24×60×60)≈1150条/s1亿条 / (24 \times 60 \times 60) \approx 1150条/s1亿条/(24×60×60)≈1150条/s(估算值,大概 1000 多条/s)
- 平均速率 = 1150条/s×1k≈1M/s1150条/s \times 1k \approx 1M/s1150条/s×1k≈1M/s
- 高峰时段(晚上下班 7-12 点左右):
- 高峰条数 = 平均 20 倍 ≈23000条/s\approx 23000条/s≈23000条/s(大概 2 万条/s 左右)
- 高峰速率 ≈20M/s\approx 20M/s≈20M/s
- 低谷时段(半夜 3-4 点):
- 低谷条数 = 几十条/s
- 低谷速率 = 几十 k/s
七、如何预估 Topic 的分区数
压测:使用官方压测脚本
- 单分区生产峰值速率 Tp=30T_p = 30Tp=30 多 M/s
- 单分区消费峰值速率 Tc=50T_c = 50Tc=50 M/s
- 期望支撑的吞吐高峰 Tt=100T_t = 100Tt=100 M/s
- 估算的分区数 = Ttmin(Tp,Tc)=10030≈3\frac{T_t}{\min(T_p, T_c)} = \frac{100}{30} \approx 3min(Tp,Tc)Tt=30100≈3 个分区
八、Kafka 的数据一致性问题
1. 数据是否重复
- 生产者开启
幂等性:保证单分区单会话不重复。- 标识结构:
<PID, partition, SequenceNum>PID:Kafka 给生产者分配的 Producer IDpartition:发往的分区号SequenceNum:生产者给数据打上的自增序列号
- 标识结构:
- 消费者重复消费引起的重复
- 可能场景:
- 消费者消费了数据,但是并未提交 offset。
- 消费者使用自动提交 offset,但在尚未提交时发生了 rebalance,新的消费者重新消费。
- 解决办法:
- 下游数据消费支持幂等写入
- 下游数据消费使用事务写入
- 可能场景:
2. 数据是否丢失(Producer 写入时丢失)
- Producer 端:
acks参数设置为-1
0:不用等待数据落盘,只管不断发送数据1:等待 leader 响应落盘-1:等待 leader 和 follower 都落盘后响应
- Broker 端:
- 数据副本和最小同步副本数(
min.insync.replicas)都至少设置为 2。 - 数据副本为 2 保证降低数据丢失风险,最小同步副本数为 2 保证满足条件才算数据同步完成。
- 数据副本和最小同步副本数(
九、Kafka 的数据有序性保证
1. 单分区有序
- 开启幂等性:
新来数据的
SeqNum= broker 保存的最大SeqNum + 1才接收;如果> 1则打回重新排序,按照一个并发发送。 - 并发设为 1:
未开启幂等性时,将
max.in.flight.requests.per.connection设置为 1。
2. 多分区无序
- 由架构本身决定。
- 解决方案:
- 方案一:使用单分区的 topic。
- 方案二:针对特定场景保证同一张表的数据有序。指定
key = 库名 + 表名(项目中使用的方式),通过配置 Maxwell 的producer_partition_by=table,保证一张表的数据进入同一个分区。 - 方案三:数据消费之后,再进行攒批全局排序。
十、Kafka 如何解决数据积压
1. 如何发现数据积压?
- 通过 Kafka 的监控工具 Kafka Eagle 查看。
2. 数据积压的危害
- 数据的时效性变低,解决积压问题后,再次消费数据会导致滞后消费,降低时效性。
- 数据积压时间可能超过 Kafka 设置的数据清理最大时长,从而导致数据丢失。
3. 解决办法
- 提高消费者的消费能力:
- 同比增加消费者数和 topic 的分区数量,提高消费吞吐。
- 提高单个消费者的消费能力,例如增大每次从 broker 抓取数据的大小
fetch.min.bytes(从默认的 100M 提高到 200M)。 - 提高每次从队列中拉取数据的条数,从默认的 500 条提高到 1000 条(
max.poll.records)。
- 提高消费者的处理能力:
- 例如当 Kafka 的下游消费为 Flink 时,可以针对 Flink 的反压(Backpressure)机制进行优化。
十一、Kafka 高效的原因
1. 读写操作
- Kafka 本身分布式,采用数据分区写入。
- 使用内核的页缓存技术(Page Cache)(读和写),减少磁盘 I/O。如果 Kafka producer 的生产速率与 consumer 的消费速率相差不大,就能几乎只靠对 broker page cache 的读写完成整个生产-消费过程。
2. 写数据(主要针对生产者)
- 顺序写磁盘,使得 Kafka 使用普通机械硬盘和固态硬盘写数据的速率相近。
3. 读数据(主要针对消费者)
- 每个分区内的数据逻辑分片并在区内建立稀疏索引,可以快速定位消费的数据。
- 稀疏索引每隔 4KB 数据建立一个索引:
(offset, 文件中的绝对位置)。 - 此外
.timeindex文件也采用了稀疏索引(timestamp, offset),因此 Kafka 可以根据具体的 offset 或时间戳指定消费。
- 稀疏索引每隔 4KB 数据建立一个索引:
- 零拷贝技术(Zero-Copy):不用在读取数据时重复在内核态与用户态应用之间拷贝数据。
- 普通的读数据流程:
- Kafka 零拷贝流程:
- 普通的读数据流程:
十二、Kafka 的优化(提升吞吐量)
1. Producer 端
batch.size:缓冲区一批数据的最大值,默认 16k。适当增加该值可以提高吞吐量,但如果设置过大,会导致数据传输延迟增加。linger.ms:如果数据迟迟未达到batch.size,sender 等待linger.ms之后就会发送数据。buffer.memory:RecordAccumulator 缓冲区总大小,默认 32m。compress.type:生产者发送的所有数据的压缩方式。默认是none,不压缩。
2. Broker 端
- 增加分区数量。
3. Consumer 端
fetch.min.bytes:默认 1 个字节。消费者获取服务器端一批消息最小的字节数。fetch.max.bytes:消费者获取服务器端一批消息最大的字节数。max.poll.records:一次 poll 拉取数据返回消息的最大条数,默认是 500 条。
版权声明