Confluent Avro 포맷

Confluent Avro 포맷 (Confluent Avro Format)

Format: Serialization Schema / Deserialization Schema

Avro Schema Registry(avro-confluent) 포맷은 io.confluent.kafka.serializers.KafkaAvroSerializer로 직렬화된 레코드를 읽고, io.confluent.kafka.serializers.KafkaAvroDeserializer로 다시 읽을 수 있는 레코드를 쓸 수 있게 해줍니다.

출처: 문서

본문

이 포맷으로 레코드를 읽을(역직렬화) 때 Avro writer 스키마는 레코드에 인코딩된 스키마 버전 id를 기준으로 설정된 Confluent Schema Registry에서 가져오고, reader 스키마는 테이블 스키마에서 추론합니다.

이 포맷으로 레코드를 쓸(직렬화) 때 Avro 스키마는 테이블 스키마에서 추론되며, 데이터와 함께 인코딩할 스키마 id를 얻는 데 사용됩니다. 조회는 avro-confluent.subject에 주어진 subject 아래의 설정된 Confluent Schema Registry에서 수행됩니다.

Avro Schema Registry 포맷은 Apache Kafka SQL 커넥터 또는 Upsert Kafka SQL 커넥터와 함께만 사용할 수 있습니다.

의존성 (Dependencies)

Avro Schema Registry 포맷을 사용하려면 빌드 자동화 도구(Maven, SBT 등)를 쓰는 프로젝트와 SQL Client(SQL JAR 번들) 모두에 다음 의존성이 필요합니다.

<dependency>
  <groupId>org.apache.flink</groupId>
  <artifactId>flink-avro-confluent-registry</artifactId>
  <version>2.3.0</version>
</dependency>
<dependency>
  <groupId>org.apache.flink</groupId>
  <artifactId>flink-avro</artifactId>
  <version>2.3.0</version>
</dependency>
다운로드
Download

