도움이 필요하면 리포지토리에 이슈를 등록하거나 ClickHouse 공개 Slack에서 질문하십시오.
라이선스
환경 요구 사항
버전 호환성 매트릭스
주요 기능
- 별도 설정 없이 정확히 한 번 처리 의미 체계를 제공합니다. 이는 KeeperMap이라는 새로운 ClickHouse 핵심 기능(커넥터의 상태 저장소로 사용됨)을 기반으로 하며, 간결한 아키텍처를 구현할 수 있게 합니다.
- 타사 상태 저장소를 지원합니다. 현재는 기본적으로 In-memory를 사용하지만 KeeperMap도 사용할 수 있습니다(Redis 지원도 곧 추가될 예정입니다).
- 핵심 통합: ClickHouse가 직접 개발, 유지 관리 및 지원합니다.
- ClickHouse Cloud를 대상으로 지속적으로 테스트합니다.
- 명시된 스키마 기반 및 스키마리스 데이터 삽입을 지원합니다.
- ClickHouse의 모든 데이터 타입을 지원합니다.
설치 안내
연결 정보 확인
ClickHouse Cloud 서비스의 연결 정보는 ClickHouse Cloud 콘솔에서 확인할 수 있습니다.
서비스를 선택한 다음 Connect를 클릭하십시오.
HTTPS를 선택하십시오. 연결 정보가 예시
curl 명령으로 표시됩니다.
자가 관리형 ClickHouse를 사용하는 경우 연결 정보는 ClickHouse 관리자가 설정합니다.
일반 설치 지침
- ClickHouse Kafka Connect Sink 리포지토리의 Releases 페이지에서 커넥터 JAR 파일이 포함된 ZIP 아카이브를 다운로드합니다.
- ZIP 파일의 내용을 추출한 후 원하는 위치에 복사합니다.
- Confluent Platform이 플러그인을 찾을 수 있도록 Connect 속성 파일의 plugin.path 구성에 플러그인 디렉터리가 있는 경로를 추가합니다.
- 구성에 토픽 이름, ClickHouse 인스턴스 호스트명, 비밀번호를 지정합니다.
- Confluent Platform을 다시 시작합니다.
- Confluent Platform을 사용하는 경우 Confluent Control Center UI에 로그인하여 사용 가능한 커넥터 목록에 ClickHouse Sink가 표시되는지 확인합니다.
구성 옵션
- 연결 정보: 호스트명(필수) 및 포트(선택 사항)
- 사용자 자격 증명: 비밀번호(필수) 및 사용자 이름(선택 사항)
- 커넥터 클래스:
com.clickhouse.kafka.connect.ClickHouseSinkConnector(필수) - topics 또는 topics.regex: 폴링할 Kafka 토픽 - 토픽 이름은 테이블 이름과 일치해야 합니다(필수)
- 키 및 값 컨버터: 토픽의 데이터 유형에 따라 설정합니다. worker 구성에 이미 정의되어 있지 않은 경우 필수입니다.
대상 테이블
사전 처리
지원되는 데이터 타입
-
(1) - JSON은 ClickHouse 설정에서
input_format_binary_read_json_as_string=1일 때만 지원됩니다. 이는 RowBinary 포맷 계열에서만 동작하며, 이 설정은 삽입 요청의 모든 컬럼에 영향을 미치므로 해당 컬럼은 모두 문자열이어야 합니다. 이 경우 커넥터는 STRUCT를 JSON 문자열로 변환합니다. -
(2) - struct에
oneof와 같은 union이 있는 경우, 컨버터는 필드 이름에 prefix/suffix를 추가하지 않도록 구성해야 합니다.generate.index.for.unions=false는ProtobufConverter설정입니다.
구성 예시
기본 구성
localhost:8443에서 ClickHouse 서버가 실행 중이며, 데이터는 스키마가 없는 JSON 형식이라고 가정합니다.
위의 커넥터 구성에서는 워커 구성에서
connector.client.config.override.policy=All을 설정해 클라이언트 재정의를 활성화해야 합니다. 자세한 내용은 Kafka Connect 문서를 참조하십시오.여러 토픽을 사용하는 기본 구성
DLQ를 포함한 기본 구성
여러 데이터 포맷과 함께 사용하기
Avro 스키마 지원
Avro 타입 매핑
io.confluent.connect.avro.AvroConverter에서 정의합니다. 변환 로직에 대한 자세한 내용은 Kafka Connect 문서를 참조하십시오.
✅: 지원
❌: 미지원
️⚠️: 부분 지원
Kafka Connect 타입과 ClickHouse 타입 간의 매핑은 지원되는 데이터 타입을 참조하십시오.
지원되지 않는 Avro 스키마
fixeddecimal논리 유형
- 널 허용 유니온
- 레코드 유니온
Protobuf 스키마 지원
Protobuf 타입 매핑
io.confluent.connect.protobuf.ProtobufConverter에서 정의합니다. 변환 로직에 관한 고급 정보는 Kafka Connect docs를 참조하십시오.
✅: 지원
❌: 지원되지 않음
️⚠️: 부분 지원
Kafka Connect 타입과 ClickHouse 타입 간 매핑은 지원되는 데이터 타입을 참조하십시오.
oneof 필드를 ClickHouse 컬럼으로 변환할 때 참고 사항
oneof)을 ClickHouse Variant 타입으로 변환하는 기능을 지원하지 않습니다. 대신 oneof 필드를 ClickHouse 테이블 스키마에 각각의 널 허용 필드로 나열하십시오.
예시:
지원되지 않는 Protobuf 스키마
- 다중 메시지 유니온 (CH 버전 26.1 이전)
allow_experimental_nullable_tuple_type=1로 설정하면 이 스키마가 지원됩니다(이 문서 페이지 참조).
JSON 스키마 지원
String 컨버터 지원
내부 버퍼링
poll() 호출에서 레코드를 누적한 뒤, 이를 더 큰 배치로 묶어 ClickHouse에 플러시할 수 있습니다. 이렇게 하면 각 poll()이 파티션별로 작은 배치를 많이 생성하는 워크로드에서 처리량을 높일 수 있습니다.
주요 동작:
bufferCount는 플러시하기 전에 버퍼링할 레코드 수를 제어합니다.bufferFlushTime은 버퍼링된 레코드를 플러시하기 전까지의 최대 대기 시간(밀리초)을 설정합니다.bufferFlushTime은bufferCount > 0일 때만 유효합니다.bufferCount=0및bufferFlushTime=0이면 버퍼링이 비활성화된 상태로 유지됩니다(기본 동작).exactlyOnce=true일 때는 버퍼링이 지원되지 않습니다.
exactlyOnce=false로 exactly-once 모드를 비활성화하거나, bufferCount=0으로 버퍼링을 비활성화하십시오.
예시:
로깅
모니터링
ClickHouse 전용 메트릭
Kafka 프로듀서/컨슈머 메트릭
records-sent-total: 토픽으로 전송된 총 레코드 수bytes-sent-total: 토픽으로 전송된 총 바이트 수record-send-rate: 초당 전송된 레코드의 평균 수byte-rate: 초당 전송된 평균 바이트 수compression-rate: 달성된 압축률
records-sent-total: 파티션으로 전송된 총 레코드 수bytes-sent-total: 파티션으로 전송된 총 바이트 수records-lag: 파티션의 현재 지연records-lead: 파티션의 현재 선행 정도replica-fetch-lag: 레플리카의 지연 정보
connection-creation-total: Kafka 노드에 생성된 총 연결 수connection-close-total: 종료된 총 연결 수request-total: 노드로 전송된 총 요청 수response-total: 노드에서 수신한 총 응답 수request-rate: 초당 평균 요청 수response-rate: 초당 평균 응답 수
- 처리량: 데이터 수집 속도 추적
- 지연: 병목과 처리 지연 식별
- 압축: 데이터 압축 효율 측정
- 연결 상태: 네트워크 연결 및 안정성 모니터링
Kafka Connect Framework 메트릭
task-count: 커넥터의 전체 작업 수running-task-count: 현재 실행 중인 작업 수paused-task-count: 현재 일시 중지된 작업 수failed-task-count: 실패한 작업 수destroyed-task-count: 제거된 작업 수unassigned-task-count: 할당되지 않은 작업 수
running, paused, failed, destroyed, unassigned
오류 메트릭:
deadletterqueue-produce-failures: 실패한 DLQ 쓰기 수deadletterqueue-produce-requests: 전체 DLQ 쓰기 시도 수last-error-timestamp: 마지막 오류의 타임스탬프records-skip-total: 오류로 인해 건너뛴 전체 레코드 수records-retry-total: 재시도한 전체 레코드 수errors-total: 발생한 전체 오류 수
offset-commit-failures: 실패한 offset commit 수offset-commit-avg-time-ms: offset commit의 평균 소요 시간offset-commit-max-time-ms: offset commit의 최대 소요 시간put-batch-avg-time-ms: Batch 처리 평균 시간put-batch-max-time-ms: Batch 처리 최대 시간source-record-poll-total: poll한 전체 레코드 수
모니터링 모범 사례
- Consumer lag 모니터링: 처리 병목을 식별할 수 있도록 파티션별
records-lag를 추적합니다 - 오류율 추적: 데이터 품질 문제를 감지할 수 있도록
errors-total및records-skip-total을 확인합니다 - 작업 상태 관찰: 작업이 정상적으로 실행되는지 확인할 수 있도록 작업 상태 메트릭을 모니터링합니다
- 처리량 측정: 수집 성능을 추적하기 위해
records-send-rate및byte-rate를 사용합니다 - 연결 상태 모니터링: 네트워크 문제를 확인하기 위해 노드 수준의 연결 메트릭을 점검합니다
- 압축 효율 추적: 데이터 전송을 최적화하기 위해
compression-rate를 사용합니다
제한 사항
- 삭제는 지원되지 않습니다.
- 배치 크기는 Kafka Consumer 속성을 따릅니다.
- exactly-once에 KeeperMap을 사용하는 경우 오프셋이 변경되거나 되돌려지면 해당 토픽의 KeeperMap 내용을 삭제해야 합니다. (자세한 내용은 아래 문제 해결 가이드를 참조하십시오)
성능 튜닝 및 처리량 최적화
성능 튜닝은 언제 필요합니까?
- 고처리량 워크로드: Kafka 토픽에서 초당 수백만 건의 이벤트를 처리하는 경우
- Consumer lag: 커넥터가 데이터 생성 속도를 따라가지 못해 lag가 계속 증가하는 경우
- 리소스 제약: CPU, 메모리 또는 네트워크 사용량을 최적화해야 하는 경우
- 여러 토픽: 대용량 토픽 여러 개를 동시에 소비하는 경우
- 작은 메시지 크기: 서버 측 배칭의 이점을 얻을 수 있도록 작은 메시지를 대량으로 처리하는 경우
- 낮거나 중간 수준의 처리량(< 10,000 messages/second)을 처리하는 경우
- 사용 사례에서 Consumer lag가 안정적이고 허용 가능한 수준인 경우
- 기본 커넥터 설정만으로도 필요한 처리량 요구 사항을 이미 충족하는 경우
- ClickHouse 클러스터가 유입되는 부하를 무리 없이 처리할 수 있는 경우
데이터 흐름 이해하기
- Kafka Connect Framework가 백그라운드에서 Kafka 토픽의 메시지를 가져옵니다
- 커넥터는 폴링을 통해 프레임워크의 내부 버퍼에서 메시지를 가져옵니다
- 커넥터는 메시지를 배치로 묶어 폴링 크기에 따라 처리합니다
- ClickHouse는 HTTP/S를 통해 배치된 삽입을 수신합니다
- ClickHouse는 삽입을 처리합니다(동기식 또는 비동기식)
Kafka Connect 배치 크기 튜닝
fetch 설정
fetch.min.bytes: 프레임워크가 값을 커넥터에 전달하기 전에 필요한 최소 데이터 양(기본값: 1 byte)fetch.max.bytes: 단일 요청으로 fetch할 수 있는 최대 데이터 양(기본값: 52428800 / 50 MB)fetch.max.wait.ms:fetch.min.bytes조건이 충족되지 않으면 데이터를 반환하기 전까지 대기하는 최대 시간(기본값: 500 ms)
Confluent Cloud에서는 이러한 설정을 조정하려면 Confluent Cloud 지원 케이스를 열어야 합니다.
폴링 설정
max.poll.records: 한 번의 폴링으로 반환되는 최대 레코드 수(기본값: 500)max.partition.fetch.bytes: 파티션별 최대 데이터 크기(기본값: 1048576 / 1 MB)
Confluent Cloud에서는 이러한 설정을 조정하려면 Confluent Cloud에서 지원 케이스를 접수해야 합니다.
고처리량을 위한 권장 설정
위 속성을 사용하려면 워커 구성에서
connector.client.config.override.policy=All을 통해 클라이언트 재정의를 활성화해야 합니다. 자세한 내용은 Kafka Connect 문서를 참조하십시오.- 더 큰 배치 = 더 나은 ClickHouse 수집 성능, 더 적은 파트, 더 낮은 오버헤드
- 더 큰 배치 = 더 높은 메모리 사용량, 종단 간 지연 시간 증가 가능성
- 너무 큰 배치 = timeout, OutOfMemory 오류 또는
max.poll.interval.ms초과 위험
비동기 삽입
async 삽입을 사용해야 하는 경우
- 작은 배치가 많은 경우: 커넥터가 작은 배치(< 배치당 1000행)를 자주 전송하는 경우
- 높은 동시성: 여러 커넥터 작업이 동일한 테이블에 쓰기 작업을 수행하는 경우
- 분산 배포: 서로 다른 호스트에서 많은 커넥터 인스턴스를 실행하는 경우
- 파트 생성 오버헤드: “too many parts” 오류가 발생하는 경우
- 혼합 워크로드: 실시간 수집과 쿼리 워크로드를 함께 실행하는 경우
- 이미 큰 배치(배치당 10,000행 초과)를 제어된 빈도로 전송하는 경우
- 데이터가 즉시 표시되어야 하는 경우(쿼리에서 데이터를 즉시 확인해야 함)
wait_for_async_insert=0을 사용하는 정확히 한 번 처리 의미 체계가 요구 사항과 충돌하는 경우- 대신 클라이언트 측 배칭 개선의 이점을 얻을 수 있는 사용 사례인 경우
비동기 삽입의 작동 방식
- 커넥터로부터 삽입 쿼리를 받습니다
- 데이터를 메모리 버퍼에 기록합니다(즉시 디스크에 기록하지 않음)
- 커넥터에 성공을 반환합니다 (
wait_for_async_insert=0인 경우) - 다음 조건 중 하나가 충족되면 버퍼를 디스크로 플러시합니다:
- 버퍼 크기가
async_insert_max_data_size에 도달함(기본값: 100 MB) - 첫 번째 삽입 후
async_insert_busy_timeout_ms밀리초가 경과함(기본값: 1000 ms) - 누적된 쿼리 수가 최대치에 도달함(
async_insert_max_query_number, 기본값: 100)
- 버퍼 크기가
async 삽입 활성화
clickhouseSettings 구성 매개변수에 추가합니다:
async_insert=1: 비동기 삽입을 활성화합니다wait_for_async_insert=1(권장): 커넥터가 확인 응답을 보내기 전에 데이터가 ClickHouse 스토리지에 플러시될 때까지 기다립니다. 전송 보장을 제공합니다.wait_for_async_insert=0: 버퍼링 직후 커넥터가 즉시 확인 응답을 보냅니다. 성능은 더 좋지만 플러시되기 전에 서버가 충돌하면 데이터가 손실될 수 있습니다.
async insert 동작 조정
async_insert_max_data_size(기본값: 104857600 / 100 MB): 플러시 전 최대 버퍼 크기async_insert_busy_timeout_ms(기본값: 1000): 플러시 전 최대 시간(ms)async_insert_stale_timeout_ms(기본값: 0): 마지막 삽입 후 플러시되기까지의 시간(ms)async_insert_max_query_number(기본값: 100): 플러시 전 최대 쿼리 수
- 장점: 파트 수 감소, 향상된 머지 성능, 낮은 CPU 오버헤드, 높은 동시성에서 더 나은 처리량
- 고려 사항: 데이터를 즉시 조회할 수 없고, 엔드 투 엔드 지연 시간이 약간 증가합니다
- 위험 요소:
wait_for_async_insert=0인 경우 서버 장애 시 데이터가 손실될 수 있으며, 버퍼가 크면 메모리 사용량 압박이 발생할 수 있습니다
정확히 한 번 처리 의미 체계의 async 삽입
exactlyOnce=true를 사용하는 경우:
wait_for_async_insert=1을 사용하십시오.
async 삽입에 대한 자세한 내용은 ClickHouse async 삽입 문서를 참조하십시오.
커넥터 병렬성
커넥터별 작업 수
- 효과적인 최대 작업 수 = 토픽 파티션 수
- 각 작업은 ClickHouse와의 자체 연결을 유지합니다
- 작업 수가 많을수록 오버헤드가 커지고 리소스 경합이 발생할 가능성이 높아집니다
tasks.max를 토픽 파티션 수와 같게 설정하여 시작한 다음, CPU 및 처리량 메트릭을 기준으로 조정하세요.
배칭 시 파티션 무시
exactlyOnce=false일 때만 사용하십시오. 이 설정을 사용하면 더 큰 배치를 생성해 처리량을 높일 수 있지만, 파티션별 순서 보장은 유지되지 않습니다.
여러 고처리량 토픽
topic2TableMap을 사용해 토픽을 테이블에 매핑하며, 삽입 병목으로 인해 consumer lag이 발생하는 경우에는 토픽별로 커넥터를 하나씩 생성하는 방안을 고려하십시오.
이 문제가 발생하는 주된 이유는 현재 배치가 각 테이블에 순차적으로 삽입되기 때문입니다.
권장 사항: 여러 고처리량 토픽의 경우, 병렬 삽입 처리량을 극대화하려면 토픽별로 커넥터 인스턴스를 하나씩 배포하십시오.
ClickHouse 테이블 엔진 고려 사항
MergeTree: 대부분의 사용 사례에 가장 적합하며, 쿼리와 삽입 성능의 균형이 좋습니다ReplicatedMergeTree: 고가용성을 위해 필요하며, 복제 오버헤드가 추가됩니다- 적절한
ORDER BY를 사용하는*MergeTree: 쿼리 패턴에 맞게 최적화하십시오
연결 풀링 및 타임아웃
socket_timeout(기본값: 30000 ms): 읽기 작업의 최대 대기 시간connection_timeout(기본값: 10000 ms): 연결을 설정하는 최대 대기 시간
성능 모니터링 및 문제 해결
- Consumer lag: Kafka 모니터링 도구를 사용해 파티션별 지연을 추적합니다
- 커넥터 메트릭: JMX를 통해
receivedRecords,recordProcessingTime,taskProcessingTime를 모니터링합니다(모니터링 참조) - ClickHouse 메트릭:
system.asynchronous_inserts: async insert 버퍼 사용량을 모니터링합니다system.parts: 머지 문제를 감지할 수 있도록 파트 수를 모니터링합니다system.merges: 진행 중인 머지를 모니터링합니다system.events:InsertedRows,InsertedBytes,FailedInsertQuery를 추적합니다
모범 사례 요약
- 기본값으로 시작한 다음 실제 성능을 측정하고 그 결과에 따라 조정하세요
- 더 큰 배치를 우선하세요: 가능하면 한 번의 삽입당 10,000~100,000개 행을 목표로 하세요
- 작은 배치를 많이 보내거나 동시성(Concurrency)이 높을 때는 async 삽입을 사용하세요
- 정확히 한 번 처리 의미 체계를 위해
wait_for_async_insert=1을 항상 사용하세요 - 수평 확장하세요:
tasks.max를 파티션 수만큼 늘리세요 - 처리량을 최대화하려면 대용량 토픽마다 커넥터를 하나씩 사용하세요
- 지속적으로 모니터링하세요: consumer lag, part 개수, 머지 활동을 추적하세요
- 충분히 테스트하세요: 프로덕션 배포 전에 반드시 실제와 유사한 부하에서 구성 변경을 테스트하세요
예시: 고처리량 구성
위 커넥터 구성에서는 worker 구성에서
connector.client.config.override.policy=All을 설정해 클라이언트 재정의를 활성화해야 합니다. 자세한 내용은 Kafka Connect 문서를 참조하십시오.- 폴링 한 번에 최대 10,000개의 레코드를 처리합니다
- 더 큰 삽입을 위해 여러 파티션의 데이터를 배치로 묶습니다
- 16 MB 버퍼를 사용하는 async 삽입을 사용합니다
- 8개의 병렬 작업을 실행합니다(파티션 수에 맞춰 설정)
- 엄격한 순서 보장보다 처리량에 중점을 두고 최적화되었습니다
문제 해결
”토픽 [someTopic] 파티션 [0]의 상태가 일치하지 않음”
이 조정은 exactly-once 동작에 영향을 줄 수 있습니다.
”커넥터는 어떤 오류를 재시도합니까?”
ClickHouseException- ClickHouse에서 발생할 수 있는 일반적인 예외입니다. 일반적으로 서버에 과부하가 걸렸을 때 발생하며, 특히 다음 오류 코드는 일시적인 오류로 간주됩니다:- 3 - UNEXPECTED_END_OF_FILE
- 107 - FILE_DOESNT_EXIST
- 159 - TIMEOUT_EXCEEDED
- 164 - READONLY
- 202 - TOO_MANY_SIMULTANEOUS_QUERIES
- 203 - NO_FREE_CONNECTION
- 209 - SOCKET_TIMEOUT
- 210 - NETWORK_ERROR
- 241 - MEMORY_LIMIT_EXCEEDED
- 242 - TABLE_IS_READ_ONLY
- 252 - TOO_MANY_PARTS
- 285 - TOO_FEW_LIVE_REPLICAS
- 319 - UNKNOWN_STATUS_OF_INSERT
- 425 - SYSTEM_ERROR
- 999 - KEEPER_EXCEPTION
SocketTimeoutException- 소켓 timeout이 발생하면 이 예외가 발생합니다.UnknownHostException- 호스트를 확인할 수 없을 때 이 예외가 발생합니다.IOException- 네트워크에 문제가 있을 때 이 예외가 발생합니다.
”모든 데이터가 비어 있거나 0입니다”
_를 구분자(Delimiter)로 사용). 그러면 테이블(table)의 필드는 “field1_field2_field3” 포맷을 따르게 됩니다(예: “before_id”, “after_id” 등).
”ClickHouse에서 Kafka 키를 사용하고 싶습니다”
KeyToValue 변환을 사용하면 키를 value 필드의 새 _key 필드로 옮길 수 있습니다: