Kafka 大消息处理

系统解决 Kafka 大消息难题:生产端与 broker 上限参数的成组放大、内存与副本同步的连锁代价、压缩与批量的覆盖边界、claim check 外存引用模式、分片与消费端重组、背压控制、对象生命周期与死信处理、方案选型对比、参数一致性排查与排错

Kafka 的设计假设是「消息小、数量多」:默认单条上限 1 MB,整个存储与网络栈都为高吞吐的小记录优化。但现实里总有例外——一张图片、一段 JSON 快照、一份 PDF 报表,动辄几 MB 到几十 MB。把上限一路调大能让它「跑起来」,但代价会沿着内存、网络、副本同步、消费延迟一路传导。大消息的正确解法通常不是「把上限调大」,而是换一种承载方式。

1. 为什么 Kafka 不适合大消息

1.1 上限参数的层次

「Kafka 的消息上限」不是一个参数,而是一组参数,任何一环不够大,链路就断。而且不同环节的失败表现完全不同:

环节参数默认值超限表现
生产端请求max.request.size1 MB客户端直接抛 RecordTooLargeException
broker 接收message.max.bytes~1 MBbroker 返回 RecordTooLargeException
副本同步replica.fetch.max.bytes1 MBfollower 拉不动,副本永远落后
消费端拉取max.partition.fetch.bytes1 MB消费者卡死,poll 拿不到数据

最后一行的表现最隐蔽:不报错,只是不动。因为 max.partition.fetch.bytes 是「单次 fetch 返回的上限」,若第一条消息就超过它,broker 会返回空响应,消费者反复 poll 反复拿空,看起来像「消费停了」而不是「消息太大」。这个死锁在大消息场景里极常见。

1.2 大消息的连锁代价

即便把参数都调大,代价仍然存在:

内存压力。 broker 处理一条消息时,它会在网络缓冲区、页缓存、副本同步缓冲区里各留一份。10 MB 的消息 × 100 个并发请求,瞬时内存就是数 GB。JVM 堆必须相应放大,而堆越大 GC 停顿越长——这与 Kafka 一直推荐的「小堆 + 大页缓存」原则直接冲突。

副本同步放大。 follower 要拉同样的数据。大消息让 replica.fetch.max.bytes 与网络带宽都吃紧,ISR 收缩,写入延迟上升。

分区倾斜与队头阻塞。 一个分区里混着 1 KB 和 10 MB 的消息,大消息会拖慢它所在分区的所有后续消费——因为同一个分区只能顺序消费。这在小消息场景不存在的问题,在大消息场景会显著恶化尾部延迟。

压缩失效。 已压缩的二进制(图片、PDF、zip)再压一遍几乎无收益,却照样消耗 CPU。

磁盘与保留期放大。 一条 10 MB 的消息在副本因子 3 下占 30 MB 裸容量。若这类消息占比高,磁盘消耗速度是消息数量增长的数十倍,保留期的成本估算会彻底失准。关于存储层的容量与留存策略,可以参考 Kafka 存储内核与日志压缩 。

2. 上限参数的成组放大

2.1 生产端

Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
// 单请求上限:必须 >= 最大消息 + 批次开销
props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 20 * 1024 * 1024);
// 单批上限:Kafka 会取 max(batch.size, 单条消息大小)
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 20 * 1024 * 1024);
// 缓冲总量:至少能装下几个大消息,否则 send 阻塞
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 256 * 1024 * 1024);

注意 batch.size 的作用:Kafka 的批次会为单条大消息自动扩容,所以 batch.size 小于消息大小时不会导致发送失败,但会导致「一个批次只装一条消息」,压缩与批量收益归零。

2.2 broker 端

# broker 接收上限
message.max.bytes=20971520
# 副本同步上限:必须 >= message.max.bytes
replica.fetch.max.bytes=20971520
# 每个分区返回给消费者的上限(broker 侧软限制)
fetch.max.bytes=52428800

replica.fetch.max.bytes 是最容易漏的一个。若它小于 message.max.bytes,follower 拉不到大消息,副本永远无法追平 leader,ISR 持续收缩,最终 leader 被隔离或分区不可用。

2.3 消费端

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
// 单分区单次返回上限:必须 >= 最大消息
props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 20 * 1024 * 1024);
// 单次请求总上限
props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 52_428_800);
// 单次 poll 的记录数(软上限,不切割批次)
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 50);

max.poll.records 在大消息场景要调小。原因是它不切割批次——若一个批次含 500 条大消息,一次 poll 就会返回全部,堆内存瞬间爆掉。调小到几十条,配合 max.partition.fetch.bytes 一起限制单次返回量。

