面试专题

Kafka 消息积压如何处理

从 Consumer Lag 识别、瓶颈定位、限流止血、消费者扩容、批量处理、热点分区和毒消息治理等方面,系统讲解 Kafka 消息积压处理方案。

难度:深入更新:2026-07-24

Kafka 消息积压如何处理

面试回答

Kafka 消息积压不能一上来就盲目增加消费者,首先要确认积压范围和原因:观察各 Partition 的 Consumer Lag、生产速率、消费速率、消费者存活状态、重平衡、处理耗时和下游依赖。如果消费实例异常或下游故障,应先恢复服务;如果生产速度持续大于消费速度,需要对非核心生产流量限流,同时扩容消费者,但同一消费组的有效并行度不会超过 Partition 数量。如果单条处理过慢,应通过批量写库、减少同步 RPC、优化慢 SQL 或使用受控的异步处理提高单实例吞吐;如果只有个别 Partition 积压,要检查消息 Key 是否造成热点;如果被异常消息反复阻塞,需要有限重试后进入死信或人工补偿。恢复时间可以用“积压量 ÷(消费速率-生产速率)”估算,只有消费能力大于生产速度,积压才会真正下降。Offset 跳到最新只能在业务明确同意丢弃旧消息时使用,不能作为常规解决方案。

一句话总结:

先判断为什么消费速度落后,再从减少流入、恢复异常、提升单实例吞吐和增加有效并行度四个方向处理。

详细讲解

Kafka 消息积压的本质是:

一段时间内的生产速度
>
一段时间内的有效消费速度

积压不是单一故障类型,它可能由突发流量、消费者异常、下游变慢、分区热点或异常消息阻塞引起。不同原因的处理方式完全不同。

Kafka 消息积压排查与处理流程

一、什么是 Kafka 消息积压

每个 Partition 中的消息都有 Offset。生产者持续向日志末尾写入,消费组通过已提交 Offset 记录自己的处理进度。

监控消费组时,Consumer Lag 通常可以近似理解为:

Consumer Lag
= Partition Log End Offset
- Consumer Group 已提交 Offset

例如:

Partition Log End Offset:120000
消费组已提交 Offset:95000

Lag = 120000 - 95000 = 25000

表示这个消费组在该 Partition 上大约落后 25000 个 Offset。

需要同时观察三类指标:

指标含义判断价值
当前 Lag还落后多少 Offset反映积压规模
Lag 增长速度单位时间新增多少积压判断问题是否正在恶化
最老未处理消息年龄最老待处理消息的时间戳距现在多久反映业务已经延迟多久
积压字节量待处理消息占用的总字节数反映网络、磁盘和反序列化压力
单条处理耗时处理一条或一批消息实际需要多久用于估算积压清空速度

只看 Lag 绝对值容易误判。例如同样积压 10 万条,在每秒消费 20 万条的系统里可能很快恢复,在每秒只消费 100 条的系统里则非常严重。

Lag 是 Offset 数量差,本身只表示“落后多少条记录”,不能直接等同于时间延迟。对于相同的 10 万 Lag,消费者处理速度越快,处理完现有积压所需的时间越短:

处理现有积压所需时间
≈ Lag ÷ 消费速度

如果生产者仍在持续写入,要估算 Lag 整体降为 0 的时间,则应扣除生产速度:

Lag 归零时间
≈ Lag ÷(消费速度-生产速度)

只有消费速度大于生产速度,Lag 才会持续下降。至于当前业务已经延迟多久,应直接查看最老未处理消息的时间戳,而不是只根据 Lag 推算。

消息大小影响的是积压字节量、网络传输和反序列化成本;单条处理复杂度影响的是消费速度。它们都很重要,但不能与 Lag 条数或消息年龄混为同一个指标。

二、先判断积压发生在哪里

1. 是所有 Partition 积压,还是少数 Partition 积压

所有 Partition 的 Lag 都增长
→ 更可能是整体消费能力不足、消费者异常或下游故障

只有少数 Partition 的 Lag 很高
→ 更可能是 Key 倾斜、热点业务或个别异常消息

不能只看 Topic 的总 Lag。总数会掩盖单个热点 Partition,应同时查看每个 Partition 的 Lag、生产速率和消费速率。

2. 消费者是否正常工作

