프로토콜
프로토콜 (Protocol)
Kafka가 클라이언트와 브로커 사이에 실제로 쓰는 유선 프로토콜(wire protocol)을 설명하는 문서예요. 사용 가능한 요청들, 그 바이너리 형식, 그리고 클라이언트를 구현하려면 이 프로토콜을 올바르게 쓰는 방법을 다뤄요. Kafka의 기본 설계와 용어를 이해하고 있다고 가정해요.
출처: 문서
본문
서론 (Preliminaries)
네트워크 (Network)
Kafka는 TCP 위의 바이너리 프로토콜을 사용해요. 프로토콜은 모든 API를 요청-응답 메시지 쌍으로 정의해요. 모든 메시지는 크기로 구분되며 기본 타입들로 구성돼요.
클라이언트는 소켓 연결을 시작한 뒤 일련의 요청 메시지를 쓰고, 그에 대응하는 응답 메시지를 읽어요. 연결·해제 시 핸드셰이크는 필요 없어요. TCP 핸드셰이크 비용을 분산하려고 여러 요청에 재사용되는 지속 연결을 유지하면 TCP가 더 좋아하지만, 그 페널티 외에 연결은 꽤 저렴해요.
데이터가 파티셔닝되어 클라이언트가 자신의 데이터가 있는 서버와 통신해야 하므로, 클라이언트는 여러 브로커에 연결을 유지해야 할 수 있어요. 다만 단일 클라이언트 인스턴스에서 단일 브로커에 여러 연결을 유지하는 건 보통 필요하지 않아요(연결 풀링).
서버는 단일 TCP 연결에서 요청이 보내진 순서대로 처리되고 응답도 그 순서대로 돌아온다고 보장해요. 이 순서를 보장하려고 브로커의 요청 처리는 연결당 단일 in-flight 요청만 허용해요. 클라이언트는 논블로킹 I/O로 요청 파이프라이닝을 구현해 더 높은 처리량을 달성할 수 있어요(미결 요청은 기본 OS 소켓 버퍼에 버퍼링되므로). 모든 요청은 클라이언트가 시작하고, 별도로 언급된 경우를 제외하면 서버로부터 대응 응답 메시지를 결과로 받아요.
서버에는 요청 크기에 대한 설정 가능한 최대 한도가 있고, 이 한도를 초과하는 요청은 소켓이 끊겨요.
파티셔닝과 부트스트래핑 (Partitioning and bootstrapping)
Kafka는 파티셔닝된 시스템이라 모든 서버가 완전한 데이터 집합을 가지지 않아요. 토픽은 미리 정의된 수의 파티션 P로 나뉘고, 각 파티션은 복제 팩터 N으로 복제돼요. 토픽 파티션은 그 자체로 0, 1, ..., P-1로 번호 매겨진 정렬된 "커밋 로그"예요.
이런 시스템은 모두 특정 데이터 조각을 특정 파티션에 어떻게 배정하는지의 문제가 있어요. Kafka 클라이언트가 이 배정을 직접 제어하고, 브로커는 어떤 메시지를 어떤 파티션에 게시해야 하는지에 대한 특정 의미를 강제하지 않아요. 대신 메시지를 게시하려면 클라이언트가 메시지를 특정 파티션에 직접 주소 지정하고, 가져올 때도 특정 파티션에서 가져와요. 두 클라이언트가 같은 파티셔닝 방식을 쓰려면 같은 키→파티션 매핑 계산 방법을 사용해야 해요.
이 게시·fetch 요청은 현재 주어진 파티션의 리더 역할을 하는 브로커로 보내야 해요. 이 조건은 브로커가 강제하므로, 특정 파티션에 대한 요청을 잘못된 브로커로 보내면 NotLeaderForPartition 오류 코드가 돼요.
클라이언트는 어떤 토픽이 있고 어떤 파티션을 가지며, 어떤 브로커가 그 파티션을 호스팅하는지 어떻게 알 수 있을까요? 이 정보는 동적이므로 정적 매핑 파일로 각 클라이언트를 구성할 수는 없어요. 대신 모든 Kafka 브로커가 클러스터의 현재 상태(어떤 토픽, 그 파티션, 파티션의 리더 브로커, 브로커의 호스트·포트)를 설명하는 메타데이터 요청에 응답할 수 있어요.
즉 클라이언트는 브로커 하나를 찾아야 하고, 그 브로커가 존재하는 다른 모든 브로커와 그들이 호스팅하는 파티션을 알려줘요. 이 첫 브로커도 죽을 수 있으므로, 클라이언트 구현의 모범 사례는 부트스트랩할 23개의 URL 목록을 받는 것이에요. 사용자는 로드 밸런서를 쓰거나 클라이언트에 Kafka 호스트 23개를 정적으로 구성할 수 있어요.
클라이언트는 클러스터가 바뀌었는지 폴링할 필요 없어요. 인스턴스화할 때 메타데이터를 한 번 가져와, 메타데이터가 오래됐다는 오류를 받을 때까지 그 메타데이터를 캐시할 수 있어요. 이 오류는 두 형태로 올 수 있어요: (1) 특정 브로커와 통신할 수 없음을 나타내는 소켓 오류, (2) 그 브로커가 더 이상 데이터를 요청한 파티션을 호스팅하지 않는다는 요청 응답의 오류 코드.
- 연결할 수 있는 하나를 찾을 때까지 "bootstrap" Kafka URL 목록을 순회해요. 클러스터 메타데이터를 가져와요.
- 보내거나 가져오는 토픽/파티션에 따라 적절한 브로커로 fetch·produce 요청을 보내요.
- 적절한 오류를 받으면 메타데이터를 갱신하고 다시 시도해요.
파티셔닝 전략 (Partitioning Strategies)
앞서 언급했듯 메시지의 파티션 배정은 생산하는 클라이언트가 제어해요. 그렇다면 이 기능을 최종 사용자에게 어떻게 노출해야 할까요?
파티셔닝은 Kafka에서 두 가지 목적을 서빙해요:
- 브로커에 걸쳐 데이터와 요청 부하를 균형 맞춤.
- 컨슈머 프로세스 간 처리를 나누면서 로컬 상태를 허용하고 파티션 내 순서를 보존하는 수단. 이를 의미적 파티셔닝(semantic partitioning)이라 불러요.
단순한 부하 분산을 위해 클라이언트가 모든 브로커에 라운드 로빈으로 요청할 수도 있어요. 브로커보다 프로듀서가 훨씬 많은 환경에서는 각 클라이언트가 파티션 하나를 무작위로 골라 게시하는 방법도 있어요. 후자는 TCP 연결이 훨씬 적어요.
의미적 파티셔닝은 메시지의 어떤 키를 사용해 메시지를 파티션에 배정하는 거예요. 예를 들어 클릭 메시지 스트림을 처리한다면 user id로 스트림을 파티셔닝해 특정 사용자의 모든 데이터가 단일 컨슈머로 가게 할 수 있어요. 이렇게 하려면 클라이언트가 메시지와 연결된 키를 가져와 이 키의 어떤 해시로 메시지를 전달할 파티션을 골라요.
배칭 (Batching)
우리 API는 효율성을 위해 작은 것들을 함께 배치하도록 권장해요. 이는 매우 중요한 성능 이득이에요. 메시지를 보내는 API와 가져오는 API 모두 항상 단일 메시지가 아니라 메시지 시퀀스로 동작해 이를 장려해요. 영리한 클라이언트는 이를 활용해 개별로 보내진 메시지를 배치로 묶어 더 큰 덩어리로 보내는 "비동기" 모드를 지원할 수 있어요. 더 나아가 여러 토픽·파티션에 걸친 배칭도 허용해요 — produce 요청은 많은 파티션에 추가할 데이터를 담을 수 있고, fetch 요청은 많은 파티션에서 한 번에 데이터를 당길 수 있어요.
클라이언트 구현자는 원하면 이를 무시하고 모든 것을 한 번에 하나씩 보낼 수도 있어요.
호환성 (Compatibility)
Kafka는 "양방향" 클라이언트 호환성 정책이 있어요. 즉 새 클라이언트가 옛 서버와, 옛 클라이언트가 새 서버와 통신할 수 있어요. 이로써 사용자가 다운타임 없이 클라이언트나 서버를 업그레이드할 수 있어요.
Kafka 프로토콜은 시간이 지나며 바뀌었으므로 클라이언트와 서버는 유선으로 보내는 메시지의 스키마에 합의해야 해요. 이는 API 버전 관리를 통해 이뤄줘요. 각 요청을 보내기 전에 클라이언트는 API 키와 API 버전을 보내요. 이 두 16비트 숫자를 함께 보면 뒤따를 메시지의 스키마를 고유하게 식별해요.
클라이언트는 API 버전 범위를 지원하려는 의도예요. 특정 브로커와 통신할 때 주어진 클라이언트는 둘 다 지원하는 가장 높은 API 버전을 사용하고 그 버전을 요청에 표시해요. 서버는 지원하지 않는 버전의 요청을 거부하고, 항상 클라이언트가 요청에 포함한 버전에 기반해 기대하는 정확한 프로토콜 형식으로 응답해요. 의도된 업그레이드 경로는 새 기능을 먼저 서버에 배포하고(옛 클라이언트는 그 기능을 쓰지 않음), 새 클라이언트가 배포되면서 점진적으로 활용하는 것이에요. 지원 API 버전을 검색할 때는 서버가 다른 버전으로 응답할 수 있는 예외적인 경우가 있다는 점을 참고해요.
참고로 KIP-482 tagged fields는 버전 번호를 올리지 않고 요청에 추가될 수 있어요. 이는 호환성을 깨지 않고 메시지 스키마를 발전시키는 추가 방법을 제공해요. 태그된 필드는 설정되지 않으면 공간을 차지하지 않아요. 그래서 드물게 쓰이는 필드는 필수 스키마에 넣는 것보다 태그된 필드로 만드는 게 더 효율적이에요. 다만 태그된 필드는 모르는 수신자에게 무시되는데, 이것이 발신자가 원하는 동작이 아니면 문제가 될 수 있어요. 그런 경우 버전 올림이 더 적절할 수 있어요.
지원 API 버전 검색 (Retrieving Supported API versions)
여러 브로커 버전에 대해 동작하려면 클라이언트가 다양한 API의 어떤 버전을 브로커가 지원하는지 알아야 해요. 브로커는 KIP-35에 설명된 대로 0.10.0.0부터 이 정보를 노출해요. 클라이언트는 지원 API 버전 정보를 사용해 클라이언트와 브로커가 모두 지원하는 가장 높은 API 버전을 골라야 해요. 그런 버전이 없으면 사용자에게 오류를 보고해야 해요.
클라이언트가 브로커에서 지원 API 버전을 얻는 순서:
- 클라이언트는 브로커와 연결이 설정된 후
ApiVersionsRequest를 브로커에 보내요. SSL이 활성화되면 SSL 연결이 설정된 후에 이뤄져요. ApiVersionsRequest를 받으면 브로커는 현재 인증 상태와 무관하게 지원하는 ApiKeys와 버전의 전체 목록을 반환해요. 브로커 버전 정보 누출로 간주되면 SSL 클라이언트 인증을 쓰는 우회책이 있어요. 0.10.0.0보다 오래된 브로커는 이 API를 지원하지 않아 요청을 무시하거나 연결을 닫을 수 있어요. 클라이언트ApiVersionsRequest버전을 브로커가 지원하지 않고(클라이언트가 앞섬) 브로커가 2.4.0 이상이면, 브로커는 오류 코드UNSUPPORTED_VERSION,api_versions필드에 지원 버전을 담은 버전 0 ApiVersionsResponse로 응답해요. 이후 클라이언트가 클라이언트·브로커가 지원하는 가장 높은 버전으로 재시도해요 (KIP-511).- 브로커와 클라이언트가 API의 여러 버전을 지원하면 클라이언트는 브로커와 자신이 지원하는 최신 버전을 사용하는 게 권장돼요.
- 프로토콜 버전의 비권장은 프로토콜 문서에서 API 버전을 비권장으로 표시해 이뤄줘요.
- 브로커에서 얻은 지원 API 버전은 그 정보를 얻은 연결에만 유효해요. 연결이 끊기면 그 사이 브로커가 업그레이드/다운그레이드됐을 수 있으므로 클라이언트는 다시 얻어야 해요.
SASL 인증 순서 (SASL Authentication Sequence)
SASL 인증에 사용되는 순서:
- 클라이언트가
ApiVersionsRequest를 보내 브로커가 지원하는 요청 버전 범위를 얻을 수 있어요(선택). - 클라이언트가 인증용 SASL 메커니즘을 담은
SaslHandshakeRequest를 보내요. 요청된 메커니즘이 서버에서 활성화되지 않으면 서버는 지원 메커니즘 목록으로 응답하고 클라이언트 연결을 닫아요. 활성화되어 있으면 성공 응답을 보내고 SASL 인증을 계속해요. - 실제 SASL 인증이 이제 수행돼요.
SaslHandshakeRequest버전이 v0이면 메커니즘에 해당하는 SASL 클라이언트·서버 토큰 시리즈가 Kafka 프로토콜 헤더로 감싸지 않고 불투명 패킷으로 보내져요. v1이면SaslAuthenticate요청/응답이 사용되며 실제 SASL 토큰이 Kafka 프로토콜로 감싸져요. 브로커의 마지막 메시지의 오류 코드가 인증 성공 여부를 나타내요. - 인증이 성공하면 이후 패킷은 Kafka API 요청으로 처리돼요. 그렇지 않으면 클라이언트 연결이 닫혀요.
0.9.0.x 클라이언트와의 상호운용을 위해, 서버가 받은 첫 패킷이 유효한 Kafka 요청이 아니면 SASL/GSSAPI 클라이언트 토큰으로 처리해요. 이 패킷부터 시작해 위의 처음 두 단계를 건너뛰고 SASL/GSSAPI 인증을 수행해요.
프로토콜 (The Protocol)
프로토콜 기본 타입 (Protocol Primitive Types)
프로토콜은 다음 기본 타입들로 구성돼요.
| 타입 | 설명 |
|---|---|
| BOOLEAN | 바이트의 불리언 값. 0은 false, 1은 true. 0이 아닌 값은 true로 간주. |
| INT8 | -2⁷부터 2⁷-1까지의 정수. |
| INT16 | -2¹⁵부터 2¹⁵-1까지의 정수. 네트워크 바이트 순서(big-endian)로 2바이트 인코딩. |
| INT32 | -2³¹부터 2³¹-1까지의 정수. 네트워크 바이트 순서로 4바이트 인코딩. |
| INT64 | -2⁶³부터 2⁶³-1까지의 정수. 네트워크 바이트 순서로 8바이트 인코딩. |
| UINT16 | 0부터 65535까지의 정수. 네트워크 바이트 순서로 2바이트 인코딩. |
| UINT32 | 0부터 2³²-1까지의 정수. 네트워크 바이트 순서로 4바이트 인코딩. |
| VARINT | -2³¹부터 2³¹-1까지의 정수. Google Protocol Buffers의 가변 길이 zig-zag 인코딩을 따름. |
| VARLONG | -2⁶³부터 2⁶³-1까지의 정수. Google Protocol Buffers의 가변 길이 zig-zag 인코딩. |
| UUID | type 4 불변 범용 고유 식별자(Uuid). 네트워크 바이트 순서로 16바이트 인코딩. |
| FLOAT64 | IEEE 754 배정밀도 64비트 값. 네트워크 바이트 순서로 8바이트 인코딩. |
| STRING | 문자 시퀀스. 먼저 길이 N을 INT16으로, 그다음 N바이트의 UTF-8 인코딩 문자 시퀀스. 길이는 음수가 아니어야 함. |
| COMPACT_STRING | 문자 시퀀스. 먼저 길이 N+1을 UNSIGNED_VARINT로, 그다음 N바이트의 UTF-8 문자 시퀀스. |
| NULLABLE_STRING | 문자 시퀀스 또는 null. null이 아닌 문자열은 길이 N을 INT16으로, 그다음 N바이트의 UTF-8. null은 길이 -1로 인코딩되고 뒤따르는 바이트 없음. |
| COMPACT_NULLABLE_STRING | 문자 시퀀스. 길이 N+1을 UNSIGNED_VARINT로, 그다음 N바이트 UTF-8. null 문자열은 길이 0으로 표현. |
| BYTES | 원시 바이트 시퀀스. 길이 N을 INT32로, 그다음 N바이트. |
| COMPACT_BYTES | 원시 바이트 시퀀스. 길이 N+1을 UNSIGNED_VARINT로, 그다음 N바이트. |
| NULLABLE_BYTES | 원시 바이트 시퀀스 또는 null. null이 아닌 값은 길이 N을 INT32로. null은 길이 -1로 인코딩. |
| COMPACT_NULLABLE_BYTES | 원시 바이트 시퀀스. null 객체는 길이 0으로 표현. |
| RECORDS | BYTES로 표현된 Kafka 레코드 시퀀스. 자세한 설명은 Message Sets 참고. |
| COMPACT_RECORDS | COMPACT_BYTES로 표현된 Kafka 레코드 시퀀스. |
| NULLABLE_RECORDS | NULLABLE_BYTES로 표현된 Kafka 레코드 시퀀스. |
| COMPACT_NULLABLE_RECORDS | COMPACT_NULLABLE_BYTES로 표현된 Kafka 레코드 시퀀스. |
| ARRAY | 주어진 타입 T의 객체 시퀀스. T는 기본 타입(예: STRING)이거나 구조. 길이 N을 INT32로, 그다음 T의 인스턴스 N개. 문서에서 [T]로 표기. |
| COMPACT_ARRAY | 타입 T의 객체 시퀀스. 길이 N+1을 UNSIGNED_VARINT로. 문서에서 (T)로 표기. |
| NULLABLE_ARRAY | 타입 T의 객체 시퀀스. null 배열은 길이 -1로 표현. 문서에서 ?[T]로 표기. |
| COMPACT_NULLABLE_ARRAY | 타입 T의 객체 시퀀스. null 배열은 길이 0으로 표현. 문서에서 ?(T)로 표기. |
| STRUCT | 대문자 첫 글자 문자열로 이름 붙은, 하나 이상의 필드로 구성된 구조. 정의된 순서로 각 필드의 직렬화로 인코딩된 복합 객체. 문서에서 { } 로 둘러쌈. |
| NULLABLE_STRUCT | null 허용 구조. null이 아닌 값은 첫 바이트가 1이고 그다음 각 필드 직렬화. null은 값 -1의 바이트로 인코딩. 문서에서 ?{ }로 둘러쌈. |
요청 형식 문법을 읽는 방법
아래 BNF는 요청·응답 바이너리 형식에 대한 정확한 문맥 자유 문법을 줘요. BNF는 사람이 읽을 수 있는 이름을 주려고 의도적으로 컴팩트하지 않아요. BNF에서 생산 시퀀스는 연결을 나타내고, 여러 가능한 생산은 '|'로 구분되며 괄호로 그룹화될 수 있어요. 최상위 정의는 항상 먼저 주어지고 하위 부분은 들여쓰기돼요.
공통 요청·응답 구조 (Common Request and Response Structure)
모든 요청·응답은 다음 문법에서 나와요.
RequestOrResponse => Size (RequestMessage | ResponseMessage)
Size => int32
| 필드 | 설명 |
|---|---|
| message_size | 뒤따르는 요청·응답 메시지의 크기(바이트). 클라이언트는 이 4바이트 크기를 정수 N으로 읽은 뒤 그다음 N바이트의 요청을 읽고 파싱할 수 있어요. |
요청·응답 헤더 (Request and Response Headers)
서로 다른 요청·응답 버전은 대응하는 헤더의 다른 버전을 요구해요. 이 헤더 버전은 API 메시지 설명과 함께 아래에 지정돼요.
레코드 배치 (Record Batch)
레코드 배치 형식에 대한 설명은 원문 링크에서 확인할 수 있어요.
상수 (Constants)
오류 코드 (Error Codes)
서버에서 어떤 문제가 발생했는지 나타내는 숫자 코드를 사용해요. 클라이언트가 이를 예외나 적절한 오류 처리 메커니즘으로 번역할 수 있어요. 현재 사용 중인 오류 코드 표의 일부:
| 오류 | 코드 | 재시도 가능 | 설명 |
|---|---|---|---|
| UNKNOWN_SERVER_ERROR | -1 | False | 서버가 요청 처리 중 예상치 못한 오류를 겪음. |
| NONE | 0 | False | |
| OFFSET_OUT_OF_RANGE | 1 | False | 요청된 오프셋이 서버가 유지하는 오프셋 범위 밖. |
| CORRUPT_MESSAGE | 2 | True | 메시지가 CRC 체크섬 실패, 유효 크기 초과, 컴팩션 토픽의 null 키, 또는 그 밖의 손상. |
| UNKNOWN_TOPIC_OR_PARTITION | 3 | True | 서버가 이 토픽-파티션을 호스팅하지 않음. |
| INVALID_FETCH_SIZE | 4 | False | 요청된 fetch 크기가 유효하지 않음. |
| LEADER_NOT_AVAILABLE | 5 | True | 리더십 선거 중이라 이 토픽-파티션에 리더가 없음. |
| NOT_LEADER_OR_FOLLOWER | 6 | True | 리더 전용 요청은 브로커가 현재 리더가 아님을, 아무 복제본 전용 요청은 브로커가 토픽 파티션의 복제본이 아님을 나타냄. |
| REQUEST_TIMED_OUT | 7 | True | 요청이 타임아웃됨. |
| BROKER_NOT_AVAILABLE | 8 | False | 브로커를 사용할 수 없음. |
| REPLICA_NOT_AVAILABLE | 9 | True | 요청된 토픽-파티션에 대한 복제본을 사용할 수 없음. |
| MESSAGE_TOO_LARGE | 10 | False | 요청에 서버가 받아들일 최대 메시지 크기보다 큰 메시지 포함. |
| STALE_CONTROLLER_EPOCH | 11 | False | 컨트롤러가 다른 브로커로 이동함. |
| OFFSET_METADATA_TOO_LARGE | 12 | False | 오프셋 요청의 메타데이터 필드가 너무 큼. |
| NETWORK_EXCEPTION | 13 | True | 응답을 받기 전에 서버가 연결 해제됨. |
| COORDINATOR_LOAD_IN_PROGRESS | 14 | True | 코디네이터가 로딩 중이라 요청을 처리할 수 없음. |
| COORDINATOR_NOT_AVAILABLE | 15 | True | 코디네이터를 사용할 수 없음. |
| NOT_COORDINATOR | 16 | True | 올바른 코디네이터가 아님. |
| INVALID_TOPIC_EXCEPTION | 17 | False | 유효하지 않은 토픽에 연산을 수행하려 시도. |
| RECORD_LIST_TOO_LARGE | 18 | False | 요청에 서버의 구성된 세그먼트 크기보다 큰 메시지 배치 포함. |
| NOT_ENOUGH_REPLICAS | 19 | True | 필요한 수보다 동기화 복제본이 적어 메시지가 거부됨. |
참고: 이 페이지는 Kafka 프로토콜 가이드의 서론, 기본 타입, 공통 요청·응답 구조, 오류 코드 도입부를 다뤄요. 원문은 이어서 각 API(Produce, Fetch, Metadata, Offset, Group, Transaction 등)의 메시지 스키마와 레코드 배치·메시지 집합 형식을 정의하는 매우 방대한 유선 참조 문서(14,000여 줄)이며, 그 부분의 대부분은 보존해야 할 바이너리 형식 스펙이에요.
더 알아보기 (Learn more)
- Kafka 프로토콜은 TCP 위의 이진 프로토콜이고, 모든 요청·응답은 크기로 구분되며 API 키+버전으로 스키마를 식별해요.
- 클라이언트는
ApiVersionsRequest로 지원 버전을 얻어 둘 다 지원하는 가장 높은 버전을 써요. - 부트스트랩은 메타데이터 요청으로 시작해, 오류를 받으면 메타데이터를 갱신하며 요청을 올바른 리더 브로커로 보내요.