2.4 参数必须成组放大的原因

这组参数的关系是一条约束链:

max.partition.fetch.bytes >= message.max.bytes
replica.fetch.max.bytes   >= message.max.bytes
max.request.size          >= 最大消息 + 开销
fetch.max.bytes           >= max.partition.fetch.bytes

任何一环小于 message.max.bytes,链路就断在那一环。运维上最常见的错误是只改了 topic 的 max.message.bytes(topic 级覆盖)而忘了 broker 的 replica.fetch.max.bytes,结果是生产者能发、消费者能收,但副本同步悄悄失败。

3. 压缩与批量的作用边界

3.1 压缩对大消息的效果

压缩对大消息的帮助取决于内容类型:

  • 文本类(JSON、XML、日志):压缩比 510 倍,效果显著,一条 10 MB 的 JSON 能压到 12 MB,直接回到默认上限内。
  • 已压缩二进制(JPEG、PNG、PDF、zip、gzip):压缩比接近 1.0,纯亏 CPU。compression-rate-avg 指标若接近或大于 1,说明压缩在帮倒忙。
# 文本类大消息:zstd 兼顾压缩比与速度
compression.type=zstd
# 二进制大消息:关掉压缩
# compression.type=none

压缩发生在批次级别,因此单条大消息的压缩与批量无关,只看内容本身。这一点与「小消息靠批量提升压缩比」不同。

3.2 批量的副作用

大消息场景下 linger.ms 应当调小甚至归零。原因是大消息本身已经足够大,攒批带来的压缩收益有限,而 linger 会额外增加延迟;更重要的是,多个大消息堆在一个批次里会让单次请求体积爆炸,触发 broker 的请求上限或消费者的内存上限。

# 大消息场景:不攒批,发出去就好
linger.ms=0

4. 外存引用模式(Claim Check)

4.1 模式本质

Claim Check(行李牌)模式把消息内容与消息本身解耦:大内容放对象存储,Kafka 里只传一个指针。

生产者                    对象存储               Kafka
  │                         │                    │
  ├── 上传大文件 ──────────► │                    │
  │   (S3/MinIO)            │                    │
  │                         │                    │
  ├── 发送指针消息 ─────────────────────────────► │
  │   {"bucket":"...", "key":"...", "size":...}  │
  │                                              │
消费者 ◄──────────────────────────────────────────┘
  │  收到指针
  ├── 按指针从对象存储拉取 ◄───┘
  └── 处理内容

Kafka 里的消息只有几百字节,所有上限参数都无需调整,broker 完全无感。这是大消息场景的首选方案,尤其是内容已经是「文件」形态时(图片、视频、报表)。

4.2 生产端实现

public class ClaimCheckProducer {
    private final KafkaProducer<String, String> producer;
    private final S3Client s3;
    private final ObjectMapper mapper = new ObjectMapper();

    public void send(String topic, String key, byte[] payload) throws Exception {
        String objectKey = "kafka-payloads/" + UUID.randomUUID() + ".bin";

        // 1. 先上传大内容,失败则整个操作失败,不会产生悬空指针
        s3.putObject(b -> b.bucket("app-blobs").key(objectKey)
                            .metadata(Map.of("sha256", sha256(payload))),
                     RequestBody.fromBytes(payload));

        // 2. 再发指针消息
        ClaimCheck ref = new ClaimCheck("app-blobs", objectKey,
                                        payload.length, sha256(payload));
        producer.send(new ProducerRecord<>(topic, key, mapper.writeValueAsString(ref)),
                      (meta, ex) -> {
                          if (ex != null) {
                              // 发送失败要清理已上传的对象,避免孤儿文件
                              s3.deleteObject(b -> b.bucket("app-blobs").key(objectKey));
                          }
                      });
    }
}

顺序很重要:先传内容,再发指针。 反过来(先发指针再传内容)会让消费者拿到指针时内容还不存在,产生竞态。代价是发送失败时可能留下孤儿对象——用一个清理任务扫「超过 N 小时未被引用的对象」来兜底,比在发送路径里做复杂补偿更可靠。

4.3 消费者端与生命周期

public void consume(ConsumerRecord<String, String> record) throws Exception {
    ClaimCheck ref = mapper.readValue(record.value(), ClaimCheck.class);

    try (ResponseInputStream<GetObjectResponse> in = s3.getObject(
            b -> b.bucket(ref.bucket()).key(ref.key()))) {
        byte[] payload = in.readAllBytes();

        // 校验完整性:对象存储可能被误删或覆盖
        if (!sha256(payload).equals(ref.sha256())) {
            throw new IntegrityException("内容哈希不匹配: " + ref.key());
        }
        process(payload);
    } catch (NoSuchKeyException e) {
        // 内容已被生命周期规则删除 —— 需要死信或人工介入
        deadLetter(record, e);
    }
}