先检查:

  • 消费实例是否存活,是否频繁重启。
  • Consumer Group 中是否存在空闲或反复加入、退出的成员。
  • 是否频繁发生重平衡。
  • poll() 是否长时间没有被调用。
  • Offset 是否持续推进。
  • 消费日志中是否存在反序列化、鉴权、提交 Offset 或业务异常。

如果业务处理时间超过 max.poll.interval.ms,消费者会被认为没有正常推进消费,Partition 可能被重新分配。消费者不断处理、超时、重平衡,又可能造成更严重的积压。

3. 单条消息处理是否变慢

Kafka 消费本身往往不是最慢的部分,真正瓶颈通常在业务逻辑:

  • 数据库慢 SQL、锁等待或连接池耗尽。
  • 调用外部 RPC 超时。
  • 每条消息单独写库,缺少批处理。
  • 序列化、解压或大对象处理消耗 CPU。
  • 消费逻辑加了粒度过大的锁。
  • 下游 Elasticsearch、Redis 或第三方接口限流。

应把消费耗时拆分为:

拉取耗时
+ 反序列化耗时
+ 业务计算耗时
+ 数据库 / RPC / 下游写入耗时
+ Offset 提交耗时

只有找到主要耗时,优化才有方向。

4. 是否被“毒消息”阻塞

毒消息是指某条数据因为格式错误、业务状态异常或下游约束问题,每次消费都会失败。

如果代码采用无限重试:

消费失败
→ 立即重试
→ 再次失败
→ 一直占住当前 Partition

该 Partition 后续所有消息都可能无法继续处理,表现为单分区 Lag 持续增大。

5. Broker 或基础设施是否成为瓶颈

还要检查:

  • Broker 磁盘延迟、网络带宽和 CPU。
  • Fetch 请求延迟和错误率。
  • ISR 收缩、离线分区和副本同步异常。
  • 消费者所在机器的 CPU、内存、GC 和网络。
  • 是否发生跨机房访问或网络抖动。

如果 Broker 或网络已经饱和,单纯增加 Consumer 可能只会制造更多请求和竞争。

三、紧急处理:先止住 Lag 增长

1. 恢复异常消费者和下游

如果积压源于实例宕机、数据库故障或下游超时,应优先恢复根因:

恢复消费实例
→ 恢复数据库 / RPC / 下游
→ 确认 Offset 开始推进
→ 再评估是否需要扩容

下游仍不可用时强行增加消费者,只会放大连接数、重试和超时流量。

2. 对非核心生产流量限流或降级

如果生产速度仍在持续上升,而消费端已经满负荷,可以在业务允许时:

  • 限制非核心事件的生产速率。
  • 暂停可延迟的定时任务和批量导入。
  • 合并低价值事件。
  • 对非核心功能降级。

这不是最终解决方案,但能降低流入速度,给消费端争取恢复窗口。

这里的“合并低价值事件”,不是随意删除消息,而是对业务允许只保留最终状态或聚合结果的事件进行压缩。例如:

同一个商品在 1 秒内连续产生 20 次缓存刷新通知
→ 只保留一次刷新通知

同一设备连续上报多次进度:10%、20%、30%
→ 如果业务只关心最新进度,可以只发送 30%

大量指标增量事件
→ 在生产端按时间窗口汇总后,发送一条聚合结果

合并前必须先定义清楚业务语义,包括按什么 Key 合并、合并时间窗口多长、保留最新值还是累计值,以及异常时如何补偿。订单状态流转、支付、账户流水、库存扣减、审计日志等每条事件都具有独立业务意义,不能使用这种方式合并。

3. 确认消息保留时间是否足够

积压消息依赖 Kafka 的日志保留策略。如果预计追赶需要数小时或数天,应确认 Topic 的保留时间和磁盘空间是否足够。

如果最老未消费消息超过保留期限,旧日志可能被清理。此时即使消费者恢复,也无法从 Kafka 读取已经删除的数据。

紧急情况下可以评估临时延长保留时间,但必须同步评估 Broker 磁盘容量,避免积压尚未解决又触发磁盘告警。

四、提升消费能力

1. 增加 Consumer 实例

增加同一 Consumer Group 中的实例,可以让更多 Partition 并行消费。

但有效并行度受 Partition 数量限制:

8 个 Partition + 4 个 Consumer
→ 可以继续扩容

8 个 Partition + 8 个 Consumer
→ 已接近分区级并行上限

8 个 Partition + 12 个 Consumer
→ 约 4 个 Consumer 无 Partition 可处理

因此扩容前要同时查看当前 Partition 数量、消费者数量和分配情况。

2. 优化单实例处理吞吐

常见优化方式:

  • 把逐条写数据库改为批量写入。
  • 合并相同业务 Key 的重复更新。
  • 优化 SQL、索引和事务范围。
  • 减少消费线程中的同步远程调用。
  • 对可并行步骤使用有界线程池。
  • 复用连接和客户端,避免每条消息重复创建。
  • 降低无必要的日志和对象序列化开销。

例如:

原方案:
每条消息 → 一次数据库事务

优化后:
一批消息 → 分组校验 → 批量写入 → 提交对应 Offset

批量越大不一定越好。批次过大会增加内存占用、事务时间和失败重试成本,还可能让一次处理时间超过 max.poll.interval.ms

这里的“同步远程调用”,是指消费线程发送 RPC 或 HTTP 请求后,必须等待对方返回才能继续处理下一条消息:

消费一条消息
→ 调用远程库存服务
→ 阻塞等待 50 ms
→ 收到结果
→ 再处理下一条消息

即使本地计算只需要 1 ms,线程的大部分时间也消耗在网络等待上。如果每条消息都串行调用一次远程服务,单线程理论吞吐很容易被远程调用耗时限制。

“减少同步 RPC”可以从以下方向处理:

  • 把逐条调用改为批量接口,例如一次查询或提交 100 条。
  • 对可复用且允许短时间不一致的数据使用本地缓存,减少重复查询。
  • 把互不依赖的调用改成受控并发,但必须限制并发数,避免压垮下游。
  • 将通知、埋点等非核心副作用发送到下游 Topic,由独立消费者异步执行。
  • 合理设置超时、熔断和连接池,避免故障调用长期占住消费线程。

如果必须拿到远程结果才能保证当前消息处理正确,就不能为了吞吐简单删除该 RPC。此时更重要的是批量化、限制并发、幂等和下游容量治理。

3. 谨慎调整消费参数

常见参数包括:

max.poll.records=500
max.poll.interval.ms=300000
fetch.min.bytes=1
fetch.max.wait.ms=500
max.partition.fetch.bytes=1048576

调整原则:

  • 如果一批消息处理不完,应适当减小 max.poll.records,或优化处理速度,而不是只把 max.poll.interval.ms 调得无限大。
  • 增大 fetch.min.bytes 可以提高批量拉取效率,但会增加低流量时的等待延迟。
  • max.partition.fetch.bytes 过小可能限制大批次拉取,但调大后要评估客户端内存。
  • 参数优化只能降低 Kafka 拉取开销,无法解决数据库或 RPC 本身处理缓慢。

4. 异步处理必须有边界

一种常见设计是由 Poll 线程负责拉取,再把消息交给工作线程池处理:

Poll 线程
→ 有界队列
→ 工作线程池
→ 处理完成
→ 按 Partition 推进可安全提交的 Offset

为什么必须使用有界队列?假设 Kafka 已经积压 100 万条,Poll 线程仍然高速拉取并塞入无界队列,而工作线程每秒只能处理 100 条,那么积压只是从 Kafka 的磁盘日志转移到了 JVM 堆内存。结果可能是:

队列持续增长
→ 堆内存占用上升
→ Full GC 频繁
→ 消费进程响应变慢甚至 OOM
→ 未提交的内存任务随进程退出而丢失并被重新投递

Kafka 本身就是持久化队列,没有必要把尚未具备处理能力的消息提前搬进易失的 JVM 内存。

使用高、低水位控制 pause()resume()

可以为队列设置容量和两个阈值。例如队列容量为 10000:

高水位:8000
队列达到高水位
→ Poll 线程对相关 Partition 调用 pause()
→ 后续 poll() 暂停返回这些 Partition 的新消息

低水位:3000
工作线程处理后,队列下降到低水位
→ Poll 线程调用 resume()
→ 后续 poll() 重新拉取这些 Partition

使用两个不同阈值是为了避免队列在一个临界值附近反复暂停、恢复。阈值不能照搬固定数字,应根据单条消息大小、堆内存、工作线程数和峰值处理耗时进行容量评估。

