Kafka 如何保证数据不丢失

📅 2026/7/23 16:29:02 👁️ 阅读次数 📝 编程学习
Kafka 如何保证数据不丢失

在分布式消息系统中,数据不丢失是核心诉求之一。Kafka 通过 ACK(Acknowledgement)机制 与 ISR(In-Sync Replica)机制 的协同配合,在生产者、Broker、副本之间构建了一套完整的消息可靠性保障体系。本文将深入剖析 Kafka 是如何在吞吐量与可靠性之间做取舍,从而保证数据不丢失的。
一、ACK 机制:生产者端的确认策略
ACK 是 acknowledge 的缩写,意为"确认"。在 Kafka 中,它指的是 Producer 需要接收来自 Leader Partition 的 ACK 信息,以确认消息已成功写入。
⚠️ 重要提示:ACK 机制的开启会直接影响 Kafka 集群的吞吐量和消息可靠性。吞吐量和可靠性就像硬币的两面,两者不可兼得,只能根据业务特点进行平衡取舍。
ACK 是 Kafka 保证数据不丢失的重要手段,但配置不当也可能导致数据重复。在深入 ACK 之前,我们需要先理解另一个核心概念——ISR。
二、ISR 机制:副本同步的核心设计
背景
Kafka 中每个 Topic 的分区可以设置若干个副本(Leader + Follower),Follower 会异步同步 Leader 的数据。
最初的同步方案:必须等待所有 Follower 都完成同步后,Leader 才向生产者发送 ACK。

  • 优点:当重新选举 Leader 时,只需有 n+1 个副本即可容忍 n 台节点故障。
  • 缺点:延迟极高,必须等待所有 Follower
    同步完成。如果某个 Follower 因网络延迟迟迟不能同步,Leader 会一直阻塞等待,严重影响性能。
    解决方案:Kafka 引入了 ISR(In-Sync Replica)机制。
    什么是 ISR?
    ISR 是与 Leader 保持同步的 Follower 集合。Leader 不需要等待所有 Follower 都完成同步,只要在 ISR 中的 Follower 完成数据同步,就可以发送 ACK 给生产者。
    如果 ISR 集合里的 Follower 延迟时间超过配置参数 replica.lag.time.max.ms,就会从 ISR 中剔除。一旦 Leader 发生故障,就会从 ISR 集合里选举一个 Follower 作为新的 Leader。
    Kafka 的数据同步方式不是完全同步,也不是完全异步,而是基于 ISR 的同步机制。
    AR、ISR、OSR 的关系
    关系公式:AR = ISR + OSR
  • ISR 是 AR 的子集,由 Leader 动态维护。
  • 新加入的 Follower 会先存放在 OSR 中,追上进度后可重新加入 ISR。
  • Leader 副本天然就在 ISR 中,甚至在某些极端情况下,ISR 只有 Leader 这一个副本。
    判断标准:replica.lag.time.max.ms
    Broker 端参数 replica.lag.time.max.ms(默认值 10 秒)是判断 Follower 是否在 ISR 中的核心标准:
  • 只要 Follower 落后 Leader 的时间不连续超过 10 秒,Kafka 就认为该 Follower 与 Leader是同步的。
  • 如果同步速度持续慢于 Leader 的写入速度,超过该时间后,Follower 会被"踢出" ISR。
  • 若该副本后续追上进度,可以重新被加回 ISR,因此 ISR 是一个动态调整的集合。
    三、三种 ACK 配置详解
    ACK 是 Producer 端的配置参数,在创建 Producer 时传入:
Propertiesprops=newProperties();props.put("bootstrap.servers","localhost:9092");props.put("acks","all");// 关键配置props.put("retries",0);props.put("batch.size",16384);props.put("key.serializer",StringSerializer.class.getName());props.put("value.serializer",CustomerSerializer.class.getName());producer=newKafkaProducer<String,Customer>(props);

acks = 0:至多一次(At Most Once)

  • 机制:生产者只负责发消息,不等待任何 ACK 就立即返回。
  • 特点:延迟最低,吞吐量最高。
  • 风险:如果 Leader还未落盘就发生故障,数据会丢失。
  • 适用场景:允许少量数据丢失、对延迟要求极高的场景(如日志采集)。
    acks = 1:至多一次(At Most Once)
  • 机制:Leader 将数据落盘成功后,不管 Follower 是否同步完成,就发送 ACK。
  • 特点:兼顾了一定的可靠性和性能。
  • 风险:保证了 Leader 节点内有一份数据,但如果 Follower 还未同步时 Leader 发生故障,数据会丢失。
  • 适用场景:一般业务场景,可容忍少量数据丢失。
    acks = -1 或 acks = all:至少一次(At Least Once)
  • 机制:生产者等待 Leader 和 ISR 集合内的所有 Follower 都完成同步后,才发送 ACK。
    特点:可靠性最高,保证数据不丢失。 风险:当 Follower 同步完成后、Leader 发送 ACK 之前,如果 Leader
    发生故障:
    会重新从 ISR 内选举新 Leader;
    由于生产者没收到 ACK,会重新发送消息给新 Leader;
    此时会造成数据重复;
  • 适用场景:对数据可靠性要求极高的场景(如金融交易)。
    核心要点
  • ISR 机制 是 Kafka 在吞吐量和可靠性之间做取舍的关键设计——通过动态维护同步副本集合,避免了等待所有副本同步带来的性能瓶颈。
  • ACK 机制 让 Producer 可以根据业务需求灵活选择确认策略:
  • 追求极致性能:acks=0;
  • 平衡性能与可靠性:acks=1;
  • 追求数据绝对不丢失:acks=all;
  • 当使用 acks = all 时,虽然能保证数据不丢失,但可能引入数据重复问题。若需要同时保证"不丢失且不重复",则需要配合 幂等性(Idempotence) 和
    事务(Transaction) 机制来实现 Exactly Once 语义。


我整理了一份Claude Code命令手册,需要的回复Claude Code获取。