Kafka 的设计假设是「消息小、数量多」:默认单条上限 1 MB,整个存储与网络栈都为高吞吐的小记录优化。但现实里总有例外——一张图片、一段 JSON 快照、一份 PDF 报表,动辄几 MB 到几十 MB。把上限一路调大能让它「跑起来」,但代价会沿着内存、网络、副本同步、消费延迟一路传导。大消息的正确解法通常不是「把上限调大」,而是换一种承载方式。
1. 为什么 Kafka 不适合大消息
1.1 上限参数的层次
「Kafka 的消息上限」不是一个参数,而是一组参数,任何一环不够大,链路就断。而且不同环节的失败表现完全不同:
| 环节 | 参数 | 默认值 | 超限表现 |
|---|---|---|---|
| 生产端请求 | max.request.size | 1 MB | 客户端直接抛 RecordTooLargeException |
| broker 接收 | message.max.bytes | ~1 MB | broker 返回 RecordTooLargeException |
| 副本同步 | replica.fetch.max.bytes | 1 MB | follower 拉不动,副本永远落后 |
| 消费端拉取 | max.partition.fetch.bytes | 1 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、日志):压缩比 5
10 倍,效果显著,一条 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 决策路径
按顺序问三个问题:
- 大消息是偶发还是常态? 偶发(比如一天几百条报表)→ 调大上限最省事,但要确认所有环节的参数都放大了。
- 内容能压缩吗? 文本类先试
zstd,若压缩后能回到默认上限内,问题直接消失。 - 下游能访问对象存储吗? 能 → 优先 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 设计 中关于分区与保留策略的部分。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。