pause() 的含义是暂停指定 Partition 在后续 poll() 中返回新记录,它不会:

  • 取消当前 Partition 的分配。
  • 自动提交 Offset。
  • 主动触发 Consumer Group 重平衡。
  • 清除已经拉取并放入本地队列的消息。

暂停期间仍然必须持续调用 poll(),让 Consumer 处理心跳、协调和重平衡事件。不能因为所有分区都已暂停,就让 Poll 线程长时间睡眠或等待队列清空,否则仍可能超过 max.poll.interval.ms 并失去分区。

只能由 Consumer 所在线程控制拉取

KafkaConsumer 不是线程安全的。通常由唯一的 Poll 线程持有并调用 poll()pause()resume() 和提交 Offset;工作线程只处理业务,并把“完成结果”写回一个线程安全的完成队列。工作线程不能直接操作同一个 Consumer。

伪代码可以理解为:

Poll 线程循环:
    读取工作线程上报的完成结果
    推进每个 Partition 的连续成功 Offset
    根据队列水位执行 pause / resume
    poll() 拉取允许继续消费的 Partition
    把记录投递到对应的有界工作队列

工作线程:
    处理业务
    成功或失败结果写回完成队列
    不直接调用 KafkaConsumer

为什么要按 Partition 提交“连续成功区间”

假设同一 Partition 拉取到:

Offset 100、101、102

并发处理结果是:

100 成功
101 仍在处理
102 成功

此时最多只能提交 101,也就是声明 100 已处理完成;不能提交 103。否则进程在 101 完成前崩溃,重启后会从 103 开始,101 就可能被永久跳过。

只有 101 也成功后,100~102 形成连续成功区间,才能提交 103。若业务要求同一 Partition 严格有序,最简单的方式是同一 Partition 串行处理,不要并发执行。

重平衡时还要处理本地未完成任务

发生 Partition 撤销时,应停止向被撤销 Partition 的队列继续投递,并在回调允许的时间内提交已经连续处理成功的 Offset。尚未完成或来不及安全提交的任务,应允许新 Consumer 从旧 Offset 重新消费,因此业务处理必须具备幂等性。

需要共同遵守的约束包括:

  • KafkaConsumer 本身不是线程安全的,不能让多个工作线程同时操作同一个 Consumer。
  • 队列必须有界,不能把 Kafka 的积压无上限转移到 JVM 内存。
  • pause()resume() 应由 Poll 线程执行,工作线程通过控制信号提出请求。
  • 即使 Partition 已暂停,Poll 线程也要持续执行 poll()
  • Offset 必须在消息真正处理成功后提交。
  • 同一 Partition 如果要求严格顺序,不能让后面的消息先于前面的消息产生业务效果。
  • 并发完成后只能提交“连续成功区间”的下一个 Offset,不能跨过尚未完成或失败的消息。

异步化能提高吞吐,但会显著增加 Offset、顺序和失败处理的复杂度。

五、Partition 不够时怎么办

1. 增加 Partition 主要提升未来并行度

如果 Consumer 数量已经等于 Partition 数量,且单分区处理能力也接近上限,可以评估增加 Partition。

但要明确:

给 Topic 新增 Partition,不会把旧 Partition 中已经积压的消息自动重新分布到新 Partition。

新 Partition 主要承接之后产生的新消息。旧积压仍然留在原来的 Partition 中,仍需由对应消费者追赶。

2. 增加 Partition 可能改变 Key 路由

默认的 Key 哈希结果通常与 Partition 数量有关。增加 Partition 后,同一个 Key 的新消息可能进入不同于历史消息的 Partition,从而影响跨扩容时点的顺序假设。

例如为了便于理解,假设某个分区规则可以简化为:

目标分区 = hash(key) % Partition 数量

同一个 Key 的哈希值为 13:

原来 4 个 Partition:13 % 4 = 1,写入 P1
扩容到 8 个 Partition:13 % 8 = 5,新消息写入 P5

扩容前已经写入 P1 的旧消息不会被搬到 P5。P1 和 P5 又可能由不同 Consumer 并行处理,因此 P5 中的新消息可能先于 P1 中的旧消息完成。这个风险不只是“扩容瞬间”存在,而是会持续到旧分区中的相关历史消息全部处理完毕。

所谓“消费端能否接受跨分区处理”,实际是在问:

同一个业务 Key 的旧消息和新消息同时位于不同 Partition,并可能并行、乱序完成时,业务是否仍然正确?

