欢迎光临
我们一直在努力

Kafka如何保证消息不丢

Kafka 的“消息不丢”不是一个配置能解决的,而是要在 生产者 → Broker → 消费者 三个环节同时做保证。
下面从“为什么会丢”到“怎么解决”做一次完整拆解。

先明确:Kafka 消息丢失发生在哪几层?

一条消息的生命周期:

Producer 发送 → Broker 存储/复制 → Consumer 拉取 → Consumer 提交 offset

所以丢失可能发生在:

  • 生产者端:发送失败、确认策略太弱、缓冲丢失、重试不合理

  • Broker 端:副本数不足、ISR 门槛太低、unclean leader 选举、刷盘策略

  • 消费者端:offset 提前提交、自动提交、rebalance 时处理不当

  • 端到端语义:重复消费、消费逻辑非幂等、事务缺失

  • 生产者

    Kafka Producer(生产者)是向 Kafka Broker 发送消息的客户端,负责把业务数据封装成消息,发送到指定 Topic。核心流程:消息构建 → 序列化 → 分区选择 → 累加器缓存 → 发送线程批量发送 → Broker 应答确认。

    核心架构组件

  • 主线程(业务线程) 业务代码调用 producer.send(),不会直接发网络请求。消息先做序列化、分区计算,放入RecordAccumulator(记录累加器)缓存,send 默认异步,返回 Future。
  • RecordAccumulator 累加器(消息缓冲区) 生产者核心缓存,消息不会来一条发一条,而是按 topic‑partition 分组缓存,攒批量。
    • 每个分区对应一个双端队列,队列元素是 ProducerBatch(批量批次)
    • 一个 Batch 包含多条消息,达到大小或者时间阈值就触发发送
    • 配置参数:
      • buffer.memory:缓冲区总内存大小
      • batch.size:单个批次最大字节数
      • linger.ms:延时发送,没填满 batch 时,最多等待多久再发送,用来攒批量,提升吞吐量
  • Sender 发送线程 后台独立线程,不断扫描累加器,把就绪的 Batch 取出来,转化为网络请求,发给对应 Broker;处理 Broker 返回的响应,处理成功 / 失败、重试。
  • 分区器 Partitioner 决定这条消息发到 Topic 的哪一个分区。
  • 分区规则:

  • 如果消息指定了 partition,直接使用该分区;
  • 如果没指定分区,但有 key:对 key 做 hash 取模分区,同一个 key 一定落到同一个分区;
  • 既无 partition 也无 key:轮询策略,依次分发到各个分区,负载均衡。
  • Kafka 默认分区器:DefaultPartitioner

  • 序列化器 Serializer 把 Java/Go 对象转成字节数组,网络只能传输字节。内置:StringSerializer、ByteArraySerializer,复杂对象需要自定义序列化(JSON、Protobuf、Avro)。
  • 消费者必须使用对应的反序列化器。

  • 元数据管理器 Metadata 维护 Topic 的元数据:有哪些分区、分区对应的 leader broker、ISR 集合。 生产者发送前如果没有元数据,会先向 Broker 请求更新元数据。
  • 发送消息三种方式

    1. 异步发送(不关心结果,最常用)

    producer.send(new ProducerRecord<>("topic", "key","value"));

    直接返回,不等待 broker 响应,高吞吐;消息丢失风险取决于 acks 和重试配置。

    2. 异步 + 回调(推荐,感知成功失败)

    producer.send(record, new Callback() {
    @Override
    public void onCompletion(RecordMetadata metadata, Exception e) {
    if(e != null){
    // 发送失败
    }else{
    // 发送成功,拿到分区、offset
    }
    }
    });

    回调在 Sender 线程执行,回调逻辑不要写耗时操作,会阻塞发送线程。

    3. 同步发送

    调用 get() 阻塞等待结果:

    RecordMetadata meta = producer.send(record).get();

    吞吐量很低,适合强可靠场景。

    核心可靠性参数

    1. acks 应答级别

    控制 broker 多少副本写入成功才返回成功给生产者。

    • acks=0:生产者发完就认为成功,不等 broker 响应。吞吐量最高,可能丢消息。
    • acks=1(默认):仅 leader 副本写入成功就应答。leader 刚写完挂掉,还没同步到从副本,消息丢失。
    • acks=all /‑1:leader 和所有 ISR 内副本全部写入完成才应答。可靠性最高,吞吐量下降。

    2. retries 重试次数

    消息发送失败时生产者自动重试。配合 retry.backoff.ms 重试间隔。

    注意:重试会带来消息重复,业务要做幂等处理。

    3. max.in.flight.requests.per.connection

    单个 broker 连接上,未完成的请求最大数量。

    • Kafka <1.0:大于 1 时,重试会造成消息乱序;
    • Kafka >=1.0 开启 enable.idempotence 幂等生产者,即使大于 1 也能保证顺序。

    4. enable.idempotence 幂等生产者

    默认开启 true。 解决:重试带来消息重复。原理:生产者分配 PID,每条消息带序列号,Broker 保存序列号,重复序列号直接丢弃。

    幂等只能保证单会话、单分区不重复;生产者重启之后 PID 变化,无法跨会话幂等。

    5. 事务生产者 transactional.id

    用于跨分区、跨 Topic 的原子写入;要么全部成功,要么全部失败,配合消费者的 read_committed 隔离级别。

    消息发送完整流程

  • 业务线程调用 send (record)
  • 判断 topic 元数据是否存在,不存在则拉取元数据
  • 使用序列化器把 key、value 序列化为字节数组
  • Partitioner 计算目标分区
  • 尝试把消息写入 RecordAccumulator 的对应分区队列的 ProducerBatch
    • 如果当前 batch 未满,写入;
    • 如果满了,新建 batch;
  • 满足条件(batch.size 满 /linger.ms 超时),batch 变为待发送;
  • Sender 线程拉取待发送 batch,组装网络请求,发送给对应 Leader Broker
  • Broker 处理:写入日志,按照 acks 策略等待副本同步,返回响应
  • Sender 收到响应:
    • 成功:释放 batch 内存,回调 onCompletion 成功;
    • 失败:判断是否可以重试,放回累加器等待重试;不可重试则回调异常。
  • 重要参数汇总

    参数作用
    bootstrap.servers broker 地址列表
    key.serializer / value.serializer 序列化器
    buffer.memory 生产者缓冲区总内存
    batch.size 单个批次大小
    linger.ms 等待攒批时间
    acks 副本应答策略
    retries 重试次数
    enable.idempotence 幂等生产者开关
    max.in.flight.requests.per.connection 单连接并发请求数
    compression.type 压缩算法:none、lz4、snappy、gzip,压缩提升吞吐量

    常见问题

  • 消息丢失
    • acks=0/1,leader 宕机;缓冲区满阻塞 / 丢弃;生产者异常退出,缓冲区未发送消息丢失。 解决:acks=all,合理重试,设置回调捕获异常。
  • 消息重复 网络抖动,Broker 已经写入成功,但是响应丢了,生产者重试。 解决:开启幂等,业务消费端做幂等。
  • 消息乱序 旧版本无幂等,max.in.flight>1,前面请求失败重试,后面请求先成功,造成乱序。 解决:开启幂等生产者,或者设置 max.in.flight.requests.per.connection=1。
  • 吞吐量调优 调大 batch.size、linger.ms;开启压缩;合适的 buffer.memory;acks=1;异步发送。
  • 阻塞 send () 当 RecordAccumulator 缓冲区满了,send () 会阻塞 max.block.ms 时间,超时抛异常。
  • Kafka Producer(生产者)可能丢失消息的场景

    前提:消息还没收到 Broker 成功 ACK,一旦生产者拿到成功 ACK,消息就已经持久化到 ISR 副本,生产者侧不再丢失。

    场景 1:消息在 RecordAccumulator 缓冲区,还没发送,进程宕机

    消息写入 ProducerBatch,还没满足发送条件(batch 未满、linger.ms 没到),仍然保存在生产者内存。 生产者进程直接崩溃、kill -9,内存数据全部清空,消息永久丢失。

    • 常见:linger.ms 设置较大,攒批时间长;服务下线没有执行 flush () 和 close ()。
    • 规避:优雅停机,调用producer.flush()再 close;业务层持久化消息,发送成功后再删除本地记录。

    场景 2:acks=0,发送后不等 Broker 应答

    生产者发出 TCP 数据包直接标记发送成功,完全不等待 Broker 返回结果。 网络丢包、Broker 宕机,Broker 根本没有收到消息,生产者却认为发送成功 → 消息丢失。

    • 规避:高可靠场景禁止 acks=0。

    场景 3:acks=1,Leader 写入成功、Follower 尚未同步,Leader 宕机

    Leader 副本写入本地日志,立刻返回 ACK,此时 Follower 还没拉取同步这条消息。 Leader 机器故障宕机,集群从 ISR 选新 Leader,新 Leader 没有这条消息 → 消息丢失。

    这个丢失发生在 Broker 侧,但根源是生产者应答级别配置过低。

    • 规避:高可靠场景使用acks=all。

    场景 4:缓冲区已满,send 阻塞超时抛出异常,消息丢弃

    Broker 不可用,Sender 无法发送,累加器缓冲区持续堆积打满。 后续调用 send () 会阻塞,最长等待max.block.ms,超时抛出 TimeoutException,这条消息直接丢弃。 业务代码如果没有捕获异常、没有兜底保存消息,消息丢失。

    • 规避:捕获超时异常,将消息存入本地 / 数据库死信,后续重试;合理调大 buffer.memory。

    场景 5:不可重试异常,消息直接丢弃,不会重试

    发生这类异常时,生产者不会重试,直接回调失败。如果业务不处理,消息丢失:

  • 消息大小超过message.max.bytes;
  • key/value 序列化失败;
  • Topic 不存在、没有写入权限;
  • 分区不存在。
    • 规避:Callback 捕获异常,不可重试消息写入死信队列。

    场景 6:异步发送无 Callback,重试耗尽后静默丢失

    业务代码直接producer.send(record)不传入回调。 多次重试之后仍然发送失败,Sender 线程抛出异常,业务代码完全感知不到,消息悄悄丢失。

    这是生产中非常经典的坑。

    • 规避:所有 send 必须携带 Callback,异常打日志、告警、落死信。

    场景 7:flush () 之后立刻宕机

    flush()的作用是把累加器里所有待发送 batch 交给 Sender 线程发起网络请求,不代表 Broker 写入成功。 flush 调用完成,但是网络请求还在路上、还没收到 ack,进程宕机,消息丢失。

    很多人误解 flush = 保证落地,这是错误的。

    场景 8:元数据拉取失败超时,消息无法写入

    生产者没有 topic 元数据,请求 Broker 获取元数据长时间失败,send 阻塞 max.block.ms 超时,消息丢弃。


    生产者消息重复、乱序场景

    消息重复

    Broker 已经成功写入消息,但是 ACK 包在网络传输中丢失;生产者超时触发重试,再次发送同一条消息 → Broker 收到两条一样的消息。

    • 解决:开启enable.idempotence=true幂等生产者,PID + 序列号自动过滤重复消息;
    • 限制:幂等生产者仅单会话、单分区生效,生产者重启后 PID 变化,无法跨会话防重复。跨分区原子写入需要事务。

    消息乱序

    Kafka 1.0 之前,没有幂等生产者,max.in.flight.requests.per.connection>1:

    • 第 1 批消息发送失败,等待重试;
    • 第 2 批消息发送成功;
    • 第 1 批重试成功,消息顺序颠倒,产生乱序。
    • 解决:开启幂等生产者;或者max.in.flight.requests.per.connection=1,同一连接只允许一个未完成请求。

    生产者防丢总结

    生产者消息丢失,基本都发生在消息拿到成功 ACK 之前。

  • 内存缓冲区消息,进程异常退出丢失;
  • acks 配置过低(0/1),Leader 宕机丢失;
  • 缓冲区满、元数据超时,send 抛异常未兜底;
  • 不可重试异常、无回调,失败静默丢弃;
  • flush 不等于消息落盘。
  • 高可靠生产者配置

    acks=all
    enable.idempotence=true
    retries=2147483647
    retry.backoff.ms=100
    max.in.flight.requests.per.connection=5
    buffer.memory=67108864
    batch.size=16384
    linger.ms=5

    Broker端

    Broker 端消息丢失,几乎全部集中在 "副本复制不完整" 或 "数据未落盘" 这两个根因上—— 只要 Leader 故障时,ISR 里没有一份完整副本可用,或所有副本的数据都还在内存页缓存里,消息就会丢。下面系统拆解。

    说明:图为 Kafka 副本、应答门槛与丢失机制的通用结构示意,非特定集群测量结果。

    Broker 端可靠性依赖的三根支柱(先理解,再谈丢)

  • 多副本:每个分区有 Leader + 若干 Follower,分布在不同机器,防止单机故障。
  • ISR(In-Sync Replicas):与 Leader 保持同步的副本集合,是 "哪些副本数据完整可用" 的判据。
  • Leader 选举:Leader 故障时从副本中选新 Leader,选谁决定了数据会不会丢。
  • 一句话:只要 "选出的新 Leader 里没有完整数据",Broker 就丢消息。

    Broker 端可能丢消息的完整场景

    场景 A:acks=1,Leader 写完 ACK 后宕机(最经典)

  • Leader 收到消息,写入本地日志,立刻返回 ACK 给生产者。
  • Follower 还没来得及拉取同步这条消息,Leader 宕机。
  • 从 ISR 选出新 Leader(原 Follower),它没有这条数据。
  • 旧 Leader 恢复后变 Follower,会截断本地多出的日志对齐新 Leader → 消息永久丢失。
    • 根治:acks=all。

    场景 B:acks=all,但 ISR 收缩到只剩 Leader 一个副本

    • acks=all 的意思是 " 等 ISR 内所有副本写完 ",不是等全部副本。
  • 两个 Follower 因网络 / GC 跟不上,被踢出 ISR。ISR 只剩 Leader。
  • 生产者 acks=all,ISR 只有 Leader,Leader 写完即返回成功 → 等价于 acks=1。
  • 紧接着 Leader 宕机,ISR 无副本可用 → 丢消息。
    • 根治:min.insync.replicas(topic 级)。设 2 表示 ISR 至少 2 份才允许写入;不足时直接拒绝写入(抛异常),宁可写不进去也不丢。

    场景 C:unclean.leader.election.enable=true,选举非 ISR 旧副本当 Leader

    • 参数默认 false(只能从 ISR 里选 Leader,ISR 全挂则该分区不可写,保可用性弃可用性)。
    • 设为 true:允许选 数据不完整的非 ISR 副本 当 Leader。
  • ISR 内副本全部宕机;一个数据较旧的副本还活着。
  • unclean 开启,旧副本被选为新 Leader。
  • 它缺少一部分最新消息;旧 Leader 恢复后被截断对齐 → 生产者已 ACK 的消息丢失。
    • 根治:高可靠场景 关闭 unclean.leader.election.enable。

    场景 D:ISR 内所有副本同时断电,数据仍在 PageCache 未落盘

    • Kafka 写入:消息先写进操作系统 PageCache(页缓存) 就返回,不是立即刷盘;刷盘由 OS 后台或 log.flush.* 控制。
  • 消息已写入各副本的页缓存,还没落到物理磁盘。
  • ISR 内所有机器同时断电,页缓存数据全部丢失,磁盘上没有 → 永久丢失。
    • 说明:只要有一个副本已刷盘且存活,选举它当 Leader 就不会丢。只有 ISR 全部机器同时断电才触发。
    • 注意:不要靠频繁刷盘防丢(log.flush.interval.ms 调很小)会严重伤性能。Kafka 的可靠性靠副本,不靠刷盘。

    场景 E:磁盘损坏 / 日志文件丢失(运维事故)

    • 磁盘硬件故障、数据文件损坏、人为误删 topic 的 segment 日志。
    • 若对应副本也被破坏,消息无法恢复。

    场景 F:Follower 长期落后被踢出 ISR,副本数悄悄减少

    • Follower 因网络抖动、GC 停顿长期追不上 Leader,被移出 ISR。
    • 副本 "实际可用的" 数量下降。若此时 Leader 故障,可用的完整副本更少,丢失概率上升(本质是场景 B 的诱因)。
    • 属于事前预警,不是立即丢,但会放大其它场景的风险。

    Broker 端 "什么是安全的" 对照

    配置组合可靠性说明
    acks=0 最低 发出即成功,网络 / Broker 问题即丢
    acks=1(默认) Leader 写入即成功,Leader 宕机会丢
    acks=all + 副本 3 + min.insync.replicas=2 + unclean=false ISR≥2 才写,只从 ISR 选 Leader
    副本数 = 1 极低 无冗余,任何故障即丢

    Broker 端防丢最佳实践(可直接背)

    # Broker 端
    default.replication.factor=3 # 默认副本数 3
    min.insync.replicas=2 # topic 级:ISR 至少 2 份才写入
    unclean.leader.election.enable=false # 禁止非 ISR 选主
    # 生产者端配合
    acks=all
    enable.idempotence=true

    • 生产环境禁止 auto.create.topics.enable=true 自动建单副本 topic(默认单副本,无冗余)。
    • 多 Broker 跨机架部署,避免 ISR 副本集中在同一台物理机 / 同一机架(防场景 D)。
    • 磁盘做好 RAID 冗余,降低场景 E。

    高频面试题

    Q:acks=all 就绝对不丢吗? 不是。① ISR 收缩到只剩 Leader,会退化成 acks=1;② ISR 全部同时断电且未刷盘;③ unclean 选举非 ISR 副本。所以 acks=all 必须配合 min.insync.replicas、unclean=false。

    Q:为什么 Kafka 不靠频繁刷盘保证不丢? 因为 Kafka 的可靠性模型是 "多副本 + ISR",靠副本冗余而不是单机刷盘;强制频繁 fsync 会大幅降低吞吐,且仍无法抵御整机断电。只有 ISR 全机断电才丢,属于极端场景。

    Q:unclean.leader.election.enable 怎么权衡? false = 保数据不丢,但 ISR 全挂时分区不可写(牺牲可用性);true = 保可用性,但可能丢数据。金融 / 订单类必须 false。

    Kafka Consumer(消费者端)

    核心一句话:消费者是主动拉取模型,消息丢失 / 重复几乎全部由 offset 提交时机决定,Broker 保存 offset 到__consumer_offsets系统主题。

    一、消费者核心架构与基础概念

    1. 消费模型

    Kafka 消费者不是 Broker 推送消息,而是消费者主动调用 poll () 拉取。 消费者组 ConsumerGroup:同一个 group.id 下所有消费者构成一个消费者组。

    • 一个分区,同一时刻只能分配给组内一个消费者(保证单分区顺序消费);
    • 组内消费者数量 > 分区数:多余消费者空闲,没有分区;
    • 组内消费者数量 < 分区数:一个消费者消费多个分区。

    2. Offset(偏移量)

    offset 是分区内消息的序号,代表消费位置。offset 持久化存储在 Broker 的__consumer_offsets系统 topic。

    • 当前消费到哪条消息,由提交的 offset 决定;
    • 提交 offset,代表 Kafka 认为这条消息以及之前所有消息已经消费完成。

    重点:消息是否被处理成功,Kafka 本身完全不知道,只能相信消费者提交的 offset。 这是消费端所有问题根源。

    3. 关键角色

  • GroupCoordinator:Broker 中的协调器,管理这个消费者组,接收心跳、offset 提交、处理 rebalance。每个消费者组固定一台 Broker 作为协调器。
  • Consumer Leader:组内第一个加入的消费者,负责执行分区分配策略,生成分配方案交给 GroupCoordinator 下发给所有成员。
  • Fetcher:消费者内部组件,负责从对应分区 Leader Broker 拉取消息。
  • 二、消费者工作流程

  • 消费者启动,指定 group.id,连接 Broker,寻找 GroupCoordinator。
  • 向协调器发起加入组请求,触发 Rebalance。
  • Leader 生成分区分配方案,Coordinator 下发,消费者拿到自己负责的分区。
  • 循环调用poll():
  • Fetcher 从分区拉取一批消息(max.poll.records 控制条数);
  • 业务代码处理消息;
  • 根据 offset 提交策略,提交 offset 到 Broker 的__consumer_offsets。
  • 持续发送心跳(heartbeat)给 GroupCoordinator,证明消费者存活。
  • 发生 rebalance 时,停止消费,释放分区,等待重新分配分区。
  • 三、Offset 提交的两种方式

    1. 自动提交 enable.auto.commit=true(默认开启)

    enable.auto.commit=true
    auto.commit.interval.ms=5000 # 默认5秒自动提交

    机制:poll 拉取消息后,每到时间间隔,自动提交上一轮 poll 拉取消息的 offset。 ✅ 优点:简单,不用手动写提交代码 ❌ 巨大风险,极易消息丢失 场景:poll 拿到消息,业务正在处理,还没处理完;5s 时间到,自动提交 offset。随后业务进程崩溃,这批消息已经提交 offset,Kafka 认为消费完成,重启消费者不会重新拉取 → 消息丢失

    自动提交只适合:消息丢了影响很小的日志采集场景;订单、支付等强可靠业务严禁使用自动提交。

    2. 手动提交 enable.auto.commit=false(高可靠业务推荐)

    业务代码自己控制什么时候提交 offset,分两种:

  • 同步提交 commitSync () 阻塞调用,等待 Broker 返回提交成功 / 失败;失败会抛异常,可以重试。 缺点:阻塞消费线程,降低吞吐。
  • 异步提交 commitAsync () 非阻塞,提交请求发出去,不等待应答;支持回调接收提交结果。 优点:吞吐更高;缺点:无法即时捕获失败。
  • 最佳实践:消息全部处理成功之后,再提交 offset,保证 At-Least-Once(至少消费一次,可能重复,不会丢)。

    四、消费者端消息丢失场景

    消费者端天然不容易丢消息,错误的 offset 提交逻辑才会丢消息

    场景 1:自动提交 offset,消息未处理完成,进程崩溃(最经典)

    流程:

  • poll 拉取消息;
  • 业务正在处理;
  • 到达 auto.commit.interval.ms,自动提交 offset;
  • 业务处理中途宕机; Kafka 已经记录 offset,重启后从下一条开始消费,未处理的消息直接丢失。
  • 场景 2:手动提交,先提交 offset,后执行业务逻辑(代码写反)

    // 错误写法!!
    consumer.commitSync(); //先提交offset
    processMessage(msg); //再处理业务

    提交 offset 成功,业务执行失败 / 程序崩溃 → 消息丢失。

    场景 3:重置 offset 到更早位置以外的错误位置

    手动调用 seek 把 offset 设置到当前消息之后,跳过一批消息,直接丢失中间消息。

    五、消费者端消息重复消费场景(高频考点)

    At-Least-Once:Kafka 默认语义,消息至少被消费一次,大概率出现重复

    场景 1:业务处理成功,但 offset 提交失败

    业务逻辑执行完成,准备提交 offset,网络异常 / 进程宕机,offset 没有持久化到 Broker。消费者重启,从旧 offset 重新拉取,再次消费同一条消息。

    场景 2:Rebalance 再均衡导致重复消费(最常见)

  • poll 拉取消息,业务处理中,offset 还没有提交;
  • 触发 Rebalance,分区被分配给组内另一个消费者;
  • 新消费者从旧 offset 拉取消息,重新消费 → 重复消费。
  • 原因:rebalance 发生时,已经拉取、处理,但 offset 未提交。

    场景 3:消费超时被踢出组(max.poll.interval.ms)

    业务处理消息很慢,两次 poll 间隔超过max.poll.interval.ms,协调器判定消费者卡死,踢出组触发 rebalance,未提交 offset 消息重新分配,产生重复。

    解决办法:业务处理放到独立线程池,poll 主线程只负责拉取消息、控制 poll 调用,不要阻塞 poll 线程;减小 max.poll.records,减少单次拉取消息数量。

    六、消费者关键配置参数

    参数作用
    group.id 消费者组标识,相同 group.id 属于同一组
    enable.auto.commit 是否自动提交 offset,高可靠场景 false
    auto.commit.interval.ms 自动提交间隔
    max.poll.records 一次 poll 最多拉多少条消息,消息处理慢要调小
    max.poll.interval.ms 两次 poll 最大间隔,超时被踢出组,触发 rebalance
    session.timeout.ms 会话超时,心跳超时被踢出组
    heartbeat.interval.ms 心跳发送间隔,一般是 session.timeout.ms 的 1/3
    auto.offset.reset 分区无 offset 时的重置策略earliest:从头开始;latest:最新消息开始;none:抛异常

    七、消费语义

  • At-Most-Once(最多一次,可能丢消息) 先提交 offset,再处理业务;处理失败,消息丢失。
  • At-Least-Once(至少一次,可能重复,不丢) 先处理业务,成功后再提交 offset;业务成功,提交失败,消息重复。Kafka 默认能力。
  • Exactly-Once(精确一次) 跨分区 / 跨 topic 需要 Kafka 事务;单分区单会话可以靠生产者幂等 + 消费端业务幂等实现。
  • Kafka 本身无法直接实现 ExactlyOnce,需要业务层配合。

    八、消费者防丢 & 防重复最佳实践

  • 关闭自动提交 enable.auto.commit=false,手动提交 offset;
  • 业务处理完成后,再提交 offset;
  • 消息处理放到独立线程池,poll 主线程不能阻塞;
  • 合理设置max.poll.records,不要一次性拉取大量消息;
  • 业务增加幂等(唯一消息 ID,去重),解决重复消费问题;
  • 使用 CooperativeStickyAssignor 减少 rebalance 分区迁移,降低重复概率;
  • 滚动发布配置group.instance.id静态成员,减少不必要 rebalance。
  • 九、面试简答

    Q:消费者为什么会重复消费? A:业务处理成功,但 offset 未提交,发生 rebalance 或者消费者重启,新消费者从旧 offset 拉取消息,导致重复。

    Q:消费端怎么实现不丢消息? A:关闭自动提交,消息处理成功之后手动提交 offset;保证业务处理成功,再提交 offset。

    Q:auto.offset.reset 什么时候生效? A:分区没有保存该消费者组的 offset 的时候(第一次消费、offset 被删除),才会生效。已经有 offset,不会触发。

    赞(0)
    未经允许不得转载:171主机测评 » Kafka如何保证消息不丢
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址