AvroConfluent 형식

AvroConfluent 형식

AvroConfluent 형식은 Confluent Schema Registry(또는 API 호환 서비스)를 사용하여 Avro로 인코딩된 메시지의 읽기와 쓰기를 지원해요. 각 메시지는 매직 바이트(0x00)와 4바이트 빅엔디언 스키마 ID, 그리고 Avro 바이너리 데이터로 구성된 Confluent 와이어 형식을 사용합니다.

출처: 문서

본문

Input Output Alias

설명 (Description)

Apache Avro는 효율적인 데이터 처리를 위해 바이너리 인코딩을 사용하는 행 지향 직렬화 형식입니다. AvroConfluent 형식은 Confluent Schema Registry(또는 API 호환 서비스)를 사용하여 Avro로 인코딩된 메시지의 읽기와 쓰기를 지원합니다.

각 메시지는 Confluent 와이어 형식을 사용합니다: 매직 바이트(0x00) 다음에 4바이트 빅엔디언 스키마 ID, 그 다음에 Avro 바이너리 데이터. 읽을 때 ClickHouse는 레지스트리를 쿼리하여 스키마 ID를 해석합니다. 쓸 때 ClickHouse는 출력 컬럼에서 파생된 스키마를 등록하고 결과 ID를 각 행 앞에 붙입니다. 스키마는 최적의 성능을 위해 캐시됩니다.

데이터 타입 매핑 (Data type mapping)

아래 표는 Apache Avro 형식이 지원하는 모든 데이터 타입과 INSERT, SELECT 쿼리에서 해당 ClickHouse 데이터 타입을 보여줍니다.

Avro data type INSERT ClickHouse data type Avro data type SELECT
boolean , int , long , float , double Int(8\16\32) , UInt(8\16\32) int
boolean , int , long , float , double Int64 , UInt64 long
boolean , int , long , float , double Float32 float
boolean , int , long , float , double Float64 double
bytes , string , fixed , enum String bytes or string *
bytes , string , fixed FixedString(N) fixed(N)
enum Enum(8\16) enum
array(T) Array(T) array(T)
map(V, K) Map(V, K) map(string, K)
union(null, T) , union(T, null) Nullable(T) union(null, T)
union(T1, T2, …) ** Variant(T1, T2, …) union(T1, T2, …) **
null Nullable(Nothing) null
int (date) *** Date , Date32 int (date) ***
long (timestamp-millis) *** DateTime64(3) long (timestamp-millis) ***
long (timestamp-micros) *** DateTime64(6) long (timestamp-micros) ***
bytes (decimal) *** DateTime64(N) bytes (decimal) ***
int IPv4 int
fixed(16) IPv6 fixed(16)
bytes (decimal) *** Decimal(P, S) bytes (decimal) ***
string (uuid) *** UUID string (uuid) ***
fixed(16) Int128/UInt128 fixed(16)
fixed(32) Int256/UInt256 fixed(32)
record Tuple record
  • bytes가 기본이며 output_format_avro_string_column_pattern 설정으로 제어됩니다.

** Variant 타입은 null을 필드 값으로 암시적으로 받아들이므로, 예를 들어 Avro union(T1, T2, null)은 Variant(T1, T2)로 변환됩니다. 결과적으로 ClickHouse에서 Avro를 생성할 때는 스키마 추론 중 어떤 값이 실제로 null인지 알 수 없으므로, 항상 null 타입을 Avro union 타입 집합에 포함해야 합니다.

*** Avro 논리 타입(logical types)

지원되지 않는 Avro 논리 데이터 타입:

  • time-millis
  • time-micros
  • duration

형식 설정 (Format settings)

