第六部分:消息队列架构(第 61~70 题)
进入「第三阶段架构师能力」。覆盖 MQ 选型、不丢失、重复与幂等、顺序性、积压、死信、延迟消息、高性能调优、事务消息、可观测。每题需给出:语义保证→方案→代码→兜底。
第 61 题:Kafka vs RocketMQ vs RabbitMQ 选型 选型
- 1. 三者定位差异?2. 各适合什么场景?3. 吞吐量对比?4. 事务消息谁支持?5. 延迟消息?
第一步:分析。选型看语义需求:高吞吐日志用 Kafka;金融级事务/顺序/延迟用 RocketMQ;复杂路由/低延迟用 RabbitMQ。
第二步:挑战。场景匹配、语义支持、运维。
第三步:架构。①日志/埋点:Kafka(百万 TPS、分区并行、可丢);②订单交易:RocketMQ(事务消息、exactly-once 语义、顺序);③延迟关单:RocketMQ 延迟消息(18 级)或 RabbitMQ TTL+死信;④复杂路由:RabbitMQ(Exchange 绑定)。
第四步:选型。混合:Kafka(日志/流)+ RocketMQ(交易/延迟)+ RabbitMQ(内部复杂路由,可选)。
第五步:一致性。Kafka at-least-once + 消费幂等;RocketMQ 事务消息保证本地事务与发消息原子。
第六步:高可用。三者均多副本;Kafka ISR、RocketMQ 主从、RabbitMQ 镜像队列。
第七步:优化。按场景分集群,避免混用致相互影响。
Kafka:分区+顺序写+零拷贝,吞吐极致,但事务/延迟支持弱,适合日志流。RocketMQ:阿里系,事务消息(半消息+回查)、定时/延迟消息、顺序消费成熟,适合交易。RabbitMQ:AMQP,Exchange 灵活路由,低延迟,但吞吐与堆积能力弱于前两者。一张对比表最直观:
| 维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 吞吐 | 极高 | 高 | 中 |
| 事务消息 | 弱 | 强(半消息) | 无原生 |
| 延迟消息 | 需外挂 | 原生(18级) | TTL+死信 |
| 顺序 | 分区内 | 队列内 | 单队列 |
| 适用 | 日志/流 | 交易/金融 | 内部路由 |
// RocketMQ 事务消息(订单) Message m = new Message("order", JSON.toJSONString(order).getBytes()); SendResult r = txProducer.sendMessageInTransaction(m, order); // Kafka 日志(可丢, acks=1, 批量) props.put("acks", "1"); props.put("linger.ms", "5");
追问 1:能不能只用 Kafka?能,但事务/延迟要自研(外挂),不如 RocketMQ 原生省心;看团队。
追问 2:RabbitMQ 吞吐不够?内部低流量路由够用;高吞吐走 Kafka/RocketMQ。
追问 3:延迟消息精度?RocketMQ 18 级固定(1s/5s/...2h),非任意;任意精度用时间轮/DB 扫描。
追问 4:运维成本?Kafka/RocketMQ 集群运维重;RabbitMQ 轻但规模受限。
追问 5:混用复杂度?按域分集群,SDK 统一封装,治理成本可控。
- ❌ 只用一种 MQ 包打天下——语义不匹配。正确:按场景选。
- ❌ 交易用 Kafka 无事务——丢消息。正确:RocketMQ 事务。
- ❌ 日志用 RabbitMQ——吞吐不够。正确:Kafka。
第 62 题:消息不丢失全链路保障 可靠性
- 1. 哪些环节会丢?2. 生产者怎么保?3. Broker 怎么保?4. 消费者怎么保?5. 极端宕机?
第一步:分析。丢失在三处:生产(未确认)、Broker(未刷盘/副本)、消费(未处理就 ack)。需逐环节加固。
第二步:挑战。性能 vs 可靠权衡、确认机制、刷盘。
第三步:架构。①生产:同步发送 + 重试 + 回调确认(RocketMQ 同步 send;Kafka acks=all+重试);②Broker:同步刷盘 + 多副本(Raft/ISR);③消费:先处理业务再 ack(非自动提交),失败不 ack 触发重试;④兜底:本地消息表/事务消息防生产丢。
第四步:选型。RocketMQ(同步双写+同步刷盘)/ Kafka(acks=all+min.insync.replicas)。
第五步:一致性。at-least-once + 消费幂等 = 不丢不重。
第六步:高可用。多副本防 Broker 丢;消费重试防处理丢。
第七步:优化。关键链路同步刷盘,非关键异步提吞吐。
为什么先处理后 ack?若先 ack 再处理,处理崩=丢消息;先处理后 ack,处理失败重投(at-least-once)。acks=all/min.insync.replicas:Leader 等多数副本写入才确认,防单点丢。代价:同步刷盘/多副本降吞吐,关键链路值得。
// Kafka: acks=all + 重试, 防生产/Broker丢 props.put("acks", "all"); props.put("min.insync.replicas", "2"); props.put("retries", "5"); // 消费: 手动提交, 处理完再 ack(关闭 auto-commit) consumer.subscribe("order"); for (Record r : consumer.poll(1000)) { handle(r); consumer.commitSync(); // 处理成功才提交 }
追问 1:同步刷盘太慢?关键链路同步,非关键异步;或 SSD + 批量刷盘提速。
追问 2:消费处理成功但 ack 失败?消息重投 → 靠幂等(见 Q63)防重复副作用。
追问 3:Broker 全挂?多副本+跨机架;极端丢用本地消息表补发。
追问 4:重复还是丢失优先?业务多选「不丢」(at-least-once)+ 幂等,绝不丢。
追问 5:性能影响多大?同步确认+刷盘降 30%~50% 吞吐,关键链路可接受。
- ❌ 发完就当成功——未确认丢。正确:同步+确认。
- ❌ 先 ack 再处理——处理崩丢。正确:处理后 ack。
- ❌ 异步刷盘单副本——Broker 丢。正确:同步+多副本。
第 63 题:消息重复与幂等消费 幂等
- 1. 为什么会有重复?2. 幂等怎么实现?3. 数据库幂等?4. Redis 幂等?5. 幂等键怎么选?
第一步:分析。MQ 是 at-least-once(不丢),重复必然存在。应对=消费幂等(同消息处理多次=一次效果)。
第二步:挑战。重复来源、幂等键、并发、存储。
第三步:架构。①幂等键:消息带业务唯一 ID(如支付流水号);②去重:a) 数据库唯一索引/冲突更新(INSERT ... ON DUPLICATE);b) Redis SETNX 幂等表(带 TTL);c) 状态机(已处理则跳过);③并发:唯一键保证原子;④兜底:对账修复。
第四步:选型。唯一索引(最稳)+ Redis 幂等表(高性能)+ 状态机。
第五步:一致性。幂等保证重复处理不重复副作用;与「不丢」互补。
第六步:高可用。幂等存储高可用;失败重试不破坏幂等。
第七步:优化。幂等表设 TTL 防无限增长;热点用 Redis。
为什么重复不可避免?at-least-once 语义下,网络/ack 超时必重投。幂等键:必须用稳定业务 ID(支付流水号),不能用 MQ 自带的 offset(会变)。唯一索引最稳:INSERT 冲突即重复,DB 层原子保证;Redis SETNX 高性能但有 TTL 窗口需对账兜底。
// DB 唯一索引幂等(支付流水号唯一) try { pointMapper.add(txId, uid, score); } // 唯一键冲突=重复 catch (DuplicateKeyException e) { return; } // 跳过 // Redis 幂等表(高性能) if (!redis.setnx("dedup:"+txId, "1", 24h)) return; // 已处理跳过 handle(msg); // 状态机: if(order.status==PAID) return;
追问 1:幂等表 TTL 内过期?窗口外重复靠对账/唯一索引兜底;TTL 设足够长(如 7 天)。
追问 2:并发同 ID 同时到?唯一索引/SETNX 原子,后者冲突即重复。
追问 3:幂等键冲突(不同业务同 ID)?幂等键 = 业务域+流水号,避免跨域冲突。
追问 4:消费部分失败?事务内「处理+记幂等」原子;失败回滚重投。
追问 5:Redis 幂等表丢?Redis 不可靠时以 DB 唯一索引为准,Redis 加速。
- ❌ 不处理重复——重复副作用。正确:幂等。
- ❌ 用 offset 做幂等键——会变。正确:业务 ID。
- ❌ 仅 Redis 无兜底——丢后重复。正确:DB 唯一索引。
第 64 题:消息顺序性保障 顺序
- 1. 为什么乱序?2. 怎么保证分区内有序?3. 全局有序?4. 并发消费与顺序矛盾?5. 乱序怎么修复?
第一步:分析。MQ 只保证分区/队列内有序。全局有序代价极高(单分区/单消费者)。按业务键路由到同一分区即可。
第二步:挑战。路由、并发、全局序。
第三步:架构。①路由:同一订单号 hash 到同分区(Kafka partitionKey=orderId;RocketMQ 同 queue);②单分区单消费线程:保证顺序处理;③全局序:仅单分区单消费者(牺牲并行,仅极端需要);④修复:状态机拒绝非法跃迁 + 重排序缓冲。
第四步:选型。Kafka(partitionKey)/ RocketMQ(MessageQueue 顺序消费)。
第五步:一致性。顺序消息保证状态流转正确;失败重试保持顺序(不跳)。
第六步:高可用。顺序消费失败阻塞该分区(或特殊处理),不影响其他分区。
第七步:优化。仅对需顺序的 topic 用顺序,其余并行提吞吐。
分区内有序:Kafka 同 partition 内 offset 递增,单消费者线程顺序读即有序。路由关键:orderId 作 key 保证同订单同分区。全局序代价:单分区单消费者=无并行,吞吐极低,几乎不用。状态机防乱序:即使极少乱序,状态机拒绝「发货前无支付」,可缓冲重排。
// Kafka: 同 orderId 同分区 producer.send(new ProducerRecord<>("order", orderId, msg)); // RocketMQ: 同 key 同 queue + 顺序消费监听器 SendResult r = producer.send(msg, (q,m)-> q.get( Math.abs(orderId.hashCode())%q.size()), msg); // 消费: 状态机校验 if (!order.canTransit(toStatus)) { buffer.reorder(msg); return; }
追问 1:分区扩容后顺序乱?扩容不改 key 路由逻辑,旧数据已落旧分区;新数据按新分区,同 key 仍同分区(同算法)。
追问 2:顺序消费很慢?仅必要 topic 顺序;按 key 细分分区提并行度(同 key 内仍序)。
追问 3:消费失败卡住?顺序消费失败会阻塞该队列;需快速失败+死信,避免全卡。
追问 4:跨服务顺序?用 correlationId + 状态机 + 缓冲重排,跨服务难严格序。
追问 5:要不要全局序?几乎不需;业务设计成「最终状态正确」比全局序更优。
- ❌ 多分区并发要全局序——乱。正确:key 路由同分区。
- ❌ 顺序消费不用单线程——乱。正确:单分区单线程。
- ❌ 无状态机——非法跃迁。正确:状态机兜底。
第 65 题:消息积压治理 积压
- 1. 积压怎么发现?2. 怎么快速消费?3. 消费者扩容?4. 紧急丢弃/转存?5. 预防?
第一步:分析。积压=生产 > 消费。治理=「提消费能力 + 临时扩容 + 降级非核心 + 修复根因」。
第二步:挑战。消费速度、扩容、下游扛不住、根因。
第三步:架构。①监控:lag 告警(Kafka consumer lag);②提能力:临时增加分区+消费者实例(Kafka 分区数=最大并行度);③降级:非核心消息跳过/转冷存后补;④修复:定位慢查询/锁;⑤紧急:抽样转存(如转对象存储)避免 broker 爆。
第四步:选型。Kafka(增分区+消费组扩容)+ 监控(Burrow/Prometheus)+ 降级开关。
第五步:一致性。积压消息仍按顺序(同分区);降级需保证可补。
第六步:高可用。消费扩容不冲击下游(限流);转存防 broker 满。
第七步:优化。批量消费+并行;修复根因防再积。
为什么增分区?Kafka 并行度受分区数限制,消费者数 ≤ 分区数才有效,故先增分区再扩实例。下游扛不住:消费快了但下游(DB)被打,需限流或批量。根因优先:不修 bug 只扩容量,积压会再来;先止血(扩容)再修根因。
// 批量消费提吞吐 @KafkaListener(topics="order", concurrency=32) // 并发=分区数 public void on(List<ConsumerRecord> batch){ db.batchInsert(batch); // 批量写 consumer.commitSync(); } // 监控: kafka-consumer-groups.sh --describe --lag
追问 1:分区数已最大?无法再扩并行;改批量/优化消费逻辑,或临时多消费组分治。
追问 2:下游 DB 被打?消费限流/批量+合并写;或削峰(再入 MQ 缓冲)。
追问 3:积压消息过期?设长 retention;过期丢失需补(DB 扫差)。
追问 4:非核心消息?降级开关跳过或转冷存,优先保核心。
追问 5:预防积压?容量评估+消费能力预留+lag 告警+压测。
- ❌ 只扩消费者不增分区——无效。正确:先增分区。
- ❌ 不修根因——再积。正确:止血+修因。
- ❌ 消费猛冲打垮下游——二次故障。正确:限流。
第 66 题:死信队列与重试机制 死信
- 1. 重试策略?2. 死信队列是什么?3. 重试阻塞怎么办?4. 死信怎么处理?5. 指数退避?
第一步:分析。失败消息需有限重试 + 隔离,避免阻塞正常消费。死信队列(DLQ)承接最终失败消息。
第二步:挑战。重试风暴、阻塞、死信处理。
第三步:架构。①重试:指数退避 + 上限(如 3~5 次),非业务错(4xx)不重试;②隔离:重试用独立重试队列/延迟重试,不占主消费;③死信:超上限进 DLQ,人工/定时补偿,不阻塞主队列;④告警:DLQ 有消息即告警。
第四步:选型。RocketMQ(重试队列+死信自动)/ RabbitMQ(DLX+ TTL)/ Kafka(自建重试 topic)。
第五步:一致性。死信保留消息可补;重试保证最终处理。
第六步:高可用。死信隔离保主队列流畅;补偿可重投。
第七步:优化。分类重试(瞬时 vs 永久错);死信人工台。
为什么死信?持续失败的消息若一直在主队列重试,占资源+阻塞同队列其他消息。进 DLQ 隔离,主队列继续。指数退避:1s→2s→4s 避免立即重试打爆下游。业务错不重试:参数错(4xx)重试无用,直接死信。
// RocketMQ: 重试(默认16次退避) + 死信(%DLQ%topic) @RocketMQMessageListener(topic="notify", consumerGroup="g") public void on(Msg m){ try { sms.send(m); } catch (RetryableException e) { throw e; } // 触发重试 catch (BizException e) { // 不重试 deadLetter.save(m); return; } } // 死信消费组: 告警+人工补偿
追问 1:死信一直堆积?告警+人工台;定时补偿重投(修复后)。
追问 2:重试退避太长?按错误类型:瞬时短退避,持久长退避/直死信。
追问 3:Kafka 没原生死信?自建重试 topic + DLQ topic,消费失败转发。
追问 4:重试幂等?重试消息需幂等(Q63),避免重复副作用。
追问 5:死信误判?区分可重试/不可重试;不可重试直死信,省资源。
- ❌ 失败无限重试——阻塞。正确:上限+死信。
- ❌ 业务错也重试——无用。正确:直死信。
- ❌ 死信不处理——丢消息。正确:告警+补偿。
第 67 题:延迟消息与定时任务 延迟
- 1. 延迟消息原理?2. 18 级不够怎么办?3. 任意精度方案?4. 关单与支付并发?5. 大量延迟消息性能?
第一步:分析。延迟消息=「到时间才消费」。精度需求决定方案:固定档用 MQ 延迟,任意精度用时间轮/DB 扫描/RocketMQ 定时。
第二步:挑战。精度、并发、性能、可靠。
第三步:架构。①固定档:RocketMQ 延迟消息(如 30 分钟档);②任意精度:自研时间轮(Netty HashedWheelTimer)+ DB 持久化,或 RocketMQ 5.0 定时消息(任意);③关单并发:消费时查订单状态(已付则跳过),幂等;④性能:延迟消息走独立 topic,批量。
第四步:选型。RocketMQ 定时消息(5.0+任意)/ 时间轮 + DB / 或 DB 扫描(低精度)。
第五步:一致性。关单前校验状态,已付不关;失败可重投。
第六步:高可用。延迟消息持久化;时间轮故障可重建(DB 兜底)。
第七步:优化。时间轮 O(1) 触发;DB 扫描低峰批量。
MQ 延迟原理:消息暂存内部 Schedule 主题,到点才转正常队列。18 级局限:仅固定档(1s/5s/.../2h),「37 分钟」无法表达。RocketMQ 5.0 支持任意定时。时间轮:Hash 分桶,O(1) 入/出,适合海量定时;需 DB 持久化防丢。状态校验:延迟消息到时订单可能已付,必须查状态再关。
// RocketMQ 5.0 任意定时消息 Message m = new Message("order-delay", body); m.setDelayTimeSec(30*60); // 任意秒 producer.send(m); // 消费: 状态校验(幂等关单) public void on(OrderMsg m){ Order o = orderMapper.get(m.orderId); if (o.status == UNPAID) orderMapper.close(m.orderId); // 已付跳过 }
追问 1:时间轮故障丢定时?DB 持久化定时任务,时间轮重启重建;或 RocketMQ 定时(持久)。
追问 2:百万定时性能?时间轮 O(1);分片多时间轮实例;DB 扫描低峰批量。
追问 3:关单与支付并发?状态机(UNPAID→PAID 与 UNPAID→CLOSED 互斥),后到者失败。
追问 4:RocketMQ 延迟精度?5.0 前 18 级固定;5.0 任意;旧版用时间轮补。
追问 5:DB 扫描方案?订单加 expire_time 索引,定时扫未付且过期;简单但实时性差。
- ❌ 全表定时扫——慢且锁。正确:延迟消息/时间轮。
- ❌ 关单不校验状态——误关已付。正确:状态校验。
- ❌ 18 级当任意精度——不准。正确:时间轮/5.0定时。
第 68 题:Kafka 高性能调优 调优
- 1. 吞吐瓶颈在哪?2. 生产者怎么调?3. 服务端怎么调?4. 消费怎么调?5. 零拷贝?
第一步:分析。Kafka 高性能=顺序写 + 零拷贝 + 批量 + 压缩。瓶颈常在未批量/未压缩/单分区/磁盘。
第二步:挑战。批量、压缩、分区、磁盘 IO。
第三步:架构。①生产:批量(batch.size=64k, linger.ms=5)、压缩(snappy/lz4)、acks=1(非关键);②服务端:多分区、顺序写(SSD)、页缓存;③消费:批量拉(max.poll.records)、并行;④零拷贝:sendfile 消费传输。
第四步:选型。Kafka 参数调优 + 分区规划 + SSD。
第五步:一致性。acks=1 降可靠换吞吐;关键用 all。
第六步:高可用。多副本;调优不破坏副本。
第七步:优化。压测定位瓶颈(生产/网络/磁盘/消费)。
批量+linger:攒批发降低网络往返与开销,linger.ms 等一小会凑批。压缩:snappy/lz4 降网络与磁盘。顺序写:Kafka 追加写日志,磁带式顺序 IO 极快;页缓存由 OS 管理。零拷贝:消费时 sendfile 直接内核→socket,免用户态拷贝。
// 生产者调优 props.put("batch.size", "65536"); props.put("linger.ms", "5"); props.put("compression.type", "lz4"); props.put("acks", "1"); // 消费: 批量拉 props.put("max.poll.records", "500"); // 服务端: num.partitions 多 + log.segment 调优 + SSD
追问 1:acks=1 丢?Leader 落即确认,副本异步;丢概率低,关键链路用 all。
追问 2:linger 增延迟?5~20ms 换吞吐,日志场景可接受。
追问 3:分区太多?分区多→文件多→句柄/内存涨;按吞吐需求定,非越多越好。
追问 4:磁盘瓶颈?顺序写 SSD/多磁盘;监控 IO wait。
追问 5:消费跟不上?增分区+并发消费+批量+优化处理逻辑。
- ❌ 单条发——吞吐低。正确:批量+linger。
- ❌ 不压缩——网络满。正确:lz4。
- ❌ 分区太少——并行受限。正确:多分区。
第 69 题:事务消息与本地事务一致 事务消息
- 1. 为什么普通发消息不原子?2. 事务消息原理?3. 半消息与回查?4. 回查失败?5. 替代方案?
第一步:分析。本地 DB 事务与发 MQ 是两个操作,无法原子。事务消息(半消息+回查)保证「本地成则消息发,本地败则消息不发」。
第二步:挑战。原子性、回查、可靠性。
第三步:架构。①半消息:先发半消息(对消费者不可见);②本地事务:执行 DB 写;③提交/回滚:成则 commit(半消息变可见),败则 rollback;④回查:Broker 未收到 commit/rollback 时回调生产者查本地状态;⑤替代:本地消息表(见第二部分 Outbox)。
第四步:选型。RocketMQ 事务消息(半消息+回查)/ 本地消息表(Kafka)。
第五步:一致性。事务消息保证最终一致;回查兜底防悬挂。
第六步:高可用。回查保证中间状态可决断;消费者幂等。
第七步:优化。回查限频;幂等消费。
为什么半消息?若先发消息再写库,写库失败但消息已出→不一致。半消息先占坑,本地事务成再「转正」。回查关键:网络丢 commit/rollback 时,Broker 不知状态,回调生产者查 DB(订单在=提交,不在=回滚),防消息悬挂。替代:本地消息表(同库写消息,CDC 发),Kafka 无原生事务消息时用。
// RocketMQ 事务消息 TransactionSendResult r = txProducer.sendMessageInTransaction( new Message("order", body), orderId); // 本地事务 public LocalTransactionState executeLocal(Message m, Object o){ return orderMapper.insert((Order)o) ? COMMIT : ROLLBACK; } // 回查: 查DB订单是否存在 public LocalTransactionState checkLocal(Message m, Object o){ return orderMapper.exists(orderId) ? COMMIT : ROLLBACK; }
追问 1:回查也失败?限频重试+告警;最终可人工决断;消息不丢(半消息在)。
追问 2:Kafka 无事务消息?用本地消息表(同库写消息记录,CDC/定时发)。
追问 3:消息重复发?半消息仅一条;commit 后消费者幂等。
追问 4:本地事务很长?事务消息要求本地事务短;长事务用 Saga/本地消息表异步。
追问 5:多下游?一条事务消息,各下游独立消费+幂等,互不影响。
- ❌ 先发消息再写库——不一致。正确:半消息+回查。
- ❌ 不回查——悬挂。正确:回查决断。
- ❌ 消费者不幂等——重复。正确:幂等。
第 70 题:消息轨迹与可观测性 可观测
- 1. 观测看什么?2. 消息轨迹怎么实现?3. TraceID 透传?4. 积压/延迟监控?5. 告警?
第一步:分析。MQ 可观测=「生产成功?Broker 收到?消费耗时?积压多少?」全链路指标+轨迹。
第二步:挑战。轨迹存储、Trace 透传、指标。
第三步:架构。①指标:生产 TPS、ack 率、Broker lag、消费 TPS/耗时、重试/死信数(Prometheus);②轨迹:消息带 msgId+TraceID,生产/消费各记一条(RocketMQ 消息轨迹/自建);③Trace 透传:TraceID 随消息头传下游,串联链路;④告警:lag>阈值、消费失败率、死信。
第四步:选型。RocketMQ 轨迹 / Kafka+Prometheus / OpenTelemetry + Jaeger。
第五步:一致性。轨迹与业务解耦,异步记录不阻塞。
第六步:高可用。轨迹存储独立;不影响主链路。
第七步:优化。采样轨迹降成本;关键消息全量。
为什么轨迹?消息异步,出问题难定位「丢在哪一环」。轨迹记录每环时间+状态,秒级定位。TraceID 透传:把请求 TraceID 放进消息头,消费时取出继续透传,MQ 成为链路一环(OTel 支持)。lag 是关键指标:直接反映消费健康。
// 发消息带 TraceID Message m = new Message(topic, body); m.putUserProperty("traceId", Tracing.current()); // 消费取出继续透传 public void on(Msg m){ Tracing.start(m.getUserProperty("traceId")); // 串联链路 handle(m); } // 指标: Prometheus 暴露 lag/耗时, Grafana 告警
追问 1:轨迹存储成本?采样(如 10%);关键消息全量;存 ES/时序库。
追问 2:Trace 跨 MQ 断?消息头透传 TraceID + Baggage,消费方接续。
追问 3:Kafka 没原生轨迹?自建:生产/消费埋点上报;或用 Kafka 监控+业务日志关联 msgId。
追问 4:告警怎么设?lag 持续涨、消费失败率、死信数、生产 ack 率。
追问 5:定位「丢了」?轨迹显示「生产成功但消费无记录」=消费丢;「生产无记录」=生产丢。
- ❌ 无轨迹——排查盲。正确:全链路轨迹。
- ❌ Trace 不进 MQ——链路断。正确:头透传。
- ❌ 不监控 lag——积压无感。正确:lag 告警。