Kafka底层布局
1. Kafka 架构总览Kafka 是一个分布式消息队列,接纳**发布-订阅(Pub-Sub)**模式,焦点组件包罗:
[*]Producer(生产者): 负责向 Kafka 发送消息。
[*]Broker(Kafka 服务器): 负责存储和管理消息。
[*]Topic(主题): 消息的分类单位。
[*]Partition(分区): Topic 的物理分片,进步吞吐量。
[*]Consumer(消耗者): 订阅并消耗消息。
[*]Consumer Group(消耗者组): 消耗者的逻辑分组,支持并行消耗。
[*]Zookeeper: 负责 Kafka 集群的元数据管理、Leader 推选等。
2. 底层存储布局
Kafka 接纳次序写入日志文件的方式存储数据,底层存储接纳**Segment(日志分段)+ Index(索引)**的方式管理数据。
2.1 日志存储
[*] 每个 Partition 对应一个日志目次,目次布局如下:
/kafka-logs/
├── topic-1/
│ ├── 0/# 分区0
│ │ ├── 00000000000000000000.log# 日志文件
│ │ ├── 00000000000000000000.index# 索引文件
│ │ ├── 00000000000000000000.timeindex# 时间索引文件
│ │ ├── leader-epoch-checkpoint# 领导者任期记录
[*] 日志分段(Segment)
[*]Kafka 不会将全部消息存入一个文件,而是拆分成多个段文件(Segment),每个 Segment 都是一个固定巨细(默认1GB)的日志文件。
[*]新消息总是追加到当前生动段(Active Segment),当文件到达肯定巨细后,Kafka 会新建一个段文件。
[*] 索引机制
[*]索引文件(.index): 记录消息在日志文件中的偏移量和物理位置。
[*]时间索引(.timeindex): 通过期间戳查找近来的消息,进步查询服从。
2.2 日志清算(Log Retention & Compaction)
Kafka 提供两种清算战略:
[*]日志生存(Retention): Kafka 按照**时间(log.retention.hours)或巨细(log.retention.bytes)**删除旧数据,默认存储 7 天。
[*]日志压缩(Log Compaction): 仅生存最新的 Key-Value 记录,实用于 幂等性数据存储 场景。
3. 生产者消息投递
生产者(Producer)负责将消息发送到 Kafka,Kafka 接纳以下机制包管消息可靠性:
[*] 分区战略(Partitioning)
[*]轮询(Round-Robin): 生产者将消息匀称分配到差别的分区。
[*]按 Key 选择(Keyed Partitioning): 生产者根据 Key 盘算 Hash 值,映射到固定分区,包管雷同 Key 的消息进入同一个分区。
[*]自界说战略(Custom Partitioning): 用户可以自界说分区规则。
[*] 消息确认机制(Acknowledgment)
[*]acks=0:不期待确认,大概丢失数据。
[*]acks=1:只需 Leader 记录消息,大概丢失数据(Leader 瓦解)。
[*]acks=all:全部副本都写入后才确认,包管最高可靠性。
[*] 批量发送(Batching)
[*]Kafka 生产者默认支持批量发送(Batch),进步吞吐量。
[*]通过参数 batch.size 控制批量巨细。
[*] 压缩(Compression)
[*]Kafka 支持 GZIP、Snappy、LZ4、Zstd 压缩方式,淘汰带宽占用。
4. 消耗者消耗机制
消耗者从 Kafka 拉取数据,接纳 Consumer Group(消耗者组) 机制包管数据分发:
[*]每个分区只能被一个组内的消耗者消耗,包管同一条消息不会被组内多个消耗者重复消耗。
[*]差别的 Consumer Group 可以并行消耗同一 Topic,进步并发本领。
4.1 消息拉取方式
Kafka 接纳Pull(拉取)模式,而非传统的Push(推送)模式:
[*]Push 模式:生产者自动推送数据,大概导致消耗者过载。
[*]Pull 模式:消耗者自主决定拉取频率,克制过载标题,进步吞吐量。
4.2 消耗者偏移量(Offset)管理
Kafka 利用消耗者位移(Offset) 记录消耗进度:
[*]自动提交(enable.auto.commit=true): 消耗者定期提交偏移量,大概丢失数据。
[*]手动提交: 通过 commitSync() 或 commitAsync() 提交偏移量,包管消耗的可靠性。
4.3 Rebalance 机制
当消耗者加入/退出消耗者组,Kafka 会举行重新分配分区(Rebalance):
[*]Rebalance 触发条件
[*]新消耗者加入
[*]消耗者故障
[*]分区数厘革
5. 分区副本(Replication)机制
Kafka 接纳副本机制(Replication) 包管数据高可用:
[*] 每个分区都有多个副本(Replica),此中:
[*]Leader 副本 负责读写数据。
[*]Follower 副本 仅做同步,供故障转移利用。
[*] ISR(In-Sync Replicas)同步机制
[*]Kafka 维护同步副本聚集(ISR),存储最新的同步副本。
[*]仅 ISR 内的副本能当选 Leader,包管数据同等性。
[*] 副本推选(Leader Election)
[*]当 Leader 瓦解,Kafka 会自动推选新的 Leader,包管服务可用。
6. 高吞吐计划
Kafka 接纳多种优化战略,进步吞吐本领:
[*]零拷贝(Zero-Copy)
[*]接纳 sendfile 体系调用,克制数据在用户态和内核态之间拷贝,进步性能。
[*]次序写入
[*]Kafka 接纳次序写入磁盘,淘汰随机 IO,提升写入速率。
[*]批量处置处罚
[*]生产者批量发送消息,淘汰网络开销,进步吞吐量。
7. Zookeeper 在 Kafka 中的作用
Kafka 依赖 Zookeeper 举行集群管理,紧张包罗:
[*]存储元数据
[*]记录 Topic、分区、副本等信息。
[*]推选 Kafka Controller
[*]控制分区的 Leader 推选,维护集群状态。
[*]消耗者 Rebalance
[*]调和 Consumer Group,触发 Rebalance。
同等性计划
1. 生产者同等性包管
生产者同等性紧张涉及数据是否乐成写入 Kafka,而且不会丢失或重复,Kafka 提供以下机制来包管生产者同等性:
1.1 ACK 确认机制
Kafka 生产者在发送消息时,依赖 acks 参数来确认数据是否乐成写入 Kafka:
[*]acks=0:不期待确认,最快,但大概会丢失数据(差别等)。
[*]acks=1:只期待Leader 副本确认,存在 Leader 瓦解导致数据丢失的风险。
[*]acks=all(或 acks=-1):期待全部 ISR 副本确认,确保数据不会丢失,但写入延伸较高。
✅ 最佳实践:
[*]对于高同等性要求,发起利用 acks=all。
[*]可联合 min.insync.replicas 设置,确保至少有 N 个副本 乐成写入后才确认。
1.2 生产者重试机制
Kafka 生产者大概因网络标题、Broker 宕机等缘故起因发送失败。Kafka 通过重试机制进步数据同等性:
[*]retries=N:指定重试次数。
[*]retry.backoff.ms:两次重试之间的时间隔断。
⚠️ 注意:
[*]若 retries > 0,但 max.in.flight.requests.per.connection > 1,大概导致消息乱序。
[*]办理方案:
[*]包管消息次序:设置 max.in.flight.requests.per.connection=1。
✅ 最佳实践:
[*]对于幂等性包管,需共同 enable.idempotence=true(见下一节)。
[*]retries 设为较大值(如 retries=5),克制短期故障导致数据丢失。
1.3 幂等性(Idempotency)
Kafka 生产者默认环境下大概会在重试过程中导致重复消息,可以启用幂等性包管数据同等性:
[*]enable.idempotence=true:Kafka 生产者端启用幂等性,确保同一条消息只写入一次,纵然发生重试。
Kafka 通过Producer ID(PID)+ Sequence Number 组合,确保雷同 Producer 发送的消息不会被重复写入。
✅ 最佳实践:
[*]猛烈发起在高同等性场景下启用幂等性 enable.idempotence=true。
[*]acks=all + enable.idempotence=true 可实现**"Exactly Once"(精准一次)** 语义。
1.4 变乱包管(Exactly-Once)
Kafka 生产者支持变乱(Transactional),确保跨分区或跨批次的消息要么全部乐成,要么全部失败。
启用变乱时:
[*]生产者调用 initTransactions() 初始化变乱。
[*]生产者调用 beginTransaction() 开始变乱。
[*]生产者发送消息。
[*]生产者调用 commitTransaction() 提交变乱,或 abortTransaction() 回滚变乱。
✅ 最佳实践:
[*]变乱实用于涉及多个 Topic 或多个分区的消息处置处罚场景,如金融体系、订单体系。
[*]变乱模式下,必须启用 acks=all 和 enable.idempotence=true。
2. 消耗者同等性包管
Kafka 消耗者同等性紧张涉及:
[*]消息不丢失(At-Least-Once)
[*]消息不重复(At-Most-Once)
[*]精准一次消耗(Exactly-Once)
Kafka 通过消耗偏移量(Offset)管理和变乱消耗等机制实现差别级别的同等性包管。
2.1 消耗者偏移量(Offset)管理
Kafka 接纳**偏移量(Offset)**来记录消耗者消耗的进度,Kafka 提供三种消耗语义:
语义表明偏移提交机会大概的标题At-Most-Once(最多一次)大概丢失消息,但不重复在消耗条件交失败后消息丢失At-Least-Once(至少一次)确保不丢失,但大概重复在消耗后提交失败大概导致重复消耗Exactly-Once(精准一次)消耗恰好一次变乱消耗 + 幂等性须要变乱支持✅ 最佳实践:
[*]默认 Kafka 消耗是 At-Least-Once,即消耗后提交偏移量,大概导致重复消耗。
[*]克制重复消耗:
[*]可联合 幂等性 机制(如数据库 UPSERT 利用)。
[*]利用变乱消耗(见下一节)。
2.2 变乱消耗(Exactly-Once)
Kafka 变乱消耗(Exactly-Once Processing,EoS)包管消耗者端的精准一次处置处罚:
[*]read_process_commit 原子性
[*]变乱包管了读取、处置处罚和提交偏移量要么全部完成,要么全部失败。
[*]Kafka Streams API
[*]Kafka Streams 提供内置Exactly-Once 语义,自动处置处罚变乱提交。
✅ 最佳实践:
[*]利用 Kafka Streams 举行 EoS 消耗(保举)。
[*]如果用平凡消耗者:
[*]enable.auto.commit=false,手动提交偏移量。
[*]联合 commitSync() 和变乱 commitTransaction() 共同确保同等性。
2.3 Rebalance 影响同等性
当 Consumer Group 发生**Rebalance(重新分配分区)**时,大概导致:
[*]重复消耗:如果 Rebalance 发生在偏移量提交前,大概导致部分消息重复消耗。
[*]消息丢失:如果 Rebalance 发生后,某些未提交偏移量的消息未处置处罚完。
✅ 最佳实践:
[*]利用 StickyAssignor 或 CooperativeStickyAssignor,淘汰 Rebalance 影响。
[*]手动提交偏移量(commitSync),确保处置处罚完数据后才提交。
总结
机制生产者同等性消耗者同等性ACK 机制acks=all 确保数据写入乐成-重试机制retries>0,克制瞬时失败-幂等性enable.idempotence=true,防止重复写入-变乱beginTransaction() + commitTransaction()read_process_commit 变乱消耗偏移量管理-enable.auto.commit=false + commitSync()Rebalance 处置处罚-利用 StickyAssignor 方案淘汰影响✅ 终极保举方案:
[*]生产者:acks=all + enable.idempotence=true + transactional.id
[*]消耗者:enable.auto.commit=false + commitSync() + 变乱消耗
这些优化方案可以包管 Kafka "Exactly-Once"(精准一次) 语义,确保生产者和消耗者数据同等性。
页:
[1]