如果事件是独立日志或幂等的最终状态覆盖,业务可能可以接受;如果事件依赖严格顺序,例如“创建订单 → 支付成功 → 关闭订单”,就不能直接接受。

如果业务依赖同一 Key 的严格顺序,常见选择有:

方案一:排空旧消息后再扩分区

限制或暂停生产
→ 等待旧 Topic 的相关消息处理完成
→ 增加 Partition
→ 恢复生产

这种方式逻辑最清楚,但需要维护窗口,而且在高流量系统中可能很难完全停写。

方案二:新建 Topic 并进行受控切换

例如创建具有目标分区数和新路由规则的 order-events-v2

1. 创建并验证新 Topic
2. 定义明确的切换时间、版本号或业务水位
3. 新消息从切换点开始写入 v2
4. v1 消费者继续处理切换点以前的旧消息
5. 确认某个 Key 或全部旧消息越过安全水位后,再允许 v2 对应消息生效
6. 完成核对后停止 v1,并保留回滚方案

新建 Topic 的价值是把新旧路由、配置和消费组隔离开,便于验证与回滚;代价是必须设计双 Topic 切换、重复消息、遗漏检查和新旧 Offset 衔接。它不是 Kafka 自动迁移,而是业务自己完成的版本化切换。

因此扩分区前应明确回答:

  • 分区器如何计算目标 Partition。
  • 旧消息是否已经处理完成。
  • 同一 Key 跨新旧 Partition 并行处理是否会破坏业务顺序。
  • 是选择维护窗口排空旧消息,还是新建 Topic 做受控切换。

3. 极端情况下建立临时加速链路

当原 Topic 分区数严重不足、积压量巨大时,可以设计经过评审的临时方案:

原 Topic
→ 临时搬运程序
→ 更高分区数的中转 Topic
→ 临时扩容后的消费组
→ 业务处理

这会引入顺序、重复、Offset 衔接和回切问题,只适合经过完整设计和演练的场景,不能在生产事故中临时拍脑袋执行。

六、热点 Partition 如何处理

如果只有少数 Partition 积压,需要检查消息 Key 分布。

例如所有消息都使用同一个固定 Key:

key = "default"
→ 所有消息进入同一 Partition
→ 其他 Partition 空闲
→ 单分区成为瓶颈

解决方向包括:

  • 选择分布更均匀的业务 Key。
  • 对非顺序场景取消固定 Key。
  • 对热点业务 Key 做可控分片,例如 userId + shardNo
  • 单独拆分热点租户或热点业务到独立 Topic。
  • 优化该 Partition 对应消费者的单条处理耗时。

Key 拆分会改变顺序语义。原来同一 Key 的全局顺序,拆分后通常只能保证每个子 Key 内部有序。

七、毒消息和重试如何处理

不要让一条永远失败的消息无限阻塞 Partition。常见处理策略是:

第一次失败
→ 记录错误并有限重试
→ 仍然失败
→ 进入重试 Topic / 死信 Topic
→ 主消费链路继续
→ 告警和人工补偿

需要保留:

  • 原始 Topic、Partition 和 Offset。
  • 消息 Key、事件 ID 和原始内容。
  • 异常原因和重试次数。
  • 首次失败和最后失败时间。

如果业务要求同一 Key 严格顺序,就不能简单跳过失败消息继续处理后续消息。此时应暂停对应 Partition,优先修复或补偿该消息;这是吞吐与顺序之间的业务取舍。

八、如何估算多久能消化完积压

设:

B = 当前积压消息数
P = 当前生产速度(条/秒)
C = 扩容后的稳定消费速度(条/秒)

只有:

C > P

积压才会下降。理论恢复时间约为:

恢复时间 T = B / (C - P)

例如:

当前积压 B = 360 万条
生产速度 P = 2000 条/秒
消费速度 C = 5000 条/秒

净消化速度 = 5000 - 2000 = 3000 条/秒
理论恢复时间 = 3600000 / 3000 = 1200 秒
≈ 20 分钟

实际还要为重平衡、失败重试、下游波动和流量峰值预留余量。

如果 C <= P,无论运行多久都追不上,必须降低生产速率或提高有效消费能力。

九、Offset 跳到最新为什么不是常规方案

把消费组 Offset 重置到最新位置,确实可以让 Lag 迅速变成 0,但其本质是:

