Pulsar 바이너리 프로토콜 사양
Pulsar 바이너리 프로토콜 사양 (Pulsar Binary Protocol Specification)
Pulsar는 프로듀서/소비자와 브로커 사이의 통신에 커스텀 바이너리 프로토콜을 사용해요. 이 프로토콜은 확인(acknowledgment)과 흐름 제어(flow control) 같은 필수 기능을 지원하면서도 전송과 구현 효율을 최대화하도록 설계됐어요. 클라이언트나 프로토콜을 직접 구현하려면 이 프레임 구조와 명령 구성을 이해해야 해요. 이 글에서는 바이너리 프로토콜의 프레이밍, 메타데이터, 주요 상호 작용을 정리해 드릴게요.
출처: 문서
본문
Pulsar는 프로듀서/소비자와 브로커 사이의 통신에 커스텀 바이너리 프로토콜을 사용해요. 이 프로토콜은 확인과 흐름 제어 같은 필수 기능을 지원하면서 전송 및 구현 효율을 최대화하도록 설계됐어요.
클라이언트와 브로커는 서로 *명령(command)*을 교환해요. 명령은 바이너리 프로토콜 버퍼(protobuf) 메시지 형식으로 작성돼요. protobuf 명령의 형식은 PulsarApi.proto 파일에 지정되어 있고, 아래 Protobuf 인터페이스 섹션에도 문서화되어 있어요.
연결 공유 (Connection sharing) 서로 다른 프로듀서와 소비자에 대한 명령은 제약 없이 인터리브되어 같은 연결을 통해 전송될 수 있어요.
Pulsar 프로토콜과 관련된 모든 명령은 BaseCommand protobuf 메시지에 포함돼요. 이 메시지에는 모든 가능한 하위 명령이 선택적 필드로 포함된 Type enum이 있어요. BaseCommand 메시지는 하위 명령을 하나만 지정할 수 있어요.
프레이밍 (Framing)
protobuf는 어떤 메시지 프레임도 제공하지 않으므로, Pulsar 프로토콜의 모든 메시지에는 프레임의 크기를 지정하는 4바이트 필드가 앞에 붙어요. 단일 프레임의 최대 허용 크기는 5MB예요.
Pulsar 프로토콜은 두 가지 유형의 명령을 허용해요:
- 메시지 페이로드를 담지 않는 단순 명령.
- 메시지 게시나 전달 시 사용되는 페이로드를 담는 페이로드 명령. 페이로드 명령에서 protobuf 명령 데이터 뒤에 protobuf 메타데이터가 오고, 그 다음에 raw 형식으로 protobuf 외부에서 전달되는 페이로드가 와요. 모든 크기는 4바이트 부호 없는 빅엔디안 정수로 전달돼요.
메시지 페이로드는 효율을 위해 protobuf 형식이 아닌 raw 형식으로 전달돼요.
단순 명령
단순(페이로드 없는) 명령은 이 기본 구조를 가져요:
| 구성 요소 | 설명 | 크기(바이트) |
|---|---|---|
totalSize |
그 뒤의 모든 것을 포함한 프레임의 크기(바이트) | 4 |
commandSize |
protobuf 직렬화된 명령의 크기 | 4 |
command |
protobuf 직렬화된 명령 |
메시지 명령
페이로드 명령은 이 기본 구조를 가져요:
| 구성 요소 | 필수 또는 선택 | 설명 | 크기(바이트) |
|---|---|---|---|
totalSize |
필수 | 그 뒤의 모든 것을 포함한 프레임의 크기(바이트) | 4 |
commandSize |
필수 | protobuf 직렬화된 명령의 크기 | 4 |
command |
필수 | protobuf 직렬화된 명령 | |
magicNumberOfBrokerEntryMetadata |
선택 | 브로커 엔트리 메타데이터를 식별하는 2바이트 바이트 배열(0x0e02) 참고: magicNumberOfBrokerEntryMetadata, brokerEntryMetadataSize, brokerEntryMetadata는 함께 사용해야 해요. |
2 |
brokerEntryMetadataSize |
선택 | 브로커 엔트리 메타데이터의 크기 | 4 |
brokerEntryMetadata |
선택 | 바이너리 protobuf 메시지로 저장된 브로커 엔트리 메타데이터 | |
magicNumber |
필수 | 현재 형식을 식별하는 2바이트 바이트 배열(0x0e01) |
2 |
checksum |
필수 | 그 뒤의 모든 것에 대한 CRC32-C 체크섬 | 4 |
metadataSize |
필수 | 메시지 메타데이터의 크기 | 4 |
metadata |
필수 | 바이너리 protobuf 메시지로 저장된 메시지 메타데이터 | |
payload |
필수 | 프레임에 남은 것은 무엇이든 페이로드로 간주되며 모든 바이트 시퀀스가 될 수 있어요 |
브로커 엔트리 메타데이터
브로커 엔트리 메타데이터는 직렬화된 protobuf 메시지로 메시지 메타데이터 옆에 저장돼요. 메시지가 브로커에 도착했을 때 브로커가 생성하며, 구성된 경우 변경 없이 소비자에게 전달돼요.
| 필드 | 필수 또는 선택 | 설명 |
|---|---|---|
broker_timestamp |
선택 | 메시지가 브로커에 도착한 타임스탬프 (1970년 1월 1일 UTC 이후 밀리초 수) |
index |
선택 | 메시지의 인덱스. 브로커가 할당해요. |
브로커에 대해 브로커 엔트리 메타데이터를 사용하려면 broker.conf 파일에서 brokerEntryMetadataInterceptors 파라미터를 구성해요.
소비자에 대해 브로커 엔트리 메타데이터를 사용하려면:
- 18 이상의 클라이언트 프로토콜 버전을 사용해요.
broker.conf파일에서brokerEntryMetadataInterceptors파라미터를 구성하고exposingBrokerEntryMetadataToClientEnabled파라미터를true로 설정해요.
메시지 메타데이터
메시지 메타데이터는 애플리케이션이 지정한 페이로드 옆에 직렬화된 protobuf 메시지로 저장돼요. 메타데이터는 프로듀서가 생성하며 변경 없이 소비자에게 전달돼요.
| 필드 | 필수 또는 선택 | 설명 |
|---|---|---|
producer_name |
필수 | 메시지를 게시한 프로듀서의 이름 |
sequence_id |
필수 | 프로듀서가 할당한 메시지의 시퀀스 ID |
publish_time |
필수 | Unix 시간의 게시 타임스탬프 (1970년 1월 1일 UTC 이후 밀리초 수) |
properties |
필수 | 키/값 쌍의 시퀀스 (KeyValue 메시지 사용). Pulsar에게 특별한 의미가 없는 애플리케이션 정의 키와 값이에요. |
replicated_from |
선택 | 메시지가 복제되었음을 나타내고 메시지가 원래 게시된 클러스터의 이름을 지정해요 |
partition_key |
선택 | 파티션 토픽에서 게시할 때 키가 있으면 키의 해시로 어떤 파티션을 선택할지 결정해요. 파티션 키는 메시지 키로 사용돼요. |
compression |
선택 | 페이로드가 압축되었고 어떤 압축 라이브러리로 압축됐는지 알려줘요 |
uncompressed_size |
선택 | 압축을 사용하면 프로듀서가 uncompressed size 필드를 원본 페이로드 크기로 채워야 해요 |
num_messages_in_batch |
선택 | 이 메시지가 실제로 여러 엔트리의 배치라면 이 필드를 배치의 메시지 수로 설정해야 해요 |
배치 메시지
배치 메시지를 사용할 때 페이로드는 각각 개별 메타데이터를 가진 엔트리 목록을 포함해요. 개별 메타데이터는 SingleMessageMetadata 객체로 정의돼요.
단일 배치에서 페이로드 형식은 다음과 같아요:
| 필드 | 필수 또는 선택 | 설명 |
|---|---|---|
metadataSizeN |
필수 | 직렬화된 Protobuf 단일 메시지 메타데이터의 크기 |
metadataN |
필수 | 단일 메시지 메타데이터 |
payloadN |
필수 | 애플리케이션이 전달한 메시지 페이로드 |
각 메타데이터 필드는 다음과 같아요:
| 필드 | 필수 또는 선택 | 설명 |
|---|---|---|
properties |
필수 | 애플리케이션 정의 속성 |
partition key |
선택 | 특정 파티션으로의 해싱을 나타내는 키 |
payload_size |
필수 | 배치 안의 단일 메시지 페이로드 크기 |
압축이 활성화되면 전체 배치가 한 번에 압축돼요.
상호 작용
연결 설정
브로커에 TCP 연결(보통 포트 6650)을 연 후, 클라이언트가 세션을 시작할 책임이 있어요.
브로커로부터 Connected 응답을 받은 후 클라이언트는 연결을 사용할 준비가 된 것으로 간주할 수 있어요. 반면 브로커가 클라이언트 인증을 검증하지 못하면 브로커는 Error 명령으로 응답하고 TCP 연결을 닫아요.
예시:
message CommandConnect {
"client_version" : "Pulsar-Client-Java-v1.15.2",
"auth_method_name" : "my-authentication-plugin",
"auth_data" : "my-auth-data",
"protocol_version" : 6
}
필드:
- client_version: 문자열 기반 식별자. 형식은 강제되지 않아요.
- auth_method_name: (선택) 인증이 활성화된 경우 인증 플러그인의 이름.
- auth_data: (선택) 플러그인 특정 인증 데이터.
- protocol_version: 클라이언트가 지원하는 프로토콜 버전. 브로커는 프로토콜의 최신 개정에서 도입된 명령은 보내지 않아요. 브로커가 최소 버전을 강제할 수 있어요.
- original_principal: 프록시가 추가해요. 일반 클라이언트는 이 값을 제공할 것으로 기대되지 않아요. 설정되고 권한 부여가 활성화되면 auth_data는 broker.conf의 proxyRoles 중 하나와 매핑되어야 해요.
- original_auth_method: 프록시가 추가해요. 일반 클라이언트는 이 값을 제공할 것으로 기대되지 않아요.
- original_auth_data: 구성된 경우 프록시가 추가해요. 일반 클라이언트는 이 값을 제공할 것으로 기대되지 않아요.
- proxy_version: 프록시가 추가해요. 프록시는 이 필드가 있는 Connect 명령을 거부해요. 일반 클라이언트는 이 값을 제공할 것으로 기대되지 않아요. 브로커에서 인증·권한 부여가 활성화되면 구성된 proxyRoles 중 하나만 이 필드를 제공할 권한이 있어요. 이전 버전과의 호환성을 위해 브로커는 proxyRole이 이 필드를 제공하도록 요구하지 않아요.
message CommandConnected {
"server_version" : "Pulsar-Broker-v1.15.2",
"protocol_version" : 6
}
필드:
- server_version: 브로커 버전의 문자열 식별자.
- protocol_version: 브로커가 지원하는 프로토콜 버전. 클라이언트는 프로토콜의 최신 개정에서 도입된 명령을 보내려 시도하면 안 돼요.
Keep Alive
장기간의 네트워크 파티션 또는 원격 끝에서 TCP 연결을 끊지 않고 머신이 충돌하는 경우(전원 차단, 커널 패닉, 하드 리부트 등)를 식별하기 위해, 원격 피어의 가용성 상태를 프로빙하는 메커니즘을 도입했어요.
클라이언트와 브로커 모두 주기적으로 Ping 명령을 보내고, 타임아웃(브로커 기본 60s) 안에 Pong 응답을 받지 못하면 소켓을 닫아요.
유효한 Pulsar 클라이언트 구현은 Ping 프로브를 보낼 필요가 없지만, 원격 측이 TCP 연결을 강제로 닫지 않도록 브로커로부터 Ping을 받은 후에는 신속하게 응답해야 해요.
프로듀서
메시지를 보내려면 클라이언트가 프로듀서를 설정해야 해요. 프로듀서를 만들 때 브로커는 먼저 이 특정 클라이언트가 토픽에 게시할 권한이 있는지 확인해요.
클라이언트가 프로듀서 생성 확인을 받으면, 이전에 협상한 프로듀서 ID를 참조해 브로커에 메시지를 게시할 수 있어요.
클라이언트가 프로듀서 생성 성공이나 실패를 나타내는 응답을 받지 못하면, 프로듀서 생성을 재시도하는 명령을 보내기 전에 먼저 원래 프로듀서를 닫는 명령을 보내야 해요.
참고 프로듀서를 만들거나 연결하기 전에 먼저 토픽 조회를 수행해야 해요.
Command Producer
message CommandProducer {
"topic" : "persistent://my-property/my-cluster/my-namespace/my-topic",
"producer_id" : 1,
"request_id" : 1
}
필드:
- topic: 프로듀서를 만들려는 완전한 토픽 이름.
- producer_id: 클라이언트가 생성한 프로듀서 식별자. 같은 연결 안에서 고유해야 해요.
- request_id: 이 요청의 식별자. 응답을 원래 요청과 매칭하는 데 사용돼요. 같은 연결 안에서 고유해야 해요.
- producer_name: (선택) 프로듀서 이름을 지정하면 그 이름이 사용되고, 그렇지 않으면 브로커가 고유한 이름을 생성해요. 생성된 프로듀서 이름은 전역적으로 고유함이 보장돼요. 구현은 프로듀서를 처음 만들 때는 브로커가 새 프로듀서 이름을 생성하게 하고, 재연결 후 프로듀서를 다시 만들 때는 그것을 재사용하도록 기대돼요.
브로커는 ProducerSuccess 또는 Error 명령으로 응답해요.
Command ProducerSuccess
message CommandProducerSuccess {
"request_id" : 1,
"producer_name" : "generated-unique-producer-name"
}
필드:
- request_id: CreateProducer 요청의 원래 ID.
- producer_name: 생성된 전역적으로 고유한 프로듀서 이름 또는 클라이언트가 지정한 이름(있는 경우).
Command Send
Send 명령은 이미 존재하는 프로듀서의 맥락 안에서 새 메시지를 게시하는 데 사용돼요. 연결에 대한 프로듀서가 아직 생성되지 않았다면 브로커가 연결을 종료해요. 이 명령은 명령과 메시지 페이로드를 모두 포함하는 프레임에서 사용되며, 완전한 형식은 메시지 명령 섹션에 지정돼 있어요.
message CommandSend {
"producer_id" : 1,
"sequence_id" : 0,
"num_messages" : 1
}
필드:
- producer_id: 기존 프로듀서의 ID.
- sequence_id: 각 메시지에는 연관된 시퀀스 ID가 있으며, 0부터 시작하는 카운터로 구현될 것으로 예상돼요. 메시지의 효과적인 게시를 확인하는 SendReceipt가 그 시퀀스 ID로 참조해요.
- num_messages: (선택) 한 번에 메시지 배치를 게시할 때 사용돼요.
Command SendReceipt
메시지가 구성된 수의 복제본에 영구 저장된 후 브로커가 프로듀서에게 확인 영수증을 보내요.
message CommandSendReceipt {
"producer_id" : 1,
"sequence_id" : 0,
"message_id" : {
"ledgerId" : 123,
"entryId" : 456
}
}
필드:
- producer_id: 보내기 요청을 일으킨 프로듀서의 ID.
- sequence_id: 게시된 메시지의 시퀀스 ID.
- message_id: 게시된 메시지에 시스템이 할당한 메시지 ID. 단일 클러스터 안에서 고유해요. 메시지 ID는 ledgerId와 entryId라는 2개의 long으로 구성되며, 이 고유 ID가 BookKeeper ledger에 추가될 때 할당됨을 반영해요.
Command CloseProducer
참고 이 명령은 프로듀서 또는 브로커가 보낼 수 있어요.
CloseProducer 명령을 받으면 브로커는 프로듀서에 대해 더 이상 메시지를 받지 않고, 모든 보류 중인 메시지가 영구 저장될 때까지 기다린 후 클라이언트에 Success로 응답해요.
클라이언트가 타임아웃 안에 Producer 명령에 대한 응답을 받지 못하면, 다른 Producer 명령을 보내기 전에 먼저 CloseProducer 명령을 보내야 해요. 클라이언트는 다음 Producer 명령을 보내기 전에 CloseProducer 명령에 대한 응답을 기다릴 필요가 없어요.
브로커는 정상 페일오버를 수행할 때(예: 브로커가 재시작 중이거나, 부하 분산기가 토픽을 다른 브로커로 옮기기 위해 언로드 중일 때) 클라이언트에 CloseProducer 명령을 보낼 수 있어요.
CloseProducer를 받으면 클라이언트는 서비스 디스커버리 조회를 다시 수행하고 프로듀서를 다시 만들어야 할 것으로 예상돼요. TCP 연결은 영향을 받지 않아요.
소비자
소비자는 구독에 연결해 메시지를 소비하는 데 사용돼요. 매 재연결 후 클라이언트는 토픽을 구독해야 해요. 구독이 아직 없으면 새 구독이 생성돼요.
참고 소비자를 만들거나 연결하기 전에 먼저 토픽 조회를 수행해야 해요.
클라이언트가 소비자 생성 성공이나 실패를 나타내는 응답을 받지 못하면, 소비자 생성을 재시도하는 명령을 보내기 전에 먼저 원래 소비자를 닫는 명령을 보내야 해요.
흐름 제어
소비자가 준비된 후 클라이언트는 브로커가 메시지를 푸시하도록 권한을 부여해야 해요. 이는 Flow 명령으로 이루어져요.
Flow 명령은 소비자에게 메시지를 보낼 수 있는 추가 *퍼밋(permit)*을 부여해요. 전형적인 소비자 구현은 애플리케이션이 소비할 준비가 되기 전에 이 메시지들을 큐에 쌓아두기 위해 큐를 사용해요.
애플리케이션이 큐에 있는 메시지의 절반을 디큐한 후에 소비자는 브로커에게 더 많은 메시지를 요청하는 퍼밋(큐 메시지의 절반과 같은 수)을 보내요.
예를 들어 큐 크기가 1000이고 소비자가 큐에 있는 500개 메시지를 소비한다면, 소비자는 브로커에 500개 메시지를 요청하는 퍼밋을 보내요.
Command Subscribe
message CommandSubscribe {
"topic" : "persistent://my-property/my-cluster/my-namespace/my-topic",
"subscription" : "my-subscription-name",
"subType" : "Exclusive",
"consumer_id" : 1,
"request_id" : 1
}
필드:
- topic: 소비자를 만들려는 완전한 토픽 이름.
- subscription: 구독 이름.
- subType: 구독 유형: Exclusive, Shared, Failover, Key_Shared.
- consumer_id: 클라이언트가 생성한 소비자 식별자. 같은 연결 안에서 고유해야 해요.
- request_id: 이 요청의 식별자. 응답을 원래 요청과 매칭하는 데 사용돼요. 같은 연결 안에서 고유해야 해요.
- consumer_name: (선택) 클라이언트가 소비자 이름을 지정할 수 있어요. 이 이름은 stats에서 특정 소비자를 추적하는 데 사용될 수 있어요. 또한 Failover 구독 유형에서 이 이름은 어떤 소비자가 마스터(메시지를 받는 소비자)로 선출될지 결정하는 데 사용돼요. 소비자가 소비자 이름으로 정렬되고 첫 번째가 마스터로 선출돼요.
Command Flow
message CommandFlow {
"consumer_id" : 1,
"messagePermits" : 1000
}
필드:
- consumer_id: 이미 설정된 소비자의 ID.
- messagePermits: 더 많은 메시지를 푸시하도록 브로커에 부여할 추가 퍼밋 수.
Command Message
Message 명령은 브로커가 주어진 퍼밋의 한도 내에서 기존 소비자에게 메시지를 푸시하는 데 사용돼요.
이 명령은 메시지 페이로드도 포함하는 프레임에서 사용되며, 완전한 형식은 메시지 명령 섹션에 지정돼 있어요.
message CommandMessage {
"consumer_id" : 1,
"message_id" : {
"ledgerId" : 123,
"entryId" : 456
}
}
Command Ack
Ack는 주어진 메시지가 애플리케이션에 의해 성공적으로 처리되었고 브로커가 버릴 수 있다는 것을 브로커에 알리는 데 사용돼요.
또한 브로커는 확인된 메시지에 기반해 소비자 위치를 유지해요.
message CommandAck {
"consumer_id" : 1,
"ack_type" : "Individual",
"message_id" : {
"ledgerId" : 123,
"entryId" : 456
}
}
필드:
- consumer_id: 이미 설정된 소비자의 ID.
- ack_type: 확인 유형: Individual 또는 Cumulative.
- message_id: 확인할 메시지의 ID.
- validation_error: (선택) 소비자가 다음 이유로 메시지를 버렸음을 나타냄: UncompressedSizeCorruption, DecompressionError, ChecksumMismatch, BatchDeSerializeError.
- properties: (선택) 예약된 구성 항목.
- txnid_most_bits: (선택) 트랜잭션 코디네이터 ID와 동일. txnid_most_bits와 txnid_least_bits가 트랜잭션을 고유하게 식별해요.
- txnid_least_bits: (선택) 트랜잭션 코디네이터에서 연 트랜잭션의 ID. txnid_most_bits와 txnid_least_bits가 트랜잭션을 고유하게 식별해요.
- request_id: (선택) 응답과 타임아웃 처리를 위한 ID.
Command AckResponse
AckResponse는 클라이언트가 보낸 확인 요청에 대한 브로커의 응답이에요. 요청에서 보낸 consumer_id를 포함해요. 트랜잭션을 사용하면 요청에서 보낸 트랜잭션 ID와 요청 ID를 모두 포함해요. 클라이언트는 요청 ID에 따라 특정 요청을 완료해요. error 필드가 설정되면 요청이 실패했음을 나타내요.
리다이렉션이 있는 AckResponse 예시:
message CommandAckResponse {
"consumer_id" : 1,
"txnid_least_bits" = 0,
"txnid_most_bits" = 1,
"request_id" = 5
}
Command CloseConsumer
참고 이 명령은 프로듀서 또는 브로커가 보낼 수 있어요.
이 명령은 CloseProducer와 동일하게 동작해요.
클라이언트가 타임아웃 안에 Subscribe 명령에 대한 응답을 받지 못하면, 다른 Subscribe 명령을 보내기 전에 먼저 CloseConsumer 명령을 보내야 해요. 클라이언트는 다음 Subscribe 명령을 보내기 전에 CloseConsumer 명령에 대한 응답을 기다릴 필요가 없어요.
Command RedeliverUnacknowledgedMessages
소비자는 브로커에게 그 특정 소비자에게 푸시되었지만 아직 확인되지 않은 보류 중인 메시지의 일부 또는 전체를 재전달하도록 요청할 수 있어요.
protobuf 객체는 소비자가 재전달되길 원하는 메시지 ID 목록을 받아요. 목록이 비어 있으면 브로커는 모든 보류 중인 메시지를 재전달해요.
재전달 시 메시지는 같은 소비자에게 보내지거나, 공유 구독의 경우 사용 가능한 모든 소비자에 퍼질 수 있어요.
Command ReachedEndOfTopic
이 명령은 토픽이 "종료"되고 구독의 모든 메시지가 확인될 때마다 브로커가 특정 소비자에게 보내요.
클라이언트는 이 명령을 사용해 소비자로부터 더 이상 메시지가 오지 않는다는 것을 애플리케이션에 알려야 해요.
Command ConsumerStats
이 명령은 클라이언트가 브로커로부터 Subscriber 및 Consumer 수준의 stats를 검색하기 위해 보내요.
필드:
- request_id: 요청의 ID. 요청과 응답을 연관시키는 데 사용돼요.
- consumer_id: 이미 설정된 소비자의 ID.
Command ConsumerStatsResponse
이것은 클라이언트의 ConsumerStats 요청에 대한 브로커의 응답이에요. 요청에서 보낸 consumer_id의 Subscriber 및 Consumer 수준 stats를 포함해요. error_code 또는 error_message 필드가 설정되면 요청이 실패했음을 나타내요.
Command Unsubscribe
이 명령은 클라이언트가 연관된 토픽에서 consumer_id를 구독 해제하기 위해 보내요.
필드:
- request_id: 요청의 ID.
- consumer_id: 구독 해제해야 할 이미 설정된 소비자의 ID.
서비스 디스커버리
토픽 조회
토픽 조회는 클라이언트가 프로듀서나 소비자를 만들거나 재연결할 때마다 수행해야 해요. 조회는 곧 사용할 토픽을 서비스하는 특정 브로커를 발견하는 데 사용돼요.
조회는 admin API 문서에 설명된 대로 REST 호출로 수행할 수 있어요.
Pulsar-1.16부터는 바이너리 프로토콜 안에서도 조회를 수행할 수 있어요.
예를 들어 pulsar://broker.example.com:6650에서 실행 중인 서비스 디스커버리 컴포넌트가 있다고 가정해 볼게요.
개별 브로커는 pulsar://broker-1.example.com:6650, pulsar://broker-2.example.com:6650, ...에서 실행될 거예요.
클라이언트는 디스커버리 서비스 호스트에 대한 연결을 사용해 LookupTopic 명령을 발행할 수 있어요. 응답은 연결할 브로커 호스트명이거나, 조회를 재시도할 브로커 호스트명일 수 있어요.
LookupTopic 명령은 이미 Connect/Connected 초기 핸드셰이크를 거친 연결에서만 사용해야 해요.
message CommandLookupTopic {
"topic" : "persistent://my-property/my-cluster/my-namespace/my-topic",
"request_id" : 1,
"authoritative" : false
}
필드:
- topic: 조회할 토픽 이름.
- request_id: 응답과 함께 전달될 요청의 ID.
- authoritative: 초기 조회 요청은 false를 사용해야 해요. 리다이렉트 응답을 따를 때 클라이언트는 응답에 포함된 것과 같은 값을 전달해야 해요.
LookupTopicResponse
성공적인 조회의 응답 예시:
message CommandLookupTopicResponse {
"request_id" : 1,
"response" : "Connect",
"brokerServiceUrl" : "pulsar://broker-1.example.com:6650",
"brokerServiceUrlTls" : "pulsar+ssl://broker-1.example.com:6651",
"authoritative" : true
}
리다이렉션이 있는 조회 응답 예시:
message CommandLookupTopicResponse {
"request_id" : 1,
"response" : "Redirect",
"brokerServiceUrl" : "pulsar://broker-2.example.com:6650",
"brokerServiceUrlTls" : "pulsar+ssl://broker-2.example.com:6651",
"authoritative" : true
}
이 두 번째 경우에는 broker-2.example.com에 LookupTopic 명령 요청을 다시 발행해야 하고, 이 브로커가 조회 요청에 최종 답변을 할 수 있어요.
파티션 토픽 디스커버리
파티션 토픽 메타데이터 디스커버리는 토픽이 "파티션 토픽"인지와 몇 개의 파티션이 설정되었는지 알아내는 데 사용돼요.
토픽이 "partitioned"로 표시되면 클라이언트는 각 파티션에 대해 partition-X 접미사를 사용해 프로듀서나 소비자 여러 개를 만들어야 할 것으로 예상돼요.
이 정보는 프로듀서나 소비자를 처음 만들 때만 검색하면 돼요. 재연결 후에는 할 필요가 없어요.
파티션 토픽 메타데이터 디스커버리는 토픽 조회와 매우 유사하게 동작해요. 클라이언트가 서비스 디스커버리 주소로 요청을 보내고 응답에 실제 메타데이터가 포함돼요.
Command PartitionedTopicMetadata
message CommandPartitionedTopicMetadata {
"topic" : "persistent://my-property/my-cluster/my-namespace/my-topic",
"request_id" : 1
}
필드:
- topic: 파티션 메타데이터를 확인할 토픽.
- request_id: 응답과 함께 전달될 요청의 ID.
Command PartitionedTopicMetadataResponse
메타데이터가 있는 응답 예시:
message CommandPartitionedTopicMetadataResponse {
"request_id" : 1,
"response" : "Success",
"partitions" : 32
}
Protobuf 인터페이스
모든 Pulsar Protobuf 정의는 여기에서 찾을 수 있어요.
더 알아보기 (Learn more)
- 개발자를 위한 Pulsar — 다른 개발 문서를 살펴봐요.
- 클라이언트 라이브러리 — 공식 클라이언트 구현을 이해해요.
- 토픽 관리 — 토픽 소유 브로커 조회를 다뤄요.