Kafka에서 ClickHouse로
ClickHouse Cloud를 사용하는 경우에는 대신 ClickPipes 사용을 권장합니다. ClickPipes는 프라이빗 네트워크 연결을 네이티브로 지원하고, 수집 및 클러스터 리소스를 서로 독립적으로 스케일링할 수 있으며, Kafka 데이터를 ClickHouse로 스트리밍할 때 포괄적인 모니터링 기능을 제공합니다.
개요
절차
1. 준비
2. ClickHouse 구성
conf.d/ 디렉터리 아래의 새 파일에 넣거나 기존 설정 파일에 병합할 수 있습니다. 구성 가능한 설정은 여기를 참조하십시오.
이 튜토리얼에서 사용할 KafkaEngine이라는 데이터베이스도 생성하겠습니다:
3. 대상 테이블 생성
4. 토픽 생성 및 데이터 입력
github 토픽을 생성할 수 있습니다:
5. Kafka 테이블 엔진 생성
JSONEachRow를 사용한다는 점에 유의하십시오. github와 clickhouse 값은 각각 토픽 이름과 컨슈머 그룹 이름을 나타냅니다. 실제로 토픽은 값 목록으로 지정할 수도 있습니다.
github_queue에 간단한 SELECT를 수행하면 몇 개의 행을 읽을 수 있어야 합니다. 이 작업은 consumer 오프셋을 앞으로 이동시키므로 reset하지 않으면 이 행들을 다시 읽을 수 없다는 점에 유의하십시오. LIMIT와 필수 매개변수 stream_like_engine_allow_direct_select에도 유의하십시오.
6. materialized view 생성
7. 행이 삽입되었는지 확인합니다
일반적인 작업
메시지 소비 중지 및 재시작
Kafka 메타데이터 추가
_ 접두사로 시작합니다.
가상 컬럼의 전체 목록은 여기에서 확인할 수 있습니다.
테이블에 가상 컬럼을 반영하도록 업데이트하려면 materialized view를 삭제하고, Kafka 엔진 테이블을 다시 ATTACH한 다음, materialized view를 다시 생성해야 합니다.
Kafka 엔진 설정 수정
문제 디버깅
잘못된 형식의 메시지 처리
- 메시지 필드를 문자열로 취급하십시오. 필요한 경우 materialized view 구문에서 함수를 사용해 정제와 CAST를 수행할 수 있습니다. 이는 운영 환경용 해결책으로 보기는 어렵지만, 일회성 수집에는 도움이 될 수 있습니다.
- 토픽에서 JSON을 읽고 JSONEachRow 포맷을 사용하는 경우
input_format_skip_unknown_fields설정을 사용하십시오. 데이터를 쓸 때 ClickHouse는 기본적으로 입력 데이터에 대상 테이블(target table)에 없는 컬럼이 포함되어 있으면 예외를 발생시킵니다. 하지만 이 옵션을 활성화하면 이러한 초과 컬럼은 무시됩니다. 다시 말해, 이 역시 운영 환경 수준의 해결책은 아니며 다른 사용자를 혼란스럽게 할 수 있습니다. kafka_skip_broken_messages설정도 고려하십시오. 이 설정을 사용하려면 잘못된 형식의 메시지에 대해 block당 허용 수준을 지정해야 하며, 이는kafka_max_block_size를 기준으로 판단됩니다. 이 허용치를 초과하면(절대 메시지 수 기준) 일반적인 예외 동작으로 돌아가며, 다른 메시지는 건너뛰게 됩니다.
전달 시맨틱과 중복 문제
쿼럼 기반 삽입
ClickHouse에서 Kafka로
단계
1. 행을 직접 삽입
2. materialized views 사용
github_out 또는 이에 상응하는 토픽을 생성하십시오. Kafka 테이블 엔진 github_out_queue가 이 토픽을 가리키도록 하십시오.
github_out_mv를 생성해 GitHub 테이블을 가리키도록 하고, 트리거될 때 위의 engine에 행이 삽입되도록 합니다. 그러면 GitHub 테이블에 추가된 내용이 새 Kafka 토픽으로 전송됩니다.
github_out 토픽을 읽어 보면 메시지가 정상적으로 전달되었는지 확인할 수 있습니다.
클러스터와 성능
ClickHouse 클러스터 사용하기
성능 튜닝
- 성능은 메시지 크기, 포맷, 대상 테이블 유형에 따라 달라집니다. 단일 테이블 엔진에서 초당 10만 행은 달성 가능한 수준으로 볼 수 있습니다. 기본적으로 메시지는 블록 단위로 읽으며, 이는
kafka_max_block_size매개변수로 제어됩니다. 이 값의 기본값은 max_insert_block_size이며, 기본 설정은 1,048,576입니다. 메시지가 극히 크지 않다면 이 값은 거의 항상 늘리는 편이 좋습니다. 500k~1M 범위의 값도 흔히 사용됩니다. 처리량에 미치는 영향을 테스트하고 평가하십시오. - 테이블 엔진의 컨슈머 수는
kafka_num_consumers로 늘릴 수 있습니다. 그러나 기본적으로kafka_thread_per_consumer를 기본값 1에서 변경하지 않으면 삽입은 단일 스레드에서 직렬화됩니다. 플러시가 병렬로 수행되도록 하려면 이 값을 1로 설정하십시오. 또한 컨슈머가 N개인 Kafka 엔진 테이블(kafka_thread_per_consumer=1)을 생성하는 것은, 각각 materialized view가 있고kafka_thread_per_consumer=0인 Kafka 엔진 N개를 생성하는 것과 논리적으로 동일합니다. - 컨슈머 수를 늘리는 데에는 비용이 따릅니다. 각 컨슈머는 자체 버퍼와 스레드를 유지하므로 서버 오버헤드가 증가합니다. 따라서 가능하면 먼저 클러스터 전반으로 선형 확장하고, 컨슈머로 인한 오버헤드를 함께 고려하십시오.
- Kafka 메시지 처리량의 변동이 크고 지연이 허용된다면, 더 큰 블록이 플러시되도록
stream_flush_interval_ms를 늘리는 것을 고려하십시오. - background_message_broker_schedule_pool_size는 백그라운드 작업을 수행하는 스레드 수를 설정합니다. 이 스레드들은 Kafka streaming에 사용됩니다. 이 설정은 ClickHouse 서버 시작 시 적용되며 사용자 세션에서는 변경할 수 없고, 기본값은 16입니다. 로그에서 timeout이 보인다면 이 값을 늘리는 것이 적절할 수 있습니다.
- Kafka와 통신할 때는
librdkafka라이브러리를 사용하며, 이 라이브러리도 자체적으로 스레드를 생성합니다. 따라서 Kafka 테이블이나 컨슈머 수가 많아지면 context switch도 크게 늘어날 수 있습니다. 가능하다면 이 부하를 클러스터 전체에 분산하고 대상 테이블만 복제하거나, 여러 토픽을 읽는 하나의 테이블 엔진을 사용하는 것도 고려하십시오. 값 목록이 지원됩니다. 하나의 테이블에서 여러 materialized view가 읽을 수 있으며, 각각 특정 토픽의 데이터만 필터링할 수 있습니다.
추가 설정
- Kafka_max_wait_ms - 재시도하기 전에 Kafka에서 메시지를 읽기 위해 대기하는 시간을 밀리초 단위로 지정합니다. 사용자 프로필 수준에서 설정하며 기본값은 5000입니다.