이 케이스에서 얻어 갈 것

상황

상품 카탈로그 동기화 파이프라인입니다. product-changed 토픽은 파티션 18개, 복제 계수 3, 브로커 3대에서 일 90만 건을 받습니다. 평균 레코드 크기는 40KB였습니다.

문제는 신규 기능에서 시작됐습니다. 상품 상세 페이지 리뉴얼로 이벤트에 상세 설명 HTML과 옵션 매트릭스가 인라인으로 포함되기 시작했습니다. 의류·가구 카테고리 일부 상품에서 레코드가 3~6MB가 되었습니다.

배포 30분 뒤 프로듀서 로그 — 예외가 원인을 정확히 알려 줍니다
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로 올려 배포했습니다. 여기서 사고가 커졌습니다.

1차 조치 — 이 한 줄만 바꿨습니다
# 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개 토픽의 발행도 함께 멈췄습니다.

case10 — 메시지 크기 설정 5종의 소속과 막히는 지점 왼쪽에서 오른쪽으로 프로듀서 앱, 프로듀서 버퍼, 브로커 수신, 토픽 override, 팔로워 복제, 컨슈머 fetch 순서로 레코드가 지나가는 경로를 그리고 각 지점에 관련 설정과 기본값, 소속을 표시했습니다. 프로듀서의 max.request.size 1048576, 브로커의 message.max.bytes 1048588, 토픽의 max.message.bytes 1048588 은 넘으면 거부하는 하드 게이트입니다. 프로듀서 쪽은 전송 전에 RecordTooLargeException 으로 실패하고 브로커 쪽은 MESSAGE_TOO_LARGE 로 거부합니다. 반면 브로커의 replica.fetch.max.bytes 1048576, 컨슈머의 max.partition.fetch.bytes 1048576, fetch.max.bytes 52428800 은 절대 상한이 아니며 첫 배치가 더 커도 그대로 반환해 복제와 소비가 멈추지 않도록 보장합니다. 토픽만 올리면 프로듀서에서 막히고 프로듀서만 올리면 브로커가 거부하므로 두 곳을 함께 올려야 합니다. 프로듀서 쪽 초과는 RecordTooLargeException 으로 즉시 실패합니다. 반면 브로커가 MESSAGE_TOO_LARGE 를 반환하면 Sender 가 배치를 쪼개 다시 보내며, 이때 재시도 횟수를 소모하지 않습니다 (Sender.java 소스 주석 명시). 그래서 요청 수가 폭증하고 버퍼가 고갈되어 무한 재시도처럼 보입니다. 레코드가 1건인 배치는 더 쪼갤 수 없어 최종 실패합니다. case10 — 1.5MB 레코드가 어디서 막히는가 (설정별 소속과 성격) 레코드가 지나가는 순서대로 상한 설정을 늘어놓았습니다. 위에서 아래로 읽습니다. 1 프로듀서 앱 send(1.5MB) 기본 — 애플리케이션 2 프로듀서 버퍼 max.request.size 기본 1048576 프로듀서 ▲ 하드 — 전송 전 거부 3 브로커 수신 message.max.bytes 기본 1048588 브로커 ■ 하드 — 브로커가 거부 4 토픽 override max.message.bytes 기본 1048588 토픽 ● 하드 — 브로커 값을 덮어씀 5 팔로워 복제 replica.fetch.max.bytes 기본 1048576 브로커 ■ 소프트 — 진행 보장 6 컨슈머 fetch max.partition.fetch.bytes 기본 1048576 컨슈머 ◆ 소프트 — 진행 보장 6 컨슈머 fetch fetch.max.bytes 기본 52428800 컨슈머 ◆ 소프트 — 요청 전체 예산 하드 게이트 — 넘으면 거부합니다 max.request.size (프로듀서) : 전송 전에 RecordTooLargeException 으로 즉시 실패 message.max.bytes / max.message.bytes (브로커·토픽) : 브로커가 MESSAGE_TOO_LARGE 로 거부 소프트 상한 — 진행을 막지 않습니다 replica.fetch.max.bytes · max.partition.fetch.bytes · fetch.max.bytes 는 절대 상한이 아닙니다. 첫 배치가 이 값보다 커도 그대로 반환해 복제와 소비가 멈추지 않게 보장합니다 (문서 명시). 어긋나는 지점과 결과 토픽 max.message.bytes 만 10MB 로 올리고 프로듀서를 그대로 두면 프로듀서에서 먼저 막힙니다. 반대로 프로듀서만 올리면 브로커가 MESSAGE_TOO_LARGE 로 거부합니다 — 두 곳을 함께 올려야 합니다. 프로듀서 쪽 초과는 RecordTooLargeException 으로 즉시 실패합니다 (ApiException — 재시도 불가). 브로커 쪽 MESSAGE_TOO_LARGE 는 Sender 가 배치를 쪼개 재전송하며 재시도 예산을 소모하지 않습니다. 처방 큰 레코드를 허용하려면 프로듀서 max.request.size 와 브로커·토픽 쪽 상한을 함께 올립니다. 복제·소비 쪽은 진행이 보장되므로 필수는 아니지만 처리량을 위해 함께 검토합니다. 확인: kafka-configs.sh --bootstrap-server :9092 --describe --entity-type topics --entity-name X 가능하면 큰 본문은 오브젝트 스토리지에 두고 Kafka 에는 참조 키만 보냅니다 (claim check 패턴). message.max.bytes 는 압축 후 배치 크기 기준입니다 — 압축을 켜면 실제 통과량이 달라집니다.
메시지 크기 설정 5개가 만드는 관문 — 프로듀서 max.request.size, 브로커 message.max.bytes, 복제 replica.fetch.max.bytes, 컨슈머 fetch 설정을 순서대로 통과해야 하고 하나만 올리면 다음 관문에서 막히는 구조

관측된 증상

메트릭이 어떻게 보였는가

프로듀서 로그 — 분할 재시도의 지문

SenderMESSAGE_TOO_LARGE를 받았을 때 배치에 레코드가 2건 이상이면 배치를 분할해 재전송합니다. 이때 남기는 로그가 "splitting and retrying"입니다. 이 문자열이 이 사고의 확정 증거입니다.

catalog-publisher 로그 — 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.

브로커 로그

브로커 쪽에서는 레코드 검증 단계에서 거부되며, 프로토콜 레벨 오류로 응답합니다. 브로커 로그에 남는 형태는 다음과 같습니다.

kafka-1 server.log
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단계 — 배치 분할 재시도가 폭증을 만든다

SenderMESSAGE_TOO_LARGE 처리 분기 조건은 세 가지가 모두 참일 때입니다.

이 조건이 참이면 배치를 분할해 다시 큐에 넣고, 재시도 횟수를 차감하지 않습니다. 조건이 거짓이면(예: 레코드 1건) 일반 오류 처리로 넘어가 실패합니다.

프로듀서 파이프라인·배치·재시도의 전체 동작은 4장 Producer 심화에서 다룹니다.

재현 방법

단일 노드 KRaft 클러스터로 재현됩니다. 아래 compose는 Apache Kafka의 공식 단일 노드 예제입니다 (3노드 구성은 예제 1).

docker-compose.yml — 단일 노드 KRaft (기본 크기 한계 유지)
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개 설정을 정합하게 올립니다 최후 수단 — 아래 참조
claim-check 패턴 — 이벤트는 작게, 본문은 스토리지에
// 변경 전 — 상세 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), 모든 관문을 여유를 두고 정합하게 올립니다. 컨슈머 쪽은 성능을 위해 함께 조정합니다.

변경 후 — 목표 최대 8MiB (8388608)
# ── 프로듀서 (비압축 기준 상한)
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           # 기본값. 파티션 수 × 위 값이 이 값을 크게 넘지 않게 균형

애플리케이션 — 재시도 불가 예외를 재시도하지 않는다

예외를 분류해 처리 — 이 사고의 2차 피해(OOM)를 막습니다
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);

예방 체크리스트

시험 포인트

공식 문서 출처