Confluent Schema Registry 디코더

Confluent Schema Registry 디코더 (Confluent Schema Registry Decoders)

Confluent Schema Registry로 Kafka에서 Avro, JSON, Protobuf 메시지를 디코딩하는 방법을 설명하는 페이지예요.

출처: Confluent Schema Registry Decoders

본문

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"
  }
}

스키마 해석이 동작하는 방식

  1. 각 Confluent 직렬화 메시지는 magic byte(0x00)로 시작하고 그 뒤에 4바이트 스키마 ID가 옴
  2. 디코더가 메시지 헤더에서 스키마 ID를 추출
  3. 스키마를 Schema Registry에서 가져와 로컬에 캐시 (cached.schema.map.capacity까지)
  4. 메시지 페이로드를 해석된 스키마로 역직렬화
  5. 필드를 Pinot의 GenericRow 형식으로 추출해 수집

Confluent magic byte 접두사가 없는 메시지는 조용히 버려지고 오류로 기록돼요.

함께 보기 (See Also)

더 알아보기 (Learn more)