生命周期管理是这个模式的核心运维点。对象存储上要配规则:多久后转冷存储、多久后删除。删除时机必须晚于所有消费者可能读到该指针的时间——即 ≥ Kafka topic 的保留期。若对象先删、消息还在,消费者就会撞上 NoSuchKeyException。

一个稳妥的做法是:对象存储的 TTL 设为 topic retention × 1.5,并给 NoSuchKeyException 配死信队列而非直接重试(重试无意义,对象不会自己回来)。

5. 分片与消费端重组

5.1 什么时候需要分片

Claim Check 适合「内容本身是文件」的场景。但如果业务要求内容必须走 Kafka(比如下游是 Kafka 生态里的流处理作业,不想引入对象存储依赖),就需要把大消息切成多个小片:

原始大消息 8 MB
   ├── chunk 0: {"msgId":"m1","seq":0,"total":4,"data":"..."}  2 MB
   ├── chunk 1: {"msgId":"m1","seq":1,"total":4,"data":"..."}  2 MB
   ├── chunk 2: {"msgId":"m1","seq":2,"total":4,"data":"..."}  2 MB
   └── chunk 3: {"msgId":"m1","seq":3,"total":4,"data":"..."}  2 MB

关键约束:同一 msgId 的所有分片必须落进同一分区,否则消费端无法保证收齐。做法是让分片消息的 key 都用 msgId(Kafka 按 key 哈希分区):

String msgId = UUID.randomUUID().toString();
int chunkSize = 1 * 1024 * 1024;   // 每片 1MB,安全落在默认上限内

for (int seq = 0; seq * chunkSize < payload.length; seq++) {
    int from = seq * chunkSize;
    int to = Math.min(from + chunkSize, payload.length);
    byte[] slice = Arrays.copyOfRange(payload, from, to);

    String envelope = mapper.writeValueAsString(new Chunk(msgId, seq,
        (payload.length + chunkSize - 1) / chunkSize, slice));

    // 同一 key → 同一分区 → 顺序保证
    producer.send(new ProducerRecord<>("large-payloads", msgId, envelope));
}

分片大小要留出余量:设成 1 MB 而 max.request.size 也是 1 MB 会失败,因为信封本身还有开销。取 max.request.size × 0.7 左右比较安全。

5.2 消费端重组与背压

重组需要一个按 msgId 聚合的缓冲,这正是背压的战场:

public class ChunkAssembler {
    // msgId -> 已收到的分片;必须有上限,否则内存会被慢消费拖垮
    private final Map<String, TreeMap<Integer, byte[]>> pending =
        new LinkedHashMap<>(16, 0.75f, true) {
            @Override
            protected boolean removeEldestEntry(Map.Entry<String, TreeMap<Integer, byte[]>> e) {
                return size() > MAX_PENDING;   // 超过 1000 个未完成消息就淘汰最旧的
            }
        };

    public Optional<byte[]> accept(Chunk chunk) {
        TreeMap<Integer, byte[]> parts =
            pending.computeIfAbsent(chunk.msgId(), k -> new TreeMap<>());
        parts.put(chunk.seq(), chunk.data());

        if (parts.size() < chunk.total()) {
            return Optional.empty();   // 还没收齐
        }
        pending.remove(chunk.msgId());

        ByteArrayOutputStream out = new ByteArrayOutputStream();
        parts.values().forEach(b -> out.writeBytes(b));
        return Optional.of(out.toByteArray());
    }
}

三个必须处理的边界:

  • 缓冲上限。若某个 msgId 的分片丢了一片(比如生产者中途崩溃),这条消息永远收不齐,缓冲会一直占着。必须有 TTL 或容量上限来淘汰,并把淘汰的消息送进死信。
  • 乱序到达。同分区内 Kafka 保证顺序,所以分片一定按 seq 递增到达。但跨分区或重平衡后可能出现乱序,用 TreeMap 按 seq 排序比 List 追加更稳。
  • 重复分片。at-least-once 语义下分片可能重复投递,TreeMap.put 天然幂等(同 seq 覆盖),这是它优于 List.add 的另一个理由。

5.3 重组超时与死信

// 定时清理超时未收齐的消息
scheduler.scheduleAtFixedRate(() -> {
    long cutoff = System.currentTimeMillis() - 60_000;
    pending.entrySet().removeIf(e -> e.getValue().isEmpty() /* 或按时间戳 */);
}, 30, 30, TimeUnit.SECONDS);