主动放弃尚未消费的历史消息。

只有同时满足以下条件才可以考虑:

  • 业务明确确认旧消息可以丢弃。
  • 已评估对账、状态和下游影响。
  • 已保留必要的审计或补偿数据。
  • 操作前进行了 Dry Run,并确认 Topic、Consumer Group 和目标 Offset。
  • 操作期间停止相关消费者,避免 Offset 并发变化。

资金、订单、库存等关键链路不能把重置 Offset 当作普通清积压手段。

什么是 Dry Run

Dry Run 就是“只预演,不执行”。它先计算并展示每个 Partition 将从哪个 Offset 调整到哪个 Offset,供操作人员核对 Topic、消费组、分区范围和目标位置,但不真正修改已提交 Offset。

例如将消费组预演重置到最新位置:

bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka.example.com:9092 \
  --group order-consumer-group \
  --topic order-events \
  --reset-offsets \
  --to-latest

Kafka 4.1 的消费组工具默认只展示重置计划;确认结果无误后,显式增加 --execute 才会真正执行:

bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka.example.com:9092 \
  --group order-consumer-group \
  --topic order-events \
  --reset-offsets \
  --to-latest \
  --execute

不同 Kafka 版本的工具参数可能略有差异,生产操作前应以实际部署版本的命令帮助和官方文档为准。

为什么操作期间要停止消费者

活动消费者会继续执行 poll()、处理消息并提交 Offset。如果管理人员同时重置 Offset,就可能发生竞态:

管理员把 Offset 从 1000 重置到 5000
→ 仍在运行的消费者随后提交自己内存中的旧进度 1200
→ 重置结果被旧提交覆盖

反过来也可能是消费者刚处理到 1300,管理员基于稍早看到的状态执行重置,导致已处理消息重复消费或尚未处理消息被跳过。活动成员还可能触发重平衡,使操作时看到的分区归属和 Offset 继续变化。

因此安全流程应是:

停止该消费组的所有消费者
→ 确认消费组已无活动成员
→ 记录操作前 Offset
→ 执行 Dry Run 并双人核对
→ 使用 --execute 修改 Offset
→ 再次查询并确认结果
→ 启动消费者并观察消费位置

这不仅是避免并发覆盖。Kafka 官方的消费组 Offset 重置说明也明确要求先确保消费实例处于非活动状态。

十、恢复后要做什么

积压恢复不代表问题结束,还应完成:

  1. 复盘根因和时间线。
  2. 补充总 Lag、最大分区 Lag、时间 Lag 和增长率告警。
  3. 建立生产速度、消费速度和预计恢复时间看板。
  4. 对热点 Partition、毒消息和下游耗时建立专项监控。
  5. 按峰值流量进行容量压测,保留消费冗余。
  6. 演练消费者扩容、下游故障、重平衡和死信补偿。
  7. 核对积压期间是否发生消息过期、重复处理或业务遗漏。

十一、常见误区

误区一:积压就直接增加 Consumer

如果 Partition 数量不足、多数 Consumer 已空闲,或者瓶颈实际在数据库,增加 Consumer 不会解决问题,反而可能加剧下游压力。

误区二:把 max.poll.records 调大就能提升吞吐

拉取更多消息不等于处理更快。批次过大还可能导致处理超时、重平衡和更高的失败重试成本。

误区三:增加 Partition 会重新分配旧积压

新增 Partition 不会迁移旧 Partition 中已有的消息,只能为后续消息提供更多并行空间。

误区四:Lag 等于业务延迟

Lag 是 Offset 数量差,只表示还落后多少条记录,不等于时间延迟。相同 Lag 下,消费者处理速度越快,处理完现有积压所需的时间越短;如果生产仍在继续,则应使用 Lag ÷(消费速度-生产速度) 估算 Lag 归零时间。

当前业务已经延迟多久,应通过最老未处理消息的时间戳判断。消息大小影响积压字节量、网络与反序列化成本,单条处理耗时影响消费速度。因此排查时应把 Lag 条数、最老消息年龄、积压字节量和实际处理吞吐分开观察,不能混成一个指标。

误区五:把 Offset 跳到最新就是清理积压

这不是“处理了积压”,而是“丢弃了积压”。必须经过明确的业务授权和风险评估。

参考资料

DISCUSSION

评论与讨论

留下你的想法