Kafka 테이블 엔진
Kafka 테이블 엔진
ClickHouse Cloud를 사용 중이라면 대신 ClickPipes를 사용하는 것이 좋아요. ClickPipes는 프라이빗 네트워크 연결, 수집과 클러스터 리소스의 독립적인 확장, Kafka 데이터를 ClickHouse로 스트리밍하기 위한 종합적인 모니터링을 네이티브로 지원해요.
Kafka 엔진은 다음을 할 수 있게 해줘요.
- 데이터 흐름을 게시하거나 구독하기
- 장애 허용(fault-tolerant) 저장소 구성하기
- 스트림을 사용 가능해지면 처리하기
출처: 문서
본문
테이블 생성하기
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [ALIAS expr1],
name2 [type2] [ALIAS expr2],
...
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'host:port',
kafka_topic_list = 'topic1,topic2,...',
kafka_group_name = 'group_name',
kafka_format = 'data_format'[,]
[kafka_security_protocol = '',]
[kafka_sasl_mechanism = '',]
[kafka_sasl_username = '',]
[kafka_sasl_password = '',]
[kafka_autodetect_client_rack = '',]
[kafka_schema = '',]
[kafka_num_consumers = N,]
[kafka_max_block_size = 0,]
[kafka_skip_broken_messages = N,]
[kafka_commit_every_batch = 0,]
[kafka_client_id = '',]
[kafka_poll_timeout_ms = 0,]
[kafka_poll_max_batch_size = 0,]
[kafka_flush_interval_ms = 0,]
[kafka_consumer_reschedule_ms = 0,]
[kafka_thread_per_consumer = 0,]
[kafka_handle_error_mode = 'default',]
[kafka_commit_on_select = false,]
[kafka_consumer_acquire_timeout_ms = 30000,]
[kafka_max_rows_per_message = 1,]
[kafka_compression_codec = '',]
[kafka_compression_level = -1,]
[kafka_partition_shard_num = '',]
[kafka_shard_count = 0];
필수 매개변수:
kafka_broker_list— 브로커의 쉼표로 구분된 목록 (예:localhost:9092)kafka_topic_list— Kafka 토픽 목록kafka_group_name— Kafka 소비자 그룹. 읽기 한계(margin)는 그룹마다 별도로 추적돼요. 클러스터에서 메시지가 중복되지 않게 하려면 어디서나 같은 그룹 이름을 사용해요kafka_format— 메시지 형식. SQLFORMAT함수와 같은 표기법을 사용해요 (예:JSONEachRow). 자세한 내용은 Formats 섹션을 참고해요
선택 매개변수:
kafka_security_protocol- 브로커와 통신하는 데 사용하는 프로토콜. 가능한 값:plaintext,ssl,sasl_plaintext,sasl_sslkafka_sasl_mechanism- 인증에 사용할 SASL 메커니즘. 가능한 값:GSSAPI,PLAIN,SCRAM-SHA-256,SCRAM-SHA-512,OAUTHBEARER,AWS_MSK_IAMkafka_aws_region- MSK IAM 인증용 AWS 리전. 지정하지 않으면 브로커 주소에서 자동 감지돼요. 리전 정보를 포함하지 않는 PrivateLink 별칭이나 커스텀 DNS 호스트명을 사용할 때는 명시적으로 지정해요. 기본값: 비어 있음(자동 감지)kafka_sasl_username-PLAIN및SASL-SCRAM-..메커니즘에 사용할 SASL 사용자 이름kafka_sasl_password-PLAIN및SASL-SCRAM-..메커니즘에 사용할 SASL 비밀번호kafka_schema— 형식이 스키마 정의를 요구할 때 반드시 사용해야 하는 매개변수. 예를 들어 Cap'n Proto는 스키마 파일의 경로와 루트schema.capnp:Message객체의 이름을 요구해요kafka_schema_registry_skip_bytes— 엔벨로프 헤더가 있는 스키마 레지스트리를 사용할 때 각 메시지의 시작 부분에서 건너뛸 바이트 수 (예: 19바이트 엔벨로프를 포함하는 AWS Glue Schema Registry). 범위:[0, 255]. 기본값:0kafka_num_consumers— 테이블당 소비자 수. 한 소비자의 처리량이 부족하면 더 많은 소비자를 지정해요. 총 소비자 수는 토픽의 파티션 수를 초과하면 안 돼요 (파티션당 한 소비자만 할당될 수 있으므로), ClickHouse가 배포된 서버의 물리 코어 수보다 크면 안 돼요. 기본값:1kafka_max_block_size— poll의 최대 배치 크기(메시지 수). 기본값: max_insert_block_sizekafka_skip_broken_messages— 스키마와 맞지 않는 메시지에 대한 Kafka 메시지 파서 허용치.kafka_skip_broken_messages = N이면 엔진은 파싱할 수 없는 Kafka 메시지 N개를 건너뛰어요 (메시지 하나는 데이터 행 하나와 같음). 기본값:0kafka_commit_every_batch— 전체 블록을 쓴 뒤 한 번 커밋하는 대신 소비·처리된 각 배치를 커밋해요. 기본값:0kafka_client_id— 클라이언트 식별자. 기본값은 비어 있음kafka_poll_timeout_ms— Kafka에서 단일 poll의 타임아웃. 기본값: stream_poll_timeout_mskafka_poll_max_batch_size— 단일 Kafka poll에서 가져올 최대 메시지 수. 기본값: max_block_sizekafka_flush_interval_ms— Kafka에서 데이터를 flush하는 타임아웃. 기본값: stream_flush_interval_mskafka_consumer_reschedule_ms— Kafka 스트림 처리가 정체될 때(예: 소비할 메시지가 없을 때) 재스케줄 간격. 이 설정은 소비자가 poll 재시도를 하기 전의 지연을 제어해요.kafka_consumers_pool_ttl_ms를 초과하면 안 돼요. 기본값:500밀리초kafka_thread_per_consumer— 각 소비자에 독립적인 스레드를 제공해요. 활성화하면 각 소비자가 독립적으로 병렬로 데이터를 flush해요 (그렇지 않으면 여러 소비자의 행이 한 블록을 형성하도록 합쳐져요). 기본값:0kafka_handle_error_mode— Kafka 엔진의 오류 처리 방식. 가능한 값: default(메시지 파싱에 실패하면 예외가 발생), stream(예외 메시지와 원시 메시지가 가상 컬럼_error와_raw_message에 저장), dead_letter_queue(오류 관련 데이터가 system.dead_letter_queue에 저장)kafka_commit_on_select— select 쿼리를 만들 때 메시지를 커밋해요. 기본값:falsekafka_consumer_acquire_timeout_ms—Kafka2테이블(원거리 Keeper 기반 오프셋 저장 포함)에 대한 직접SELECT쿼리 중 Kafka 소비자를 획득하는 밀리초 단위 타임아웃. 같은 테이블에 대해 동시 직접SELECT쿼리가 여러 개 실행되면 각각 소비자가 사용 가능해질 때까지 기다려야 해요. 쿼리가 서로 다른 소비자 하위 집합을 보유할 때 타임아웃은 교착 상태를 방지해요. 기본값:30000kafka_max_rows_per_message— 행 기반 형식에서 하나의 Kafka 메시지에 쓰는 최대 행 수. 기본값:1kafka_autodetect_client_rack—librdkafka의client.rack매개변수에 자동으로 설정해 가장 가까운 Kafka 복제본을 선호해요. 지원 소스:AWS_ZONE_ID— AWS IMDSv2 가용 영역 ID, 예:euc1-az1AWS_ZONE_NAME— AWS IMDSv2 가용 영역 이름, 예:eu-central-1aGCP_ZONE— GCP 메타데이터 서비스 영역, 예:europe-central2-aCLICKHOUSE— ClickHouse 내부 감지 사용. 클라우드 메타데이터나 구성에 의존할 수 있어요AWS_ZONE_NAME_THEN_GCP_ZONE—AWS_ZONE_NAME을 시도한 다음GCP_ZONE을 시도 기본값: 빈 문자열, 비활성화. 팁: 환경마다 다른 가용 영역 형식을 사용해요. Amazon MSK는 일반적으로 영역 ID를 사용하므로AWS_ZONE_ID를 선호해요. Confluent Cloud는 일반적으로 영역 이름을 사용하므로AWS_ZONE_NAME을 선호해요. 확실하지 않으면AWS_ZONE_NAME_THEN_GCP_ZONE을 사용하거나 클러스터의broker.rack값을 확인해요. 참고: Kafka 브로커는broker.rack과replica.selector.class=org.apache.kafka.common.replica.RackAwareReplicaSelector로 구성되어야 해요
kafka_compression_codec— 메시지 생성에 사용되는 압축 코덱. 지원: 빈 문자열,none,gzip,snappy,lz4,zstd. 빈 문자열이면 테이블이 압축 코덱을 설정하지 않아서 구성 파일의 값이나librdkafka의 기본값이 사용돼요. 기본값: 빈 문자열kafka_compression_level— kafka_compression_codec이 선택한 알고리즘의 압축 수준 매개변수. 값이 높을수록 더 좋은 압축률이지만 CPU 사용량이 더 많아져요. 사용 가능한 범위는 알고리즘에 따라 다름:[0-9](gzip),[0-12](lz4),0만 (snappy),[0-12](zstd),-1= 코덱별 기본 압축 수준. 기본값:-1kafka_map_virtual_columns_on_write— 활성화하면 테이블 스키마의_key,_timestamp,_headers.name,_headers.value특수 이름 컬럼이INSERT시 해당 Kafka 메시지 메타데이터에 매핑되고 메시지 페이로드에서 제외돼요. Kafka 메시지 메타데이터로 컬럼 매핑 참고. 기본값:falsekafka_partition_shard_num— 정적 파티션-샤드 친화성을 위한 현재 샤드 번호. 1과kafka_shard_count사이여야 해요. 파티션은partition_id % kafka_shard_count == kafka_partition_shard_num - 1공식으로 할당돼요. 매크로 확장을 지원해요 (예:'{shard}').kafka_shard_count와 함께 사용해야 해요. StorageKafka2에서만 지원 (kafka_keeper_path와kafka_replica_name필요). 기본값:''(비활성화)kafka_shard_count— 소비에 참여하는 총 샤드 수.kafka_partition_shard_num과 함께 파티션을 정적으로 할당하는 데 사용돼요.kafka_partition_shard_num과 함께 사용해야 해요. StorageKafka2에서만 지원. 기본값:0(비활성화)
예시:
CREATE TABLE queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka('localhost:9092', 'topic', 'group1', 'JSONEachRow');
SELECT * FROM queue LIMIT 5;
CREATE TABLE queue2 (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka SETTINGS kafka_broker_list = 'localhost:9092',
kafka_topic_list = 'topic',
kafka_group_name = 'group1',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;
CREATE TABLE queue3 (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka('localhost:9092', 'topic', 'group1')
SETTINGS kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;
Kafka 테이블 엔진은 기본값이 있는 컬럼을 지원하지 않아요. 기본값이 있는 컬럼이 필요하다면 materialized view 수준에서 추가할 수 있어요 (아래 참고).
설명
전달된 메시지는 자동으로 추적되므로 그룹의 각 메시지는 한 번만 계산돼요. 데이터를 두 번 얻고 싶다면 다른 그룹 이름으로 테이블 복사본을 만들어요. 그룹은 유연하며 클러스터에서 동기화돼요. 예를 들어 토픽 10개와 클러스터의 테이블 복사본 5개가 있으면 각 복사본이 토픽 2개를 가져요. 복사본 수가 바뀌면 토픽이 복사본 간에 자동으로 재분배돼요. 자세한 내용은 http://kafka.apache.org/intro를 참고해요.
각 Kafka 토픽이 전용 소비자 그룹을 갖도록 권장해요. 특히 테스트나 스테이징처럼 토픽이 동적으로 생성·삭제될 수 있는 환경에서는 토픽과 그룹의 배타적 페어링을 보장해요.
SELECT는 (디버깅 외에는) 메시지 읽기에 특히 유용하지 않아요. 각 메시지는 한 번만 읽을 수 있기 때문이에요. materialized view를 사용해 실시간 스레드를 만드는 것이 더 실용적이에요. 그러려면:
- 엔진을 사용해 Kafka 소비자를 만들고 그것을 데이터 스트림으로 간주해요
- 원하는 구조로 테이블을 만들어요
- 엔진의 데이터를 변환해 앞서 만든 테이블에 넣는 materialized view를 만들어요
MATERIALIZED VIEW가 엔진에 조인되면 백그라운드에서 데이터 수집을 시작해요. 이렇게 하면 Kafka에서 메시지를 계속 받아 SELECT로 필요한 형식으로 변환할 수 있어요. 하나의 Kafka 테이블은 원하는 만큼 많은 materialized view를 가질 수 있어요. view는 Kafka 테이블에서 직접 데이터를 읽지 않고 새 레코드(블록 단위)를 받으므로, 여러 테이블에 다른 세부 수준으로 쓸 수 있어요 (그룹화-집계 포함, 포함하지 않고).
named collections를 사용해 연결 매개변수를 저장하는 예시:
CREATE NAMED COLLECTION kafka_creds AS
kafka_broker_list = 'localhost:9092',
kafka_topic_list = 'topic',
kafka_group_name = 'group1',
kafka_format = 'JSONEachRow';
CREATE TABLE queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka(kafka_creds);
CREATE TABLE daily (
day Date,
level String,
total UInt64
) ENGINE = SummingMergeTree
PARTITION BY toYYYYMM(day)
ORDER BY (day, level);
CREATE MATERIALIZED VIEW consumer TO daily
AS SELECT toDate(toDateTime(timestamp)) AS day, level, count() AS total
FROM queue GROUP BY day, level;
SELECT level, sum(total) FROM daily GROUP BY level;
성능을 향상시키기 위해 수신된 메시지는 max_insert_block_size 크기의 블록으로 그룹화돼요. 블록이 stream_flush_interval_ms 밀리초 안에 형성되지 않으면 블록의 완전성과 관계없이 데이터가 테이블로 flush돼요.
토픽 데이터 수신을 중지하거나 변환 로직을 바꾸려면 materialized view를 분리해요:
DETACH TABLE consumer;
ATTACH TABLE consumer;
ALTER로 대상 테이블을 바꾸고 싶다면 대상 테이블과 view의 데이터 사이 괴리를 피하기 위해 materialized view를 비활성화하는 것을 권장해요.
구성
GraphiteMergeTree와 유사하게 Kafka 엔진은 ClickHouse 구성 파일을 사용한 확장 구성을 지원해요. 사용할 수 있는 구성 키는 두 가지예요: 전역(<kafka> 아래)과 토픽 수준(<kafka><kafka_topic> 아래). 전역 구성이 먼저 적용되고 그 다음 토픽 수준 구성이 적용돼요(존재하는 경우).
<kafka>
<!-- Global configuration options for all tables of Kafka engine type -->
<debug>cgrp</debug>
<statistics_interval_ms>3000</statistics_interval_ms>
<kafka_topic>
<name>logs</name>
<statistics_interval_ms>4000</statistics_interval_ms>
</kafka_topic>
<!-- Settings for consumer -->
<consumer>
<auto_offset_reset>smallest</auto_offset_reset>
<kafka_topic>
<name>logs</name>
<fetch_min_bytes>100000</fetch_min_bytes>
</kafka_topic>
<kafka_topic>
<name>stats</name>
<fetch_min_bytes>50000</fetch_min_bytes>
</kafka_topic>
</consumer>
<!-- Settings for producer -->
<producer>
<kafka_topic>
<name>logs</name>
<retry_backoff_ms>250</retry_backoff_ms>
</kafka_topic>
<kafka_topic>
<name>stats</name>
<retry_backoff_ms>400</retry_backoff_ms>
</kafka_topic>
</producer>
</kafka>
가능한 구성 옵션 목록은 librdkafka 구성 참조를 참고해요. ClickHouse 구성에서는 점 대신 밑줄(_)을 사용해요. 예를 들어 check.crcs=true는 <check_crcs>true</check_crcs>가 돼요.
AWS MSK IAM 인증
AWS MSK IAM 인증은 ClickHouse가 AWS S3 지원으로 빌드되어야 해요. AWS MSK는 IAM 기반 인증을 지원하므로 별도의 사용자 이름과 비밀번호를 관리하는 대신 AWS 자격 증명으로 Kafka 클러스터에 연결할 수 있어요.
기본 설정: 테이블 설정에서 kafka_sasl_mechanism = 'AWS_MSK_IAM'을 설정해요:
CREATE TABLE msk_queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'b-1.mycluster.kafka.us-east-1.amazonaws.com:9098',
kafka_topic_list = 'my-topic',
kafka_group_name = 'my-group',
kafka_format = 'JSONEachRow',
kafka_sasl_mechanism = 'AWS_MSK_IAM';
AWS 리전은 패턴 매칭으로 브로커 엔드포인트에서 자동으로 추출돼요:
- Provisioned MSK:
b-X.cluster.kafka.<region>.amazonaws.com:9098 - Serverless MSK:
boot-X.kafka-serverless.<region>.amazonaws.com:9098 - VPC Endpoint:
vpce-X.kafka.<region>.vpce.amazonaws.com:9098
AWS 자격 증명: 자격 증명은 존재할 때 항상 ~/.aws/credentials와 ~/.aws/config(AWS 프로필 파일)에서 로드돼요. EC2 인스턴스 프로필, 환경 변수(AWS_ACCESS_KEY_ID 등), ECS 태스크 역할, 기타 자동 자격 증명 소스도 활성화하려면 서버 구성에 추가해요:
<kafka>
<use_environment_credentials>true</use_environment_credentials>
</kafka>
이 설정은 서버 관리자만 구성할 수 있어요. 기본값: false.
PrivateLink 및 커스텀 DNS: 리전 정보를 포함하지 않는 PrivateLink 별칭이나 커스텀 DNS 호스트명을 사용할 때 AWS 리전을 명시적으로 지정해요:
CREATE TABLE msk_privatelink_queue (
timestamp UInt64,
level String,
message String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'my-privatelink-alias.internal.example.com:9098',
kafka_topic_list = 'my-topic',
kafka_group_name = 'my-group',
kafka_format = 'JSONEachRow',
kafka_sasl_mechanism = 'AWS_MSK_IAM',
kafka_aws_region = 'us-east-1';
IAM 권한: 소비자 권한(메시지 읽기용):
{
"Version": "2012-10-17",
"Statement": [{
"Effect": "Allow",
"Action": [
"kafka-cluster:Connect",
"kafka-cluster:DescribeTopic",
"kafka-cluster:ReadData",
"kafka-cluster:AlterGroup",
"kafka-cluster:DescribeGroup"
],
"Resource": [
"arn:aws:kafka:REGION:ACCOUNT:cluster/CLUSTER_NAME/*",
"arn:aws:kafka:REGION:ACCOUNT:topic/CLUSTER_NAME/TOPIC_NAME/*",
"arn:aws:kafka:REGION:ACCOUNT:group/CLUSTER_NAME/CONSUMER_GROUP/*"
]
}]
}
프로듀서 권한(메시지 쓰기용):
{
"Version": "2012-10-17",
"Statement": [{
"Effect": "Allow",
"Action": [
"kafka-cluster:Connect",
"kafka-cluster:DescribeTopic",
"kafka-cluster:WriteData"
],
"Resource": [
"arn:aws:kafka:REGION:ACCOUNT:cluster/CLUSTER_NAME/*",
"arn:aws:kafka:REGION:ACCOUNT:topic/CLUSTER_NAME/TOPIC_NAME/*"
]
}]
}
Kerberos 지원
Kerberos 인식 Kafka를 처리하려면 security_protocol 자식 요소를 sasl_plaintext 값으로 추가해요. Kerberos 티켓 수여 티켓(ticket-granting ticket)이 OS 기능에 의해 얻어지고 캐시되면 충분해요. ClickHouse는 keytab 파일을 사용해 Kerberos 자격 증명을 유지할 수 있어요. sasl_kerberos_service_name, sasl_kerberos_keytab, sasl_kerberos_principal 자식 요소를 고려해요.
예시:
<!-- Kerberos-aware Kafka -->
<kafka>
<security_protocol>SASL_PLAINTEXT</security_protocol>
<sasl_kerberos_keytab>/home/kafkauser/kafkauser.keytab</sasl_kerberos_keytab>
<sasl_kerberos_principal>kafkauser/[email protected]</sasl_kerberos_principal>
</kafka>
가상 컬럼
_topic— Kafka 토픽. 데이터 타입:LowCardinality(String)_key— 메시지의 키. 데이터 타입:String_offset— 메시지의 오프셋. 데이터 타입:UInt64_timestamp— 메시지의 타임스탬프. 데이터 타입:Nullable(DateTime)_timestamp_ms— 메시지의 밀리초 단위 타임스탬프. 데이터 타입:Nullable(DateTime64(3))_partition— Kafka 토픽의 파티션. 데이터 타입:UInt64_headers.name— 메시지 헤더 키의 배열. 데이터 타입:Array(String)_headers.value— 메시지 헤더 값의 배열. 데이터 타입:Array(String)
kafka_handle_error_mode='stream'일 때 추가 가상 컬럼:
_raw_message- 성공적으로 파싱되지 못한 원시 메시지. 데이터 타입:String_error- 실패한 파싱 중 발생한 예외 메시지. 데이터 타입:String
참고: _raw_message와 _error 가상 컬럼은 파싱 중 예외가 발생한 경우에만 채워지며, 메시지가 성공적으로 파싱되면 항상 비어 있어요.
Kafka 메시지 메타데이터로 컬럼 매핑
INSERT INTO로 메시지를 생성할 때 Kafka 엔진은 테이블에 해당 컬럼이 존재하면 항상 _key(타입 String)라는 컬럼을 Kafka 메시지 키로, _timestamp(타입 DateTime)라는 컬럼을 Kafka 메시지 타임스탬프로 사용해요. 기본적으로 이 컬럼들은 다른 컬럼들과 함께 생성된 메시지 페이로드에도 나타나요.
kafka_map_virtual_columns_on_write = 1이면 동작이 바뀌어요:
_key(타입String) — Kafka 메시지 키에 매핑_timestamp(타입DateTime) — Kafka 메시지 타임스탬프에 매핑_headers.name(타입Array(String))와_headers.value(타입Array(String)) — Kafka 메시지 헤더에 매핑. 각 쌍(_headers.name[i], _headers.value[i])은 Kafka 헤더 하나가 돼요._headers.name과_headers.value가_headersNested 접두사를 공유하므로 ClickHouse는 모든 행에 대해 두 배열의 크기가 같아야 해요
이 이름을 가진 컬럼은 타입이 위에 나열된 것과 일치할 때만 메시지 페이로드에서 제외돼요. 그렇지 않으면 페이로드에 남아서, 우연히 이 이름을 관련 없는 데이터에 재사용하는 스키마도 계속 동작해요.
예시:
CREATE TABLE kafka_out
(
event_json String,
`_key` String,
`_timestamp` DateTime,
`_headers.name` Array(String),
`_headers.value` Array(String)
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'broker:9092',
kafka_topic_list = 'events',
kafka_group_name = 'events-producer',
kafka_format = 'JSONEachRow',
kafka_map_virtual_columns_on_write = 1;
INSERT INTO kafka_out VALUES
('{"a":1}', 'session-42', now(), ['source', 'trace_id'], ['api', 'abc-123']);
생성된 Kafka 메시지는 페이로드 {"event_json":"{\"a\":1}"}, 키 session-42, 현재 타임스탬프, 두 헤더 source=api와 trace_id=abc-123을 가져요.
데이터 형식 지원
Kafka 엔진은 ClickHouse에서 지원하는 모든 형식을 지원해요. 하나의 Kafka 메시지의 행 수는 형식이 행 기반인지 블록 기반인지에 따라 달라져요.
- 행 기반 형식의 경우 하나의 Kafka 메시지의 행 수를
kafka_max_rows_per_message설정으로 제어할 수 있어요 - 블록 기반 형식의 경우 블록을 더 작은 부분으로 나눌 수 없지만, 하나의 블록의 행 수는 일반 설정 max_block_size로 제어할 수 있어요
ClickHouse Keeper에 커밋된 오프셋을 저장하는 엔진
allow_kafka_offsets_storage_in_keeper가 활성화되면 Kafka 테이블 엔진에 두 가지 설정을 더 지정할 수 있어요:
kafka_keeper_path는 ClickHouse Keeper에서 테이블의 경로를 지정해요kafka_replica_name은 ClickHouse Keeper에서 복제본 이름을 지정해요
두 설정을 모두 지정하거나 둘 다 지정하지 않아야 해요. 둘 다 지정하면 새롭고 실험적인 Kafka 엔진이 사용돼요. 새 엔진은 Kafka에 커밋된 오프셋을 저장하는 데 의존하지 않고 ClickHouse Keeper에 저장해요. 여전히 Kafka에 오프셋을 커밋하려 하지만 테이블이 생성될 때만 해당 오프셋에 의존해요. 다른 모든 상황(테이블 재시작, 오류 후 복구)에서는 ClickHouse Keeper에 저장된 오프셋을 메시지 소비를 계속하기 위한 오프셋으로 사용해요. 커밋된 오프셋 외에도 마지막 배치에서 소비된 메시지 수를 저장해서, 삽입이 실패하면 같은 수의 메시지를 소비해 필요하면 중복 제거를 활성화할 수 있어요.
정적 파티션-샤드 친화성
StorageKafka2를 사용할 때 kafka_partition_shard_num과 kafka_shard_count를 지정해 정적 파티션-샤드 친화성을 선택적으로 활성화할 수 있어요. 이렇게 하면 여러 ClickHouse 인스턴스(샤드)가 같은 Kafka 토픽에서 소비하고, 각 샤드가 공식에 기반해 결정적인 파티션 하위 집합만 처리할 수 있어요:
partition_id % kafka_shard_count == kafka_partition_shard_num - 1
두 설정을 함께 지정해야 해요. 하나만 지정하면 예외가 발생해요. kafka_partition_shard_num 값은 1과 kafka_shard_count 사이여야 해요. 매크로 확장(예: '{shard}')을 지원하며 서버 시작 시마다 다시 확장돼요. 이렇게 하면 Replicated 데이터베이스의 샤드 간에 같은 테이블 메타데이터를 공유할 수 있고, 각 샤드가 자신의 값을 해석해요. 검증은 매크로 확장 후 수행돼요.
모든 샤드는 같은 kafka_keeper_path 값을 사용해야 해요. 모든 복제본은 커밋된 오프셋과 의도(intent) 크기를 공유하지만, 같은 샤드 번호를 가진 복제본만 파티션 잠금을 두고 경쟁해요 (모든 복제본이 같은 샤드 수를 갖는다고 가정).
12-파티션 토픽에서 소비하는 3개 샤드의 예시:
-- Shard 1: consumes partitions 0, 3, 6, 9
CREATE TABLE kafka_shard1 (key UInt64, value String)
ENGINE = Kafka('localhost:9092', 'my-topic', 'my-group', 'JSONEachRow')
SETTINGS
kafka_keeper_path = '/clickhouse/kafka/{database}',
kafka_replica_name = '{replica}',
kafka_partition_shard_num = '1',
kafka_shard_count = 3
SETTINGS allow_kafka_offsets_storage_in_keeper = 1;
-- Shard 2: consumes partitions 1, 4, 7, 10
CREATE TABLE kafka_shard2 (key UInt64, value String)
ENGINE = Kafka('localhost:9092', 'my-topic', 'my-group', 'JSONEachRow')
SETTINGS
kafka_keeper_path = '/clickhouse/kafka/{database}',
kafka_replica_name = '{replica}',
kafka_partition_shard_num = '2',
kafka_shard_count = 3
SETTINGS allow_kafka_offsets_storage_in_keeper = 1;
복제본과 결합하면(같은 kafka_keeper_path를 공유하는 여러 kafka_replica_name 값) 친화성 필터를 먼저 적용해 자격 있는 파티션을 결정한 다음 ZooKeeper 잠금이 그 자격 있는 파티션을 복제본 간에 분배해요.
예시:
CREATE TABLE experimental_kafka (key UInt64, value UInt64)
ENGINE = Kafka('localhost:19092', 'my-topic', 'my-consumer', 'JSONEachRow')
SETTINGS
kafka_keeper_path = '/clickhouse/{database}/{uuid}',
kafka_replica_name = '{replica}'
SETTINGS allow_kafka_offsets_storage_in_keeper=1;
알려진 제한
새 엔진은 실험적이므로 아직 프로덕션 준비가 되지 않았어요. 구현의 몇 가지 알려진 제한이 있어요.
- 테이블을 빠르게 삭제·재생성하거나 다른 엔진에 같은 ClickHouse Keeper 경로를 지정하면 문제가 생길 수 있어요. 모범 사례로
kafka_keeper_path에{uuid}를 사용해 경로 충돌을 피할 수 있어요 - 반복 가능한 읽기를 만들기 위해 단일 스레드에서 여러 파티션의 메시지를 소비할 수 없어요. 반면 Kafka 소비자는 살아있게 유지하기 위해 정기적으로 poll해야 해요. 이 두 목표의 결과로, 우리는
kafka_thread_per_consumer가 활성화된 경우에만 여러 소비자 생성을 허용하기로 결정했어요. 그렇지 않으면 소비자를 정기적으로 poll하는 것과 관련된 문제를 피하기가 너무 복잡해요 - 파티션 친화성을 사용할 때 모든 샤드는 같은
kafka_shard_count를 사용해야 해요. 그렇지 않으면 일부 파티션을 여러 샤드가 소비하거나 아무도 소비하지 않을 수 있어요
데이터 내구성(Durability)
Kafka 엔진은 삽입된 데이터가 디스크에 쓰이기 전에 OS 페이지 캐시가 버려지면 이미 소비된 행을 조용히 잃을 수 있어요. 배치가 의존하는 materialized view에 푸시된 후 소비된 오프셋이 커밋돼(브로커 또는 kafka_keeper_path가 설정되면 ClickHouse Keeper에) 소비자가 그 메시지들을 넘어 재개할 수 있어요. 그러나 삽입된 행은 대상 파트가 fsync된 후에만 내구성이 생기는데, 기본적으로 동기적으로 발생하지 않아요 (fsync_after_insert = 0). 오프셋이 커밋된 후 대상 파트가 fsync되기 전에 페이지 캐시가 손실되면 소비자가 재시작 시 그 메시지들을 넘어 재개하므로 행이 오류 없이 손실되고 count()가 단순히 더 작아져요. 일반적인 프로세스 종료는 커널이 페이지 캐시를 유지하고 결국 기록하므로 이를 드러내지 않아요. 페이지 캐시 손실은 드러내요. 예로 장치 수준 전원 손실과 비정상적인 호스트나 커널 리셋이 있어요.
권장하는 materialized-view 소비 경로(오프셋은 전체 삽입 파이프라인이 끝난 뒤에만 커밋)에서는 대상 MergeTree 테이블에 fsync_after_insert = 1(그리고 fsync_part_directory = 1)을 설정하면 삽입된 파트가 오프셋 커밋 전에 내구성이 생겨 이 창이 크게 좁아져요. 이 설정은 배치가 삽입되는 모든 MergeTree 테이블, 캐스케이드된 materialized-view 대상까지 포함해 활성화해야 해요. 기본값으로 남겨둔 테이블은 여전히 파트를 잃을 수 있어요. 비동기 중개자는 이 설정만으로 내구성을 얻지 못해요. 예를 들어 distributed_foreground_insert = 0이면 Distributed 대상이 백그라운드에서 삽입하는데, 이는 ClickHouse Cloud 외부의 기본값이므로 자체 내구성 설정이나 동기 삽입이 필요해요. 이 완화는 오프셋이 삽입 완료 전에 커밋되는 경우에도 적용되지 않아요. 그런 경우는 직접 INSERT ... SELECT ... FROM <kafka_table>과 kafka_commit_on_select = 1, 그리고 브로커 기반 엔진의 kafka_commit_every_batch = 1입니다. 후자 설정은 kafka_keeper_path가 설정되면 무시되므로 그곳에서 중간 커밋을 일으킬 수 없어요.
함께 보기