Confluent Schema Registry 디코더
Confluent Schema Registry 디코더 (Confluent Schema Registry Decoders)
Confluent Schema Registry로 Kafka에서 Avro, JSON, Protobuf 메시지를 디코딩하는 방법을 설명하는 페이지예요.
본문
Pinot은 Confluent Schema Registry로 직렬화된 Kafka 메시지를 Avro, JSON Schema, Protocol Buffers 포맷으로 디코딩하는 것을 지원해요. 이 디코더들은 레지스트리에서 스키마를 자동으로 가져와 캐시하므로, 등록된 스키마에 따라 데이터가 역직렬화되도록 보장해요.
사용 가능한 디코더
| 포맷 (Format) | 디코더 클래스 (Decoder Class) | 플러그인 (Plugin) |
|---|---|---|
| Avro | org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder |
pinot-confluent-avro |
| JSON Schema | org.apache.pinot.plugin.inputformat.json.confluent.KafkaConfluentSchemaRegistryJsonMessageDecoder |
pinot-confluent-json |
| Protocol Buffers | org.apache.pinot.plugin.inputformat.protobuf.KafkaConfluentSchemaRegistryProtoBufMessageDecoder |
pinot-confluent-protobuf |
공통 설정
모든 Confluent Schema Registry 디코더는 동일한 설정 속성을 공유해요:
| 속성 (Property) | 필수 (Required) | 기본값 (Default) | 설명 (Description) |
|---|---|---|---|
schema.registry.rest.url |
예 | — | Confluent Schema Registry REST 엔드포인트 URL |
cached.schema.map.capacity |
아니오 | 1000 | 로컬에 캐시할 최대 스키마 수 |
SSL/TLS 설정
SSL/TLS로 Schema Registry 엔드포인트에 연결하려면 schema.registry. 접두사로 속성을 추가해요:
| 속성 (Property) | 설명 (Description) |
|---|---|
schema.registry.ssl.truststore.location |
truststore 파일 경로 |
schema.registry.ssl.truststore.password |
Truststore 비밀번호 |
schema.registry.ssl.keystore.location |
keystore 파일 경로 |
schema.registry.ssl.keystore.password |
Keystore 비밀번호 |
schema.registry.ssl.key.password |
개인 키 비밀번호 |
Confluent Avro 디코더
Confluent Schema Registry에서 스키마를 관리하는 Avro 직렬화 Kafka 메시지를 디코딩해요.
{
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-avro-topic",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder",
"stream.kafka.decoder.prop.schema.registry.rest.url": "http://schema-registry:8081"
}
}
Confluent JSON Schema 디코더
Confluent의 JSON Schema 직렬화기로 직렬화된 JSON 메시지를 디코딩해요. 메시지에는 디코더가 검증을 위해 레지스트리에서 JSON Schema를 가져오는 데 사용하는 스키마 ID 헤더가 포함돼요.
{
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-json-topic",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.confluent.KafkaConfluentSchemaRegistryJsonMessageDecoder",
"stream.kafka.decoder.prop.schema.registry.rest.url": "http://schema-registry:8081"
}
}
{% hint style="info" %} JSON Schema 디코더는 들어오는 메시지를 Schema Registry에 등록된 스키마와 대조해 검증합니다. magic byte 형식과 일치하지 않는 메시지(비-Confluent 메시지)는 조용히 버려집니다. {% endhint %}
Confluent Protobuf 디코더
Confluent의 Protobuf 직렬화기로 직렬화된 Protocol Buffer 메시지를 디코딩해요. 디코더는 레지스트리에서 .proto 스키마 정의를 가져와 바이너리 페이로드를 역직렬화해요.
{
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-protobuf-topic",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.protobuf.KafkaConfluentSchemaRegistryProtoBufMessageDecoder",
"stream.kafka.decoder.prop.schema.registry.rest.url": "http://schema-registry:8081"
}
}
SSL/TLS 예시
보안 처리된 Schema Registry에 연결하려면:
{
"streamConfigs": {
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder",
"stream.kafka.decoder.prop.schema.registry.rest.url": "https://schema-registry:8082",
"stream.kafka.decoder.prop.schema.registry.ssl.truststore.location": "/path/to/truststore.jks",
"stream.kafka.decoder.prop.schema.registry.ssl.truststore.password": "changeit",
"stream.kafka.decoder.prop.schema.registry.ssl.keystore.location": "/path/to/keystore.jks",
"stream.kafka.decoder.prop.schema.registry.ssl.keystore.password": "changeit"
}
}
스키마 해석이 동작하는 방식
- 각 Confluent 직렬화 메시지는 magic byte(
0x00)로 시작하고 그 뒤에 4바이트 스키마 ID가 옴 - 디코더가 메시지 헤더에서 스키마 ID를 추출
- 스키마를 Schema Registry에서 가져와 로컬에 캐시 (
cached.schema.map.capacity까지) - 메시지 페이로드를 해석된 스키마로 역직렬화
- 필드를 Pinot의
GenericRow형식으로 추출해 수집
Confluent magic byte 접두사가 없는 메시지는 조용히 버려지고 오류로 기록돼요.
함께 보기 (See Also)
- Apache Kafka에서 수집 (Ingest from Apache Kafka) — 일반 Kafka 수집 가이드
- 스트림 수집 커넥터 (Stream Ingestion Connectors) — 전체 커넥터 설정 레퍼런스
- 지원되는 데이터 포맷 (Supported Data Formats) — 모든 지원 입력 포맷