실수 케이스 · 10
큰 메시지가 무한 재시도로 쌓였다
RecordTooLargeException이 났고, 예외 메시지가
max.request.size를 정확히 지목했습니다.
그 값을 10MB로 올렸습니다. 그러자 이번엔 브로커가 거부했고,
프로듀서는 배치를 쪼개 재시도하기 시작했습니다.
요청 수가 40배로 늘고 버퍼가 고갈되면서 서비스 전체의 발행이 멈췄습니다.
메시지 크기에 관여하는 설정은 다섯 개이고, 하나만 올리면 다음 벽에 부딪힙니다.
이 케이스에서 얻어 갈 것
- 메시지 크기 관련 5개 설정의 정확한 이름과 소속(브로커/토픽/프로듀서/컨슈머)을 헷갈리지 않게 됩니다.
message.max.bytes가 압축 후 기준이고max.request.size가 비압축 기준이라는 차이를 압니다.- 브로커가
MESSAGE_TOO_LARGE를 반환했을 때 프로듀서가 배치를 분할해 재시도하는 동작을 압니다. - 4.x 컨슈머가 큰 레코드에서도 진행을 멈추지 않는다는 것과, 그럼에도 무엇을 맞춰야 하는지 압니다.
상황
상품 카탈로그 동기화 파이프라인입니다. product-changed 토픽은
파티션 18개, 복제 계수 3, 브로커 3대에서 일 90만 건을 받습니다.
평균 레코드 크기는 40KB였습니다.
문제는 신규 기능에서 시작됐습니다. 상품 상세 페이지 리뉴얼로 이벤트에 상세 설명 HTML과 옵션 매트릭스가 인라인으로 포함되기 시작했습니다. 의류·가구 카테고리 일부 상품에서 레코드가 3~6MB가 되었습니다.
ERROR Failed to publish product-changed event productId=P-8841902
org.apache.kafka.common.errors.RecordTooLargeException: The message is 4718392 bytes when serialized which is larger than 1048576, which is the value of the max.request.size configuration.
at org.apache.kafka.clients.producer.KafkaProducer.ensureValidRecordSize(KafkaProducer.java:1163)
at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:1024)
담당자는 메시지가 지목한 대로 max.request.size를 10MB로 올려 배포했습니다.
여기서 사고가 커졌습니다.
# producer.properties
max.request.size=10485760 # 1048576 → 10MB
# 브로커 message.max.bytes 는 기본값 1048588 그대로
# 토픽 max.message.bytes 도 미설정 → 1048588
# replica.fetch.max.bytes 도 기본값 1048576 그대로
# 컨슈머 max.partition.fetch.bytes 도 기본값 1048576 그대로
배포 직후 클라이언트 측 예외는 사라졌습니다.
하지만 15:41부터 프로듀서 로그에 WARN이 초당 수천 줄씩 쏟아지기 시작했고,
15:48에는 product-changed로의 발행이 전면 중단됐습니다.
같은 프로듀서 인스턴스를 공유하던 다른 3개 토픽의 발행도 함께 멈췄습니다.
max.request.size, 브로커 message.max.bytes,
복제 replica.fetch.max.bytes, 컨슈머 fetch 설정을 순서대로 통과해야 하고
하나만 올리면 다음 관문에서 막히는 구조
관측된 증상
메트릭이 어떻게 보였는가
- 프로듀서
record-error-rate: 0 → 급등. 다만 클라이언트 측 예외는 오히려 줄어 처음엔 개선처럼 보였습니다. - 프로듀서
request-rate: 초당 900 → 초당 36,000. 배치 분할이 만든 폭증입니다. - 프로듀서
buffer-available-bytes: 32MiB → 0에 붙음. 이 지표가 0이면send()가 블로킹됩니다. - 프로듀서
record-queue-time-avg: 3ms → 수십 초. - 브로커
RequestHandlerAvgIdlePercent: 0.85 → 0.19. - 애플리케이션 스레드:
send()에서 블로킹된 스레드가 쌓여 요청 처리 스레드 풀까지 고갈. - 컨슈머 lag: 이 사고에서는 정상. 프로듀서 쪽 사고이므로 컨슈머 지표는 조용합니다.
프로듀서 로그 — 분할 재시도의 지문
Sender는 MESSAGE_TOO_LARGE를 받았을 때
배치에 레코드가 2건 이상이면 배치를 분할해 재전송합니다.
이때 남기는 로그가 "splitting and retrying"입니다. 이 문자열이 이 사고의 확정 증거입니다.
attempts left가 줄지 않는 것을 보세요WARN [Producer clientId=catalog-publisher-2] Got error produce response in correlation id 1841992 on topic-partition product-changed-7, splitting and retrying (2147483647 attempts left). Error: MESSAGE_TOO_LARGE
WARN [Producer clientId=catalog-publisher-2] Got error produce response in correlation id 1841998 on topic-partition product-changed-7, splitting and retrying (2147483647 attempts left). Error: MESSAGE_TOO_LARGE
WARN [Producer clientId=catalog-publisher-2] Got error produce response in correlation id 1842004 on topic-partition product-changed-7, splitting and retrying (2147483647 attempts left). Error: MESSAGE_TOO_LARGE
delivery.timeout.ms 초과org.apache.kafka.common.errors.TimeoutException: Expiring 1 record(s) for product-changed-7:120004 ms has passed since batch creation. The request has not been sent, or no server response has been received yet.
분할이 불가능한 단일 레코드가 브로커 한계를 넘으면
재시도 없이 실패하고 아래 형태로 전달됩니다.
프로토콜 오류 MESSAGE_TOO_LARGE(코드 10)의 정의 문장이 그대로 실려 옵니다.
org.apache.kafka.common.errors.RecordTooLargeException: The request included a message larger than the max message size the server will accept.
브로커 로그
브로커 쪽에서는 레코드 검증 단계에서 거부되며, 프로토콜 레벨 오류로 응답합니다. 브로커 로그에 남는 형태는 다음과 같습니다.
ERROR [ReplicaManager broker=1] Error processing append operation on partition product-changed-7
org.apache.kafka.common.errors.RecordTooLargeException: The record batch size in the append to product-changed-7 is 4718392 bytes which exceeds the maximum configured value of 1048588.
원인 분석
1단계 — 다섯 개 설정의 정확한 이름과 소속
이 표가 이 케이스의 핵심입니다. 이름이 비슷하고 소속이 다릅니다. 기본값은 모두 Apache Kafka 4.3 문서에서 확인한 값입니다.
| 설정 | 소속 | 기본값 | 무엇을 제한하는가 | 초과 시 |
|---|---|---|---|---|
message.max.bytes |
broker | 1048588 | Kafka가 허용하는 최대 레코드 배치 크기 — 압축이 켜져 있으면 압축 후 기준 | MESSAGE_TOO_LARGE(10) |
max.message.bytes |
topic | 1048588 | 같은 의미의 토픽 레벨 오버라이드. 서버 기본값은 message.max.bytes |
MESSAGE_TOO_LARGE(10) |
max.request.size |
producer | 1048576 | 한 요청의 최대 크기. 사실상 최대 비압축 레코드 배치 크기의 상한 | 클라이언트에서 RecordTooLargeException |
max.partition.fetch.bytes |
consumer | 1048576 | 서버가 파티션당 반환할 최대 데이터량 | 초과해도 첫 배치는 반환됨 (아래 참조) |
fetch.max.bytes |
consumer | 52428800 | 한 fetch 요청에 서버가 반환할 전체 최대 데이터량 | 초과해도 첫 배치는 반환됨 |
replica.fetch.max.bytes |
broker | 1048576 | 팔로워가 파티션당 복제로 가져오려 시도할 바이트 수 | 초과해도 첫 배치는 반환됨 |
replica.fetch.response.max.bytes |
broker | 10485760 | 복제 fetch 응답 전체의 기대 최대 바이트 | 초과해도 첫 배치는 반환됨 |
socket.request.max.bytes |
broker | 104857600 | 소켓 요청의 절대 상한 (100MiB) | 이 값을 넘는 요청은 거부됩니다 |
2단계 — 관문을 순서대로 따라간다
하나의 큰 레코드가 프로듀서에서 컨슈머까지 도달하려면 네 개의 관문을 통과해야 합니다.
① 프로듀서 클라이언트 max.request.size (1048576, 비압축 기준)
│ 실패 → RecordTooLargeException (클라이언트에서 즉시)
↓
② 브로커 리더 append message.max.bytes (1048588) 또는
│ 토픽 max.message.bytes (오버라이드, 압축 후 기준)
│ 실패 → MESSAGE_TOO_LARGE → Sender 가 배치 분할 재시도
↓
③ 팔로워 복제 replica.fetch.max.bytes (1048576)
│ replica.fetch.response.max.bytes (10485760)
│ → 4.x 는 첫 배치를 반환해 진행을 보장
↓
④ 컨슈머 fetch max.partition.fetch.bytes (1048576)
fetch.max.bytes (52428800)
→ 4.x 는 첫 배치를 반환해 진행을 보장
3단계 — 배치 분할 재시도가 폭증을 만든다
Sender의 MESSAGE_TOO_LARGE 처리 분기 조건은 세 가지가 모두 참일 때입니다.
- 오류가
MESSAGE_TOO_LARGE다 - 배치의 레코드 수가 1보다 크다 (쪼갤 수 있어야 하므로)
- 배치가 아직 완료되지 않았고, 매직 값이 V2 이상이거나 압축되어 있다
이 조건이 참이면 배치를 분할해 다시 큐에 넣고, 재시도 횟수를 차감하지 않습니다. 조건이 거짓이면(예: 레코드 1건) 일반 오류 처리로 넘어가 실패합니다.
프로듀서 파이프라인·배치·재시도의 전체 동작은 4장 Producer 심화에서 다룹니다.
재현 방법
단일 노드 KRaft 클러스터로 재현됩니다. 아래 compose는 Apache Kafka의 공식 단일 노드 예제입니다 (3노드 구성은 예제 1).
services:
broker:
image: apache/kafka:4.3.1
hostname: broker
container_name: broker
ports:
- '9092:9092'
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT'
KAFKA_LISTENERS: 'CONTROLLER://:29093,PLAINTEXT://:19092,PLAINTEXT_HOST://:9092'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://broker:19092,PLAINTEXT_HOST://localhost:9092'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@broker:29093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
CLUSTER_ID: '4L6g3nShT-eMCtK--X86sw'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs'
# message.max.bytes 는 기본값 1048588 을 그대로 둡니다 (재현 포인트)
docker compose up -d
K=/opt/kafka/bin
docker exec -it broker $K/kafka-topics.sh --create --topic big-msg \
--partitions 1 --replication-factor 1 --bootstrap-server localhost:9092
# 3MiB 짜리 한 줄 파일을 만든다
docker exec broker sh -c 'head -c 3145728 /dev/zero | tr "\0" "x" > /tmp/big.txt'
# ── 관문 ① 프로듀서 max.request.size (기본 1048576)
docker exec -i broker sh -c "$K/kafka-console-producer.sh --topic big-msg \
--bootstrap-server localhost:9092 < /tmp/big.txt"
# → org.apache.kafka.common.errors.RecordTooLargeException:
# The message is 3145729 bytes when serialized which is larger than 1048576,
# which is the value of the max.request.size configuration.
# ── 관문 ① 를 열어 준다 → 관문 ② 브로커 message.max.bytes (기본 1048588) 에서 막힘
docker exec -i broker sh -c "$K/kafka-console-producer.sh --topic big-msg \
--bootstrap-server localhost:9092 \
--producer-property max.request.size=10485760 \
--producer-property delivery.timeout.ms=15000 < /tmp/big.txt"
# → Got error produce response ... Error: MESSAGE_TOO_LARGE
# 그리고 delivery.timeout.ms 초과 후:
# org.apache.kafka.common.errors.RecordTooLargeException: The request included a message
# larger than the max message size the server will accept.
# ── 관문 ② 를 토픽 레벨로 열어 준다
docker exec -it broker $K/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name big-msg \
--alter --add-config max.message.bytes=10485760
docker exec -i broker sh -c "$K/kafka-console-producer.sh --topic big-msg \
--bootstrap-server localhost:9092 \
--producer-property max.request.size=10485760 < /tmp/big.txt"
# → 성공합니다. 브로커 message.max.bytes 를 건드리지 않고
# 토픽 max.message.bytes 만 올려도 통과한다는 점을 확인하세요.
# ── 관문 ④ 컨슈머 — 기본값 그대로도 "읽힙니다" (4.x 동작 확인)
docker exec -it broker $K/kafka-console-consumer.sh --topic big-msg \
--from-beginning --timeout-ms 20000 --bootstrap-server localhost:9092 \
--property print.value=false --property print.offset=true
# → max.partition.fetch.bytes(1048576) 보다 큰 레코드지만 정상적으로 반환됩니다.
# "컨슈머가 멈춘다"는 오래된 설명이 현행 버전에서 틀렸음을 확인하는 단계입니다.
# 토픽을 다시 기본 한계로 되돌린다
docker exec -it broker $K/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name big-msg \
--alter --delete-config max.message.bytes
# 작은 레코드 사이에 큰 레코드를 섞어 한 배치에 들어가게 만듭니다.
# linger.ms 를 크게 주면 배치가 뭉쳐져 분할 재시도가 잘 보입니다.
docker exec broker sh -c '{ \
for i in $(seq 1 50); do echo "small-$i"; done; \
head -c 2097152 /dev/zero | tr "\0" "y"; echo; \
for i in $(seq 51 100); do echo "small-$i"; done; \
} > /tmp/mixed.txt'
docker exec -i broker sh -c "$K/kafka-console-producer.sh --topic big-msg \
--bootstrap-server localhost:9092 \
--producer-property max.request.size=10485760 \
--producer-property linger.ms=1000 \
--producer-property batch.size=5242880 \
--producer-property delivery.timeout.ms=20000 < /tmp/mixed.txt" 2>&1 \
| grep -c 'splitting and retrying'
# → "splitting and retrying" 가 여러 번 나타납니다.
# 이것이 요청 폭증의 원인입니다.
해결
즉시 조치 — 출혈을 멈춘다
가장 먼저 할 일은 크기 한계를 올리는 것이 아닙니다. 분할 재시도 폭증과 버퍼 고갈을 멈춰야 합니다.
# 1. 프로듀서를 롤백한다. max.request.size 를 원래대로 돌려 클라이언트에서
# 즉시 실패하게 만듭니다. 브로커까지 가서 폭증하는 것보다 낫습니다.
kubectl rollout undo deploy/catalog-publisher
# 2. 애플리케이션의 자체 재시도 큐를 비운다.
# RecordTooLargeException 은 재시도 불가 예외입니다. 재시도하면 안 됩니다.
# 3. 큰 레코드를 만드는 기능을 끈다 (피처 플래그).
# 이것이 진짜 즉시 조치입니다.
# 4. 영향 범위를 좁힌다 — 같은 프로듀서 인스턴스를 쓰는 다른 토픽을 분리한다.
# 최소한 대형 페이로드 토픽 전용 프로듀서를 별도로 만듭니다.
근본 해결 1 — 큰 메시지를 아예 보내지 않는다 (권장)
크기 한계를 올리는 것은 마지막 선택지입니다. 먼저 페이로드를 줄일 방법을 검토하세요.
| 방법 | 내용 | 비용 |
|---|---|---|
| claim-check 패턴 | 본문을 오브젝트 스토리지에 올리고 이벤트에는 참조(URL·키)와 메타데이터만 담습니다 | 권장 — 저장소 의존이 생기지만 크기 문제가 근본적으로 사라집니다 |
| 압축 활성화 | compression.type=lz4 또는 zstd. 브로커 한계는 압축 후 기준이므로 실질 여유가 생깁니다 |
권장 — CPU 비용. HTML·JSON은 압축률이 매우 높습니다 |
| 이벤트 분해 | "상품 변경" 하나에 모든 필드를 담지 말고 변경된 부분만 담거나 여러 이벤트로 나눕니다 | 도메인 모델링 작업이 필요합니다 |
| 직렬화 형식 변경 | JSON → Avro/Protobuf. 스키마가 필드 이름을 반복하지 않아 크기가 줄어듭니다 | 스키마 관리 도입 (8장, 케이스 9) |
| 크기 한계 상향 | 5개 설정을 정합하게 올립니다 | 최후 수단 — 아래 참조 |
// 변경 전 — 상세 HTML 을 인라인으로 담았습니다 (4.7MB)
ProductChanged before = ProductChanged.newBuilder()
.setProductId(id)
.setDescriptionHtml(hugeHtml) // ← 이것이 문제
.setOptionMatrix(hugeMatrix)
.build();
// 변경 후 — 본문은 스토리지에 올리고 참조만 담습니다 (2KB)
String bodyKey = objectStore.put("product-body/" + id + "/" + revision, hugeHtml);
ProductChanged after = ProductChanged.newBuilder()
.setProductId(id)
.setRevision(revision)
.setBodyRef(bodyKey) // ← 참조
.setBodySha256(sha256(hugeHtml)) // ← 무결성 확인용
.setBodySize(hugeHtml.length())
.build();
근본 해결 2 — 올려야 한다면 다섯 개를 함께 올린다
예외 메시지가 지목한 한 개만 올렸습니다. 다음 관문에서 막히고, 그 실패는 클라이언트 예외가 아니라 브로커 왕복 후의 분할 재시도로 나타납니다.
# producer
max.request.size=10485760
# 나머지는 전부 기본값
# 브로커 message.max.bytes = 1048588
# 토픽 max.message.bytes = 1048588
# 브로커 replica.fetch.max.bytes = 1048576
# 컨슈머 max.partition.fetch.bytes = 1048576
목표 최대 크기를 정하고(여기서는 8MiB), 모든 관문을 여유를 두고 정합하게 올립니다. 컨슈머 쪽은 성능을 위해 함께 조정합니다.
# ── 프로듀서 (비압축 기준 상한)
max.request.size=10485760 # 목표보다 여유 있게
compression.type=lz4 # 브로커 한계는 압축 후 기준이므로 효과가 큽니다
buffer.memory=268435456 # 큰 레코드를 다루려면 버퍼도 키워야 합니다
max.block.ms=10000 # 무한정 블로킹되지 않게 짧게
linger.ms=5 # 4.x 기본값
# ── 토픽 (권장: 브로커 전역이 아니라 이 토픽만)
# kafka-configs --entity-type topics --entity-name product-changed \
# --alter --add-config max.message.bytes=10485760
max.message.bytes=10485760 # 압축 후 기준
# ── 브로커 (복제 경로)
replica.fetch.max.bytes=10485760
replica.fetch.response.max.bytes=20971520
# socket.request.max.bytes(104857600) 를 넘지 않아야 합니다
# ── 컨슈머 (성능 목적. 4.x 는 이 값 미달이어도 첫 배치를 받습니다)
max.partition.fetch.bytes=10485760
fetch.max.bytes=52428800 # 기본값. 파티션 수 × 위 값이 이 값을 크게 넘지 않게 균형
애플리케이션 — 재시도 불가 예외를 재시도하지 않는다
producer.send(record, (metadata, e) -> {
if (e == null) return;
if (e instanceof RecordTooLargeException) {
// 재시도 불가. 재시도 큐에 넣으면 무한히 쌓입니다.
metrics.counter("publish.too_large").increment();
dlq.send(record.key(), summarize(record)); // 본문 없이 요약만
alerts.page("record too large: " + record.topic() + " key=" + record.key());
return;
}
if (e instanceof RetriableException) {
// delivery.timeout.ms 안에서 클라이언트가 이미 재시도했습니다.
// 여기까지 왔다면 그 예산이 끝난 것입니다.
outbox.save(record);
return;
}
log.error("publish failed topic={} key={}", record.topic(), record.key(), e);
});
발행 전 크기 검사
private static final int SOFT_LIMIT = 900 * 1024; // 한계의 90% 를 경보선으로
byte[] payload = serializer.serialize(topic, event);
if (payload.length > SOFT_LIMIT) {
metrics.gauge("publish.payload_bytes", payload.length);
// 한계에 도달하기 전에 claim-check 경로로 전환합니다
return publishWithClaimCheck(event);
}
producer.send(new ProducerRecord<>(topic, event.getKey(), payload), callback);
예방 체크리스트
시험 포인트
이어서 볼 곳
공식 문서 출처
- Broker Configs —
message.max.bytes— 기본값 1048588, 압축 후 레코드 배치 기준, 토픽 레벨max.message.bytes로 오버라이드, cluster-wide 동적 설정 - Topic Configs —
max.message.bytes— 기본값 1048588, 서버 기본값은message.max.bytes - Producer Configs —
max.request.size— 기본값 1048576, 사실상 최대 비압축 배치 상한, 서버 한계와 다를 수 있음 - Consumer Configs —
max.partition.fetch.bytes·fetch.max.bytes— 기본값 1048576 / 52428800, 첫 배치가 한계를 넘어도 반환되어 진행이 보장됨 - Broker Configs —
replica.fetch.max.bytes·replica.fetch.response.max.bytes— 기본값 1048576 / 10485760, 절대 최대치가 아님 - Broker Configs —
socket.request.max.bytes— 기본값 104857600 - Producer Configs —
buffer.memory·max.block.ms·compression.type— 기본값 33554432 / 60000 /none, 버퍼 고갈 시BufferExhaustedException - Producer Configs —
delivery.timeout.ms·retries— 기본값 120000 / 2147483647 - Apache Kafka 소스 (4.3) —
KafkaProducer.ensureValidRecordSize(),Sender의 배치 분할 재시도 로직과 "재시도 횟수를 차감하지 않는다" 주석,Sender.failExpiredBatches(),Errors.MESSAGE_TOO_LARGE,RecordTooLargeException의 상속 관계