Maven, SBT, Gradle 또는 기타 빌드 자동화 도구를 사용한다면 Confluent의 maven 저장소(https://packages.confluent.io/maven/)가 프로젝트 빌드 파일에 구성되어 있는지도 확인해야 합니다.

Avro-Confluent 포맷으로 테이블을 만드는 방법

Kafka 키로는 raw UTF-8 문자열을, Kafka 값으로는 Schema Registry에 등록된 Avro 레코드를 사용하는 테이블 예제:

CREATE TABLE user_created (
  -- Kafka raw UTF-8 키에 매핑되는 열
  the_kafka_key STRING,

  -- Kafka 값의 Avro 필드에 매핑되는 몇 개의 열
  id STRING,
  name STRING,
  email STRING
) WITH (
  'connector' = 'kafka',
  'topic' = 'user_events_example1',
  'properties.bootstrap.servers' = 'localhost:9092',

  -- Kafka 키로 UTF-8 문자열을 사용, 'the_kafka_key' 테이블 열을 사용
  'key.format' = 'raw',
  'key.fields' = 'the_kafka_key',

  'value.format' = 'avro-confluent',
  'value.avro-confluent.url' = 'http://localhost:8082',
  'value.fields-include' = 'EXCEPT_KEY'
);

다음과 같이 kafka 테이블에 데이터를 쓸 수 있습니다:

INSERT INTO user_created
SELECT
  -- kafka 키에 매핑되는 열로 user id 복제
  id as the_kafka_key,
  -- 모든 값
  id, name, email
FROM some_table;

Kafka 키와 값이 모두 Schema Registry에 Avro 레코드로 등록된 테이블 예제:

CREATE TABLE user_created (
  -- Kafka 키의 'id' Avro 필드에 매핑되는 열
  kafka_key_id STRING,

  -- Kafka 값의 Avro 필드에 매핑되는 몇 개의 열
  id STRING,
  name STRING,
  email STRING
) WITH (
  'connector' = 'kafka',
  'topic' = 'user_events_example2',
  'properties.bootstrap.servers' = 'localhost:9092',

  -- 주의: Kafka 키 맥락에서의 스키마 진화는 해시 파티셔닝 때문에 거의 호환되지 않습니다.
  'key.format' = 'avro-confluent',
  'key.avro-confluent.url' = 'http://localhost:8082',
  'key.fields' = 'kafka_key_id',

  -- 이 예제에서는 Kafka 키와 값의 Avro 타입이 모두 'id' 필드를 포함하길 원함
  -- => Kafka 키 필드에 연관된 테이블 열에 접두사를 추가해 충돌 방지
  'key.fields-prefix' = 'kafka_key_',

  'value.format' = 'avro-confluent',
  'value.avro-confluent.url' = 'http://localhost:8082',
  'value.fields-include' = 'EXCEPT_KEY',

  -- subject는 Flink 1.13부터 기본값이 있지만 재정의할 수 있습니다:
  'key.avro-confluent.subject' = 'user_events_example2-key2',
  'value.avro-confluent.subject' = 'user_events_example2-value2'
);

Kafka 값이 Schema Registry에 Avro 레코드로 등록된 upsert-kafka 커넥터를 사용하는 테이블 예제:

CREATE TABLE user_created (
  -- Kafka raw UTF-8 키에 매핑되는 열
  kafka_key_id STRING,

  -- Kafka 값의 Avro 필드에 매핑되는 몇 개의 열
  id STRING,
  name STRING,
  email STRING,

  -- upsert-kafka 커넥터는 upsert 동작을 정의하기 위해 기본 키가 필요
  PRIMARY KEY (kafka_key_id) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  'topic' = 'user_events_example3',
  'properties.bootstrap.servers' = 'localhost:9092',

  -- Kafka 키로 UTF-8 문자열 사용
  -- 이 경우 'key.fields'는 테이블 기본 키가 결정하므로 지정하지 않음
  'key.format' = 'raw',

  -- 이 예제에서는 Kafka 키와 값의 Avro 타입이 모두 'id' 필드를 포함하길 원함
  'key.fields-prefix' = 'kafka_key_',

  'value.format' = 'avro-confluent',
  'value.avro-confluent.url' = 'http://localhost:8082',
  'value.fields-include' = 'EXCEPT_KEY'
);

포맷 옵션 (Format Options)

옵션 필수 전달됨 기본값 유형 설명
format required no (없음) String 사용할 포맷을 지정합니다. 여기서는 'avro-confluent'여야 합니다.
avro-confluent.basic-auth.credentials-source optional yes (없음) String Schema Registry용 Basic auth 자격 증명 소스
avro-confluent.basic-auth.user-info optional yes (없음) String schema registry용 Basic auth 사용자 정보
avro-confluent.bearer-auth.credentials-source optional yes (없음) String Schema Registry용 Bearer auth 자격 증명 소스
avro-confluent.bearer-auth.token optional yes (없음) String Schema Registry용 Bearer auth 토큰
avro-confluent.properties optional yes (없음) Map 기본 Schema Registry에 전달되는 속성 맵입니다. Flink 설정 옵션으로 공식 노출되지 않는 옵션에 유용합니다. 단, Flink 옵션의 우선순위가 더 높다는 점에 유의하세요.
avro-confluent.ssl.keystore.location optional yes (없음) String SSL keystore 위치/파일
avro-confluent.ssl.keystore.password optional yes (없음) String SSL keystore 비밀번호
avro-confluent.ssl.truststore.location optional yes (없음) String SSL truststore 위치/파일
avro-confluent.ssl.truststore.password optional yes (없음) String SSL truststore 비밀번호
avro-confluent.schema optional no (없음) String Confluent Schema Registry에 등록되었거나 등록될 스키마입니다. 스키마가 제공되지 않으면 Flink는 테이블 스키마를 avro 스키마로 변환합니다. 제공된 스키마는 테이블 스키마와 일치해야 합니다.
avro-confluent.subject optional yes (없음) String 직렬화 중 이 포맷이 사용하는 스키마를 등록할 Confluent Schema Registry subject입니다. 기본적으로 'kafka'와 'upsert-kafka' 커넥터는 이 포맷이 value 또는 key 포맷으로 사용될 때 기본 subject 이름으로 '-value' 또는 '-key'를 사용합니다. 하지만 'filesystem' 같은 다른 커넥터는 sink로 사용될 때 subject 옵션이 필수입니다.
avro-confluent.url required yes (없음) String 스키마를 가져오고/등록할 Confluent Schema Registry의 URL입니다.

데이터 타입 매핑 (Data Type Mapping)

현재 Apache Flink는 역직렬화 중 Avro reader 스키마와 직렬화 중 Avro writer 스키마를 파생하기 위해 항상 테이블 스키마를 사용합니다. Avro 스키마를 명시적으로 정의하는 것은 아직 지원되지 않습니다. Avro와 Flink DataType 사이의 매핑은 Apache Avro Format을 참고하세요.

거기에 나열된 타입 외에도 Flink는 nullable 타입의 읽기/쓰기를 지원합니다. Flink는 nullable 타입을 Avro union(something, null)로 매핑하며, 여기서 something은 Flink 타입에서 변환된 Avro 타입입니다.

Avro 타입에 대한 자세한 내용은 Avro Specification을 참고할 수 있습니다.

더 알아보기 (Learn more)