Kafka 异步链路,如何验证在 5 万 TPS 下消息不丢不重

解读

面试官真正想确认的是:

  1. 你是否理解 Kafka“至少一次”语义下仍可能因异步刷盘、重平衡、幂等未开启而出现重复或丢失;
  2. 能否把“业务正确性”指标(不丢、不重)转成可量化、可观测、可加压的测试方案;
  3. 是否具备在 5 万 TPS 这一国内头部互联网量级下,用有限成本快速暴露瓶颈并给出优化路径的能力。
    回答时要体现“测试策略 → 观测手段 → 加压模型 → 结果判定 → 风险闭环”五步完整闭环,而不是单纯堆工具。

知识点

  1. Kafka 交付语义:at-least-once、at-most-once、exactly-once(幂等 + 事务)。
  2. 丢消息根因:acks=0/1、leader 切换、ISR 收缩、页缓存未刷盘、consumer 先 commit 后业务处理。
  3. 重复根因:consumer rebalance 后重新拉取已 commit 但业务未落库的消息、producer 重试。
  4. 国内常用观测手段:
    • 消息级指纹:业务 ID + 时间戳 + 连续 sequence,写入消息头;
    • 对账队列:每条消息写 Kafka 同时写对账 Topic(双写),测试侧独立消费比对;
    • 端对端埋点:producer 在本地日志打印 msgKey,consumer 处理完后写 Redis Set,测试侧通过 Redis 集合差集快速计算丢/重。
  5. 加压模型:
    • 目标 5 万 TPS 按 1:1 读写,即消费侧也要 5 万 TPS;
    • 采用阶梯负载:5k→15k→30k→50k→60k(120% 峰值),每阶持续 15 min,观察拐点;
    • 故障注入:kill -9 leader、滚动重启 broker、扩容分区、触发 consumer group rebalance、网络 200 ms 延迟、磁盘打满 95%。
  6. 资源基线:单 broker 网卡 25 Gbps、磁盘 3 GB/s 顺序写、PageCache 50% 内存,超过即判定为环境瓶颈而非业务瓶颈。
  7. SLA 判定:丢数 0、重复 0,P99 端到端延迟 < 2 s,CPU/网卡/磁盘任一利用率 < 80%。

答案

为在 5 万 TPS 下验证“消息不丢不重”,我按“构造数据 → 双轨对账 → 加压干扰 → 指标判定 → 问题闭环”五步落地:

  1. 构造可验证数据
    在 producer 端为每条消息生成全局唯一连续 sequence(Long 型,0 自增),连同业务 key、时间戳写入消息头;同时把 sequence 同步写入本地“对账日志”和独立“对账 Topic”。测试环境独享一组对账消费者,与业务消费逻辑完全隔离。

  2. 双轨对账
    业务消费者在处理完成后,把 sequence 写回 Redis Cluster 的 Set结构(TTL 24 h)。测试侧启动“对账作业”:
    a) 实时消费对账 Topic,累积 sequence 集合 A;
    b) 每 30 s 抓取 Redis Set B;
    c) 计算 A-B 得到“丢失集合”,B-A 得到“重复集合”。
    若两集合均为空,则本窗口内“0 丢 0 重”。窗口滑动 5 min 即可覆盖 5 万×300≈1 500 万条,内存可控。

  3. 加压与干扰
    使用公司自研压测平台(基于 JMeter 改造,支持 Kafka Sampler 异步回调确认):

    • 阶梯加压到 5 万 TPS 后,持续 30 min;
    • 在 50% 时间点人工 kill 当前分区 leader,触发 ISR 收缩;
    • 在 70% 时间点滚动重启全部 consumer,制造 rebalance;
    • 在 90% 时间点使用 tc 命令对 broker 注入 200 ms 网络延迟 5 min。
      全程保持对账作业运行,实时输出丢/重曲线。
  4. 指标判定
    整个 30 min 累计 9 000 万条消息,对账结果丢数=0、重复数=0,P99 端到端延迟 1.8 s,broker CPU 峰值 78%、网卡 68%、磁盘写带宽 2.1 GB/s,均低于 80% 红线,判定通过。
    若出现丢数>0,优先检查 acks 配置、副本数、min.insync.replicas 是否满足“acks=all & min.insync.replicas≥2”;若出现重复数>0,检查是否开启幂等(enable.idempotence=true)及 consumer 处理逻辑是否先业务后 commit。

  5. 问题闭环
    测试报告给出“环境极限值”与“业务安全值”两张图:极限值 6.2 万 TPS 时 CPU 打满,安全值 5 万 TPS 对应水位 78%,建议线上按 4.5 万 TPS 设限流阈值;同时把 sequence 对账逻辑固化到线上巡检,每 10 min 跑一次差集,1 分钟内发现即告警,实现测试资产生产化。

拓展思考

  1. 如果业务已开启 Kafka 事务(exactly-once),是否还需要对账?
    仍需要。事务只能保证“Kafka 内部”不丢不重,但业务处理完写 MySQL、Redis 可能失败,导致“Kafka 已提交、业务未落库”的隐式丢失。因此端对端对账仍是金标准。

  2. 5 万 TPS 成本太高,如何用 1 万 TPS 样本外推?
    采用“等比例放大验证”策略:先用 1 万 TPS 跑 6 小时,累积 2 亿条,对账无误后,再跑 20% 扰动实验(突增到 1.5 万 TPS、断网 30 s),观察对账结果。若仍 0 丢 0 重,可结合线性外推模型(Little’s Law)预测 5 万 TPS 水位,最后以 5 万 TPS 做 15 min 点检即可,节省 70% 资源。

  3. 多云场景下,跨可用区延迟 3 ms→30 ms,如何重新设定 SLA?
    延迟升高会放大副本同步耗时,需提升 replica.lag.time.max.ms 和增加 fetch.min.bytes,避免 ISR 频繁抖动;同时把 P99 端到端延迟 SLA 放宽到 4 s,但“0 丢 0 重”红线不变,测试方案需在跨 AZ 故障注入阶段重点验证分区 leader 跨机架切换时的对账结果。