Setting Description Default
input_format_avro_allow_missing_fields 스키마에서 필드를 찾을 수 없을 때 오류를 던지는 대신 기본값을 사용할지 여부. 0
input_format_avro_null_as_default null 값을 허용하지 않는 컬럼에 null 값을 삽입할 때 오류를 던지는 대신 기본값을 사용할지 여부. 0
format_avro_schema_registry_url Confluent Schema Registry URL. 기본 인증의 경우 URL 인코딩된 자격 증명을 URL 경로에 직접 포함할 수 있어요.
format_avro_schema_registry_connection_timeout Schema Registry HTTP 클라이언트의 연결 타임아웃(초, 스키마 가져오기와 등록 모두에 사용). 0보다 커야 하며, 600(10분) 이상의 값은 599로 줄어듭니다. 1
format_avro_schema_registry_send_timeout Schema Registry HTTP 클라이언트의 전송 타임아웃(초). 0보다 커야 하며, 600(10분) 이상의 값은 599로 줄어듭니다. 1
format_avro_schema_registry_receive_timeout Schema Registry HTTP 클라이언트의 수신 타임아웃(초). 0보다 커야 하며, 600(10분) 이상의 값은 599로 줄어듭니다. 1
output_format_avro_confluent_subject 출력용: Schema Registry에 스키마가 등록되는 subject 이름. 쓸 때 필요합니다.
output_format_avro_string_column_pattern 출력용: Avro string으로 직렬화할 String 컬럼의 regexp(기본은 bytes ).

예시 (Examples)

Kafka에서 읽기 (Reading from Kafka)

Kafka 테이블 엔진을 사용하여 Avro로 인코딩된 Kafka 토픽을 읽으려면 format_avro_schema_registry_url 설정으로 스키마 레지스트리의 URL을 제공하세요.

CREATE TABLE topic1_stream
(
    field1 String,
    field2 String
)
ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'kafka-broker',
kafka_topic_list = 'topic1',
kafka_group_name = 'group1',
kafka_format = 'AvroConfluent',
format_avro_schema_registry_url = 'http://schema-registry-url';

SELECT * FROM topic1_stream;

Kafka에 쓰기 (Writing to Kafka)

Kafka 토픽에 AvroConfluent 메시지를 쓰려면 스키마 레지스트리 URL과 subject 이름을 모두 설정하세요. 스키마는 첫 쓰기 시 레지스트리에 자동 등록됩니다.

CREATE TABLE topic1_sink
(
    field1 String,
    field2 String
)
ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'kafka-broker',
kafka_topic_list = 'topic1',
kafka_format = 'AvroConfluent',
format_avro_schema_registry_url = 'http://schema-registry-url',
output_format_avro_confluent_subject = 'topic1-value';

INSERT INTO topic1_sink VALUES ('hello', 'world');
기본 인증 사용 (Using basic authentication)

스키마 레지스트리가 기본 인증을 요구하면(예: Confluent Cloud 사용 시), format_avro_schema_registry_url 설정에 URL 인코딩된 자격 증명을 제공할 수 있어요.

CREATE TABLE topic1_stream
(
    field1 String,
    field2 String
)
ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'kafka-broker',
kafka_topic_list = 'topic1',
kafka_group_name = 'group1',
kafka_format = 'AvroConfluent',
format_avro_schema_registry_url = 'https://<username>:<password>@schema-registry-url';

문제 해결 (Troubleshooting)

수집 진행 상황을 모니터링하고 Kafka 소비자 관련 오류를 디버그하려면 system.kafka_consumers 시스템 테이블을 쿼리할 수 있어요. 배포에 여러 레플리카가 있다면(예: ClickHouse Cloud) clusterAllReplicas 테이블 함수를 사용해야 합니다.

SELECT * FROM clusterAllReplicas('default',system.kafka_consumers)
ORDER BY assignments.partition_id ASC;

스키마 해석 문제가 발생하면 kafkacat과 clickhouse-local을 사용하여 문제를 해결할 수 있어요:

$ kafkacat -b kafka-broker  -C -t topic1 -o beginning -f '%s' -c 3 | clickhouse-local   --input-format AvroConfluent --format_avro_schema_registry_url 'http://schema-registry' -S "field1 Int64, field2 String"  -q 'select *  from table'
1 a
2 b
3 c

더 알아보기 (Learn more)