구현: 메시지 형식

구현: 메시지 형식 (Message Format)

이 페이지는 Kafka에서 메시지(레코드)가 디스크와 네트워크 위에서 어떤 이진 형식으로 저장·전송되는지를 다뤄요. 프로토콜 해석이나 바이너리 포맷을 직접 다룰 일이 있을 때 유용하지만, 평소 Kafka API만 쓰는 경우엔 그냥 참고용으로 가볍게 읽으면 돼요.

출처: 문서

본문

메시지(일명 레코드, Record)는 항상 배치(batch) 단위로 기록됩니다. 메시지 배치의 기술 용어는 레코드 배치(record batch)이며, 레코드 배치는 하나 이상의 레코드를 포함합니다. 극단적인 경우 단일 레코드만 담긴 레코드 배치도 가능합니다. 레코드 배치와 레코드는 각각 고유한 헤더를 갖습니다. 각각의 형식은 아래에 설명되어 있습니다.

레코드 배치 (Record Batch)

다음은 RecordBatch의 디스크 저장 형식입니다.

baseOffset: int64
batchLength: int32
partitionLeaderEpoch: int32
magic: int8 (current magic value is 2)
crc: uint32
attributes: int16
    bit 0~2:
        0: no compression
        1: gzip
        2: snappy
        3: lz4
        4: zstd
    bit 3: timestampType
    bit 4: isTransactional (0 means not transactional)
    bit 5: isControlBatch (0 means not a control batch)
    bit 6: hasDeleteHorizonMs (0 means baseTimestamp is not set as the delete horizon for compaction)
    bit 7~15: unused
lastOffsetDelta: int32
baseTimestamp: int64
maxTimestamp: int64
producerId: int64
producerEpoch: int16
baseSequence: int32
recordsCount: int32
records: [Record]

압축이 활성화된 경우에는 압축된 레코드 데이터가 레코드 수(count) 바로 뒤에 직렬화되어 이어진다는 점에 유의하세요.

batchLength는 현재 위치(batchLength 필드 바로 다음)부터 배치 끝까지의 바이트 수를 나타냅니다. 즉, 디스크에서 레코드 배치의 총 크기는 batchLength + 12 바이트이며, 여기에는 8바이트의 baseOffset과 4바이트의 batchLength 필드 자체가 포함됩니다.

CRC는 attributes부터 배치 끝까지의 데이터를 포함합니다(즉 CRC 뒤에 오는 모든 바이트). CRC는 magic 바이트 뒤에 위치하므로, 클라이언트는 배치 길이와 magic 바이트 사이의 바이트를 어떻게 해석할지 결정하기 전에 magic 바이트를 먼저 파싱해야 합니다. 파티션 리더 에포크(partition leader epoch) 필드는 CRC 계산에 포함되지 않는데, 이 필드는 브로커가 받는 모든 배치에 대해 배정되므로 CRC를 다시 계산할 필요가 없도록 하기 위함입니다. 계산에는 CRC-32C(Castagnoli) 폴리노미얼이 사용됩니다.

컴팩션 시에는 로그가 정리될 때 원래 배치에서 첫 번째와 마지막 오프셋/시퀀스 번호를 보존합니다. 이는 로그가 다시 로드될 때 프로듀서 상태를 복원할 수 있도록 하기 위해 필요합니다. 예를 들어 마지막 시퀀스 번호를 유지하지 않으면 파티션 리더 장애 후 프로듀서가 OutOfSequence 오류를 볼 수 있습니다. base 시퀀스 번호는 중복 검사를 위해 보존되어야 합니다(브로커는 들어오는 Produce 요청의 첫·마지막 시퀀스 번호가 해당 프로듀서의 마지막 시퀀스 번호와 일치하는지 확인해 중복을 검사합니다). 그 결과, 배치의 모든 레코드가 정리되었지만 프로듀서의 마지막 시퀀스 번호를 보존하기 위해 배치 자체는 유지되는 경우 로그에 빈 배치가 있을 수 있습니다. 여기서 한 가지 특이한 점은 baseTimestamp 필드는 컴팩션 중에 보존되지 않으므로, 배치의 첫 레코드가 컴팩션되어 사라지면 이 값이 변경될 수 있다는 것입니다.

컴팩션은 레코드 배치에 null 페이로드 또는 중단된 트랜잭션 마커가 포함된 경우 baseTimestamp를 수정할 수도 있습니다. 이 경우 baseTimestamp는 해당 레코드들이 삭제되어야 하는 시점의 타임스탬프로 설정되며, delete horizon 속성 비트도 함께 설정됩니다.

컨트롤 배치 (Control Batches)

컨트롤 배치는 컨트롤 레코드(control record)라는 단일 레코드를 담고 있습니다. 컨트롤 레코드는 애플리케이션에 전달되어서는 안 됩니다. 대신 컨슈머가 중단된 트랜잭션 메시지를 걸러내는 데 사용합니다.

컨트롤 레코드의 키는 다음 스키마를 따릅니다.

version: int16 (current version is 0)
type: int16 (0 indicates an abort marker, 1 indicates a commit)

컨트롤 레코드 값의 스키마는 type에 따라 달라집니다. 값은 클라이언트에게 불투명(opaque)합니다.

레코드 (Record)

각 레코드의 디스크 저장 형식은 아래에 정의되어 있습니다.

length: varint
attributes: int8
    bit 0~7: unused
timestampDelta: varlong
offsetDelta: varint
keyLength: varint
key: byte[]
valueLength: varint
value: byte[]
headersCount: varint
Headers => [Header]

레코드 헤더 (Record Header)

headerKeyLength: varint
headerKey: String
headerValueLength: varint
Value: byte[]

레코드 헤더의 키는 null이 아님이 보장되지만, 레코드 헤더의 값은 null일 수 있습니다. 레코드의 헤더 순서는 프로듀싱과 컨슈밍 시에 보존됩니다.

우리는 Protobuf와 동일한 varint 인코딩을 사용합니다. 후자에 대한 자세한 정보는 여기에서 찾을 수 있습니다. 레코드의 헤더 개수도 varint로 인코딩됩니다.

이전 메시지 형식 (Old Message Format)

Kafka 0.11 이전에는 메시지가 메시지 셋(message set)으로 전송되고 저장되었습니다. 자세한 내용은 이전 메시지 형식(Old Message Format)을 참고하세요.

더 알아보기 (Learn more)

  • 레코드 배치의 magic 버전과 CRC 검증은 클라이언트·브로커 간 호환성의 기초예요.
  • 컨트롤 배치와 트랜잭션 마커는 트랜잭션 프로토콜에서 중단된 메시지를 걸러내는 핵심 역할을 해요.