重组超时是分片方案的最大风险点:一旦超时把半成品丢掉,这条消息就永久丢失(Kafka 里没有「重发」的概念,除非生产者重试)。因此要么把超时设得足够长(覆盖最大消费延迟),要么把未收齐的消息送进死信队列人工处理。

6. 方案选型

6.1 三种方案对比

方案Kafka 内消息大小复杂度顺序性适用
调大上限原样(MB~几十 MB)低天然保证偶发大消息、内容可压缩
Claim Check指针(几百字节)中天然保证内容是文件、下游能访问对象存储
分片重组固定小片高需同 key 保证内容必须走 Kafka、无对象存储依赖

6.2 决策路径

按顺序问三个问题:

  1. 大消息是偶发还是常态? 偶发(比如一天几百条报表)→ 调大上限最省事,但要确认所有环节的参数都放大了。
  2. 内容能压缩吗? 文本类先试 zstd,若压缩后能回到默认上限内,问题直接消失。
  3. 下游能访问对象存储吗? 能 → 优先 Claim Check,它把 Kafka 完全摘出大消息问题。不能 → 只能分片。

实践中 Claim Check 覆盖了绝大多数场景,分片只在「必须走 Kafka」这一条硬约束下才用。因为分片把「消息完整性」的责任从 Kafka 转移到了应用代码,而这段代码要处理乱序、重复、超时、内存上限——每一处都可能出错。

7. 排错

7.1 症状与定位

症状一:生产者抛 RecordTooLargeException。 检查三处:客户端 max.request.size、broker message.max.bytes、topic 级 max.message.bytes(若设了)。三者取最小值生效。

症状二:消费者 poll 一直为空,但 lag 在涨。 几乎可以确定是 max.partition.fetch.bytes 小于最大消息。这是最隐蔽的一类,因为没有任何异常。

症状三:副本持续落后,ISR 反复收缩。 检查 replica.fetch.max.bytes。它与 message.max.bytes 不一致是大消息场景的经典配置错误。

症状四:broker 内存暴涨、频繁 Full GC。 大消息把堆打满。此时应减小并发请求数(num.network.threads 与 num.io.threads),或直接转向 Claim Check。

7.2 验证脚本

# 查 topic 的实际消息上限(可能被 topic 级配置覆盖)
kafka-configs.sh --bootstrap-server kafka:9092 \
  --entity-type topics --entity-name orders --describe

# 查最大消息尺寸的估算(用 kafka-run-class 的 dump-log-segment)
kafka-dump-log.sh --files /var/lib/kafka/data/orders-0/00000000000000000000.log \
  --print-data-log | awk '{print length($0)}' | sort -n | tail -1

8. 常见坑清单

  • 只调了 max.request.size,漏了 broker 与消费端,链路在中间断掉。
  • replica.fetch.max.bytes 小于 message.max.bytes,副本永远追不平。
  • max.partition.fetch.bytes 太小,消费者静默卡死,无任何异常。
  • max.poll.records 在大消息场景保持默认 500,单次 poll 内存爆掉。
  • 对已压缩的图片/PDF 开 zstd,compression-rate-avg ≥ 1,纯亏 CPU。
  • Claim Check 先发指针后传内容,消费者拿到指针时内容还不存在。
  • Claim Check 的对象 TTL 短于 topic 保留期,消费者撞 NoSuchKeyException。
  • 分片消息的 key 不用 msgId,分片散落多个分区,永远收不齐。
  • 分片重组缓冲无上限,慢消费时 OOM。
  • 重组超时后直接丢弃半成品,消息永久丢失且无人知晓。
  • 大消息场景仍设 linger.ms=100,多个大消息堆进一个批次。

9. 总结

Kafka 对大消息的「不友好」是设计取舍,不是缺陷。应对它有三条路径:偶发大消息靠成组放大上限参数(客户端、broker、副本、消费端四环一个都不能漏),文本类大消息优先靠压缩把体积压回安全区,常态大消息则应该换承载方式——Claim Check 把内容搬到对象存储,Kafka 只传指针;只有当内容必须走 Kafka 时才用分片,并严格处理乱序、重复、超时与缓冲上限。

无论选哪条路,都建议先把「最大消息尺寸」这个数字明确下来,再让所有参数围绕它对齐。参数不一致带来的故障往往没有任何报错,只是「不动了」,排查成本远高于一次性配置到位。若需要重新设计承载大消息的 topic 结构,可以参考 Kafka Topic 设计 中关于分区与保留策略的部分。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 集群升级与滚动重启实践
  2. Kafka 应用测试策略:Testcontainers 与集成测试
  3. 压缩算法选型:lz4、zstd、snappy 与 gzip