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-millistime-microsduration
형식 설정 (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