05Kafka

约 9 分钟 · 2467 字评论
请注意,本文编写于 830 天前,最后修改于 39 天前,其中某些信息可能已经过时。

Kafka 复习

一、采用 Kafka 作为消息队列的原因

  1. 兼顾实时与离线分析:出于后期打造实时分析平台和现在的离线数据分析需求,需要兼顾两者,故采用。
  2. 削峰填谷与解耦:数据采集速率和数据处理速率不同,Kafka 可动态提高消费能力,避免数据积压。
  3. 统一数据源:应用间解耦,当后期出现不同的应用需要原始数据时,Kafka 可以作为统一数据源。

二、Kafka 的基础架构分析

1. 生产者发送数据流程

生产者发送数据流程

  1. 生产者发送数据主要包含两个线程:main 线程sender 线程。main 线程中,消息通过封装后,调用 sender 方法发送给 broker。
  2. main 线程中,消息依次通过拦截器、序列化器、分区器后进入缓冲,缓冲大小默认 32M。每个分区的数据进入对应的双端队列,通过 batch.sizelinger.ms 参数控制数据发送给 sender 的时机。
  3. sender 线程中,会创建网络客户端与 broker 建立联系。默认每个 broker 只能缓存 5 个请求,向 broker 发送数据后,根据配置参数 acks,broker 会汇报数据的落盘情况。

2. Broker 的初始化流程

Broker初始化流程

  1. broker 初始化时,如果使用 ZooKeeper 作为注册中心,先向 ZooKeeper 注册的节点会当选为集群的 controller。
  2. controller 负责副本 leader 的选举、broker 的上下线等。当数据的 leader 出现故障,controller 根据 ZooKeeper 的选举机制,负责选举新的 broker 作为副本 leader,同时更新 ZooKeeper 节点的信息,其他的存活 follower 从 ZooKeeper 读取最新的动态。
  3. AR 队列(所有副本的总称)
    • ISR(In-Sync Replicas):同步的副本节点,包括 leader 节点。
    • OSR(Out-of-Sync Replicas):与 leader 失去联系的副本节点,不包括 leader 节点。

3. 消费者初始化流程

消费者初始化流程

  1. 消费者的初始化流程主要包括消费者组中消费计划的制定。
  2. 消费者建立时,会选择其中的一个消费者作为 leader,制定分区消费策略并汇报给 coordinator,再由 coordinator 分发计划给其他的消费者。

4. 消费者的消费流程

消费者的消费流程

消费时,每个消费者按照分配的消费策略,通过网络客户端对指定分区的数据发送请求,并将从 broker 抓取的数据放在队列中,再从队列中拉取。

Q: 如何保证消费的顺序?
Kafka 只能保证单分区内的数据有序性,具体顺序需求需要结合场景自定义实现。


三、生产者的分区器分类

  1. 默认的分区分配策略:Sticky Partitioner(黏性分区器)

    • 是否指定分区?指定了,直接发往指定的分区。
    • 消息是否拥有 key?拥有,对 key 算法取模,决定发往的分区。
    • 否则使用黏性分区。

    Sticky Partitioner

  2. 纯粹的黏性分区:Uniform Sticky Partitioner

    • 不管消息是否有 key,直接使用黏性分区策略。
  3. Round Robin 轮询分区策略

    • 消息对 topic 的分区进行轮询分配。

四、消费者的分区分配策略

默认策略Range + CooperativeSticky

  1. Range:对每个 topic,按 分区数量 / 消费者数 决定每个消费者将要消费的平均分区数。当消费的 topic 数量较多时,可能出现数据倾斜。
  2. Round Robin:所有的 topic 的分区和 consumer 进行排序(按 hashcode),轮询分配,但没有考虑到 rebalance 发生时带来的重分区开销。
  3. 黏性消费(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×1k1M/s
  • 高峰时段(晚上下班 7-12 点左右):
    • 高峰条数 = 平均 20 倍 ≈23000条/s\approx 23000条/s23000/s(大概 2 万条/s 左右)
    • 高峰速率 ≈20M/s\approx 20M/s20M/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=301003 个分区

八、Kafka 的数据一致性问题

1. 数据是否重复

  • 生产者开启幂等性:保证单分区单会话不重复。
    • 标识结构:<PID, partition, SequenceNum>
      • PID:Kafka 给生产者分配的 Producer ID
      • partition:发往的分区号
      • SequenceNum:生产者给数据打上的自增序列号
  • 消费者重复消费引起的重复
    • 可能场景
      1. 消费者消费了数据,但是并未提交 offset。
      2. 消费者使用自动提交 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 查看。 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 的读写完成整个生产-消费过程。 Page Cache

2. 写数据(主要针对生产者)

  • 顺序写磁盘,使得 Kafka 使用普通机械硬盘和固态硬盘写数据的速率相近。

3. 读数据(主要针对消费者)

  • 每个分区内的数据逻辑分片并在区内建立稀疏索引,可以快速定位消费的数据。
    • 稀疏索引每隔 4KB 数据建立一个索引:(offset, 文件中的绝对位置)
    • 此外 .timeindex 文件也采用了稀疏索引 (timestamp, offset),因此 Kafka 可以根据具体的 offset 或时间戳指定消费。
  • 零拷贝技术(Zero-Copy):不用在读取数据时重复在内核态与用户态应用之间拷贝数据。
    • 普通的读数据流程: 普通的读数据流程
    • Kafka 零拷贝流程: 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 条。

版权声明

本文作者:hedeoer

本文链接:/post/__gh__6

许可说明:本博客所有文章除特别声明外,均采用 许可协议。转载请注明出处!

05Kafka

0

匿名访客仅支持文本评论。

暂无评论