> ## Documentation Index
> Fetch the complete documentation index at: https://clickhouse.com/docs/llms.txt
> Use this file to discover all available pages before exploring further.

> ClickHouse에서 제공하는 공식 Kafka 커넥터입니다.

# ClickHouse Kafka Connect Sink

<Note>
  도움이 필요하면 [리포지토리에 이슈를 등록](https://github.com/ClickHouse/clickhouse-kafka-connect/issues)하거나 [ClickHouse 공개 Slack](https://clickhouse.com/slack)에서 질문하십시오.
</Note>

**ClickHouse Kafka Connect Sink**는 Kafka 토픽에서 ClickHouse 테이블로 데이터를 전달하는 Kafka 커넥터입니다.

<div id="license">
  ### 라이선스
</div>

Kafka Connector 싱크는 [Apache 2.0 라이선스](https://www.apache.org/licenses/LICENSE-2.0)에 따라 배포됩니다.

<div id="requirements-for-the-environment">
  ### 환경 요구 사항
</div>

환경에 [Kafka Connect](https://docs.confluent.io/platform/current/connect/index.html) 프레임워크 v2.7 이상이 설치되어 있어야 합니다.

<div id="version-compatibility-matrix">
  ### 버전 호환성 매트릭스
</div>

| ClickHouse Kafka Connect 버전 | ClickHouse 버전 | Kafka Connect | Confluent Platform |
| --------------------------- | ------------- | ------------- | ------------------ |
| 1.0.0                       | > 23.3        | > 2.7         | > 6.1              |

<div id="main-features">
  ### 주요 기능
</div>

* 별도 설정 없이 정확히 한 번 처리 의미 체계를 제공합니다. 이는 [KeeperMap](https://github.com/ClickHouse/ClickHouse/pull/39976)이라는 새로운 ClickHouse 핵심 기능(커넥터의 상태 저장소로 사용됨)을 기반으로 하며, 간결한 아키텍처를 구현할 수 있게 합니다.
* 타사 상태 저장소를 지원합니다. 현재는 기본적으로 In-memory를 사용하지만 KeeperMap도 사용할 수 있습니다(Redis 지원도 곧 추가될 예정입니다).
* 핵심 통합: ClickHouse가 직접 개발, 유지 관리 및 지원합니다.
* [ClickHouse Cloud](https://clickhouse.com/cloud)를 대상으로 지속적으로 테스트합니다.
* 명시된 스키마 기반 및 스키마리스 데이터 삽입을 지원합니다.
* ClickHouse의 모든 데이터 타입을 지원합니다.

<div id="installation-instructions">
  ### 설치 안내
</div>

<div id="gather-your-connection-details">
  #### 연결 정보 확인
</div>

HTTP(S)로 ClickHouse에 연결하려면 다음 정보가 필요합니다.

| 매개변수                      | 설명                                                         |
| ------------------------- | ---------------------------------------------------------- |
| `HOST` and `PORT`         | 일반적으로 TLS를 사용하는 경우 포트는 8443, TLS를 사용하지 않는 경우 8123입니다.      |
| `DATABASE NAME`           | 기본적으로 `default`라는 이름의 데이터베이스가 제공되며, 연결할 데이터베이스 이름을 사용하십시오. |
| `USERNAME` and `PASSWORD` | 기본 사용자 이름은 `default`입니다. 사용 사례에 맞는 사용자 이름을 사용하십시오.         |

ClickHouse Cloud 서비스의 연결 정보는 ClickHouse Cloud 콘솔에서 확인할 수 있습니다.
서비스를 선택한 다음 **Connect**를 클릭하십시오.

<div className="ch-image-md">
  <Frame>
    <img src="https://mintcdn.com/private-7c7dfe99/CFFsa2agBPbviR4r/images/_snippets/cloud-connect-button.webp?fit=max&auto=format&n=CFFsa2agBPbviR4r&q=85&s=ec0a298a33ca841e947fa5e8bae47362" alt="ClickHouse Cloud 서비스 연결 버튼" width="998" height="932" data-path="images/_snippets/cloud-connect-button.webp" />
  </Frame>
</div>

**HTTPS**를 선택하십시오. 연결 정보가 예시 `curl` 명령으로 표시됩니다.

<div className="ch-image-md">
  <Frame>
    <img src="https://mintcdn.com/private-7c7dfe99/CFFsa2agBPbviR4r/images/_snippets/connection-details-https.webp?fit=max&auto=format&n=CFFsa2agBPbviR4r&q=85&s=cb0fbd98aa2b5b7ca484c9f53395ee07" alt="ClickHouse Cloud HTTPS 연결 정보" width="1320" height="1184" data-path="images/_snippets/connection-details-https.webp" />
  </Frame>
</div>

자가 관리형 ClickHouse를 사용하는 경우 연결 정보는 ClickHouse 관리자가 설정합니다.

<div id="general-installation-instructions">
  #### 일반 설치 지침
</div>

이 커넥터는 플러그인을 실행하는 데 필요한 모든 클래스 파일이 포함된 단일 JAR 파일 형태로 배포됩니다.

플러그인을 설치하려면 다음 단계를 따르세요:

* ClickHouse Kafka Connect Sink 리포지토리의 [Releases](https://github.com/ClickHouse/clickhouse-kafka-connect/releases) 페이지에서 커넥터 JAR 파일이 포함된 ZIP 아카이브를 다운로드합니다.
* ZIP 파일의 내용을 추출한 후 원하는 위치에 복사합니다.
* Confluent Platform이 플러그인을 찾을 수 있도록 Connect 속성 파일의 [plugin.path](https://kafka.apache.org/documentation/#connectconfigs_plugin.path) 구성에 플러그인 디렉터리가 있는 경로를 추가합니다.
* 구성에 토픽 이름, ClickHouse 인스턴스 호스트명, 비밀번호를 지정합니다.

```yml theme={null}
connector.class=com.clickhouse.kafka.connect.ClickHouseSinkConnector
tasks.max=1
topics=<topic_name>
ssl=true
jdbcConnectionProperties=?sslmode=STRICT
security.protocol=SSL
hostname=<hostname>
database=<database_name>
password=<password>
ssl.truststore.location=/tmp/kafka.client.truststore.jks
port=8443
value.converter.schemas.enable=false
value.converter=org.apache.kafka.connect.json.JsonConverter
exactlyOnce=true
username=default
schemas.enable=false
```

* Confluent Platform을 다시 시작합니다.
* Confluent Platform을 사용하는 경우 Confluent Control Center UI에 로그인하여 사용 가능한 커넥터 목록에 ClickHouse Sink가 표시되는지 확인합니다.

<div id="configuration-options">
  ### 구성 옵션
</div>

ClickHouse Sink를 ClickHouse 서버에 연결하려면 다음 정보를 제공해야 합니다:

* 연결 정보: 호스트명(**필수**) 및 포트(선택 사항)
* 사용자 자격 증명: 비밀번호(**필수**) 및 사용자 이름(선택 사항)
* 커넥터 클래스: `com.clickhouse.kafka.connect.ClickHouseSinkConnector` (**필수**)
* topics 또는 topics.regex: 폴링할 Kafka 토픽 - 토픽 이름은 테이블 이름과 일치해야 합니다(**필수**)
* 키 및 값 컨버터: 토픽의 데이터 유형에 따라 설정합니다. worker 구성에 이미 정의되어 있지 않은 경우 필수입니다.

구성 옵션의 전체 표는 다음과 같습니다:

| 속성 이름                                           | 설명                                                                                                                                                                                       | 기본값                                                      |
| ----------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------- |
| `hostname` (필수)                                 | 서버의 호스트명 또는 IP 주소                                                                                                                                                                        | 해당 없음                                                    |
| `port`                                          | ClickHouse 포트입니다. 기본값은 8443(Cloud의 HTTPS용)이며, HTTP(자체 호스팅의 기본값)를 사용하는 경우 8123이어야 합니다                                                                                                     | `8443`                                                   |
| `ssl`                                           | ClickHouse에 대한 SSL connection을 활성화합니다                                                                                                                                                    | `true`                                                   |
| `jdbcConnectionProperties`                      | ClickHouse에 연결할 때 사용하는 connection 속성입니다. `?`로 시작해야 하며, `param=value`는 `&`로 연결해야 합니다                                                                                                      | `""`                                                     |
| `username`                                      | ClickHouse 데이터베이스 사용자 이름                                                                                                                                                                 | `default`                                                |
| `password` (필수)                                 | ClickHouse 데이터베이스 비밀번호                                                                                                                                                                   | N/A                                                      |
| `database`                                      | ClickHouse 데이터베이스 이름                                                                                                                                                                     | `default`                                                |
| `connector.class` (필수)                          | connector class(명시적으로 설정하고 기본값으로 유지)                                                                                                                                                     | `"com.clickhouse.kafka.connect.ClickHouseSinkConnector"` |
| `tasks.max`                                     | connector 작업 수                                                                                                                                                                           | `"1"`                                                    |
| `errors.retry.timeout`                          | Kafka Connect의 최대 재시도 기간(밀리초 단위)입니다. 재시도하지 않으려면 `0`으로 설정합니다. 무한 재시도는 `-1`입니다. 권장 값은 "10000"ms(10초)보다 큽니다. 시간 초과                                                                          | `"0"`                                                    |
| `exactlyOnce`                                   | Exactly Once 활성화                                                                                                                                                                         | `"false"`                                                |
| `topics` (필수)                                   | 폴링할 Kafka topic - topic 이름은 table 이름과 일치해야 합니다                                                                                                                                           | `""`                                                     |
| `key.converter` (필수\* - 설명 참고)                  | 키의 타입에 맞게 설정하세요. 키를 전달하는 경우(그리고 worker config에 정의되어 있지 않은 경우) 여기서 필수입니다.                                                                                                                 | `"org.apache.kafka.connect.storage.StringConverter"`     |
| `value.converter` (필수\* - 설명 참고)                | topic의 데이터 타입에 따라 설정하세요. 지원 포맷: JSON, String, Avro 또는 Protobuf 형식. worker config에 정의되어 있지 않으면 여기서 필수입니다.                                                                                 | `"org.apache.kafka.connect.json.JsonConverter"`          |
| `value.converter.schemas.enable`                | connector 값 컨버터 스키마 지원                                                                                                                                                                   | `"false"`                                                |
| `errors.tolerance`                              | connector 오류 허용 범위. 지원값: none, all                                                                                                                                                       | `"none"`                                                 |
| `errors.deadletterqueue.topic.name`             | 설정된 경우(errors.tolerance=all과 함께), 실패한 배치에 대해 DLQ가 사용됩니다([문제 해결](#troubleshooting) 참조)                                                                                                    | `""`                                                     |
| `errors.deadletterqueue.context.headers.enable` | DLQ에 추가 headers를 더합니다                                                                                                                                                                    | `""`                                                     |
| `clickhouseSettings`                            | ClickHouse 설정의 쉼표로 구분된 목록(예: "insert\_quorum=2, etc...")                                                                                                                                 | `""`                                                     |
| `topic2TableMap`                                | topic 이름을 table 이름에 매핑하는 쉼표로 구분된 목록(예: "topic1=table1, topic2=table2, etc...")                                                                                                           | `""`                                                     |
| `tableRefreshInterval`                          | table definition cache를 갱신하는 시간(초)                                                                                                                                                       | `0`                                                      |
| `keeperOnCluster`                               | 자체 호스팅 인스턴스의 exactly-once connect\_state 테이블에 대해 ON CLUSTER 매개변수를 구성할 수 있습니다(예: `ON CLUSTER clusterNameInConfigFileDefinition`, [분산 DDL 쿼리](/docs/ko/reference/statements/distributed-ddl) 참조 | `""`                                                     |
| `bypassRowBinary`                               | 스키마 기반 데이터(Avro, Protobuf 등)에서 RowBinary 및 RowBinaryWithDefaults 사용을 비활성화할 수 있습니다. 데이터에 누락된 컬럼이 있고 널 허용/기본값을 사용할 수 없는 경우에만 사용해야 합니다.                                                     | `"false"`                                                |
| `dateTimeFormats`                               | DateTime64 스키마 필드를 파싱할 때 사용하는 날짜/시간 포맷으로, `;`로 구분합니다(예: `someDateField=yyyy-MM-dd HH:mm:ss.SSSSSSSSS;someOtherDateField=yyyy-MM-dd HH:mm:ss`).                                           | `""`                                                     |
| `tolerateStateMismatch`                         | 커넥터가 AFTER\_PROCESSING에 저장된 현재 오프셋보다 "이전"인 레코드(예: 오프셋 5가 전송되었는데 마지막으로 기록된 오프셋이 250인 경우)를 삭제할 수 있도록 합니다. 장애 이후 수집을 복구할 때 사용해야 하며, 완료되면 다시 `"false"`로 되돌려야 합니다.                            | `"false"`                                                |
| `ignorePartitionsWhenBatching`                  | 삽입할 메시지를 수집할 때 파티션을 무시합니다(단, `exactlyOnce`가 `false`인 경우에만 적용됩니다). 성능 참고: connector 작업 수가 많을수록 작업당 할당되는 Kafka 파티션 수는 줄어들며, 이로 인해 성능 향상 효과가 점차 감소할 수 있습니다.                                 | `"false"`                                                |
| `bufferCount` (v1.3.6부터)                        | ClickHouse에 플러시하기 전에 메모리에 버퍼링할 레코드 수입니다. `0`으로 설정하면 내부 버퍼링이 비활성화됩니다. `exactlyOnce=true`에서는 버퍼링이 지원되지 않습니다.                                                                               | `"0"`                                                    |
| `bufferFlushTime` (v1.3.6부터)                    | `exactlyOnce=false`일 때 플러시 전에 레코드를 버퍼링할 수 있는 최대 시간(밀리초)입니다. `0`으로 설정하면 시간 기반 플러시가 비활성화됩니다. 기본값은 `0`입니다. 시간 기반 임계값에만 필요합니다. `bufferCount > 0`일 때만 적용됩니다.                                  | `"0"`                                                    |
| `reportInsertedOffsets` (v1.3.6부터)              | `exactlyOnce=false`일 때 `preCommit`에서 `currentOffsets` 대신 성공적으로 삽입된 오프셋만 반환하도록 합니다. 다만 `ignorePartitionsWhenBatching=true`인 경우에는 적용되지 않으며, 이때는 여전히 `currentOffsets`가 반환됩니다.               | `"false"`                                                |

<div id="target-tables">
  ### 대상 테이블
</div>

ClickHouse Connect Sink는 Kafka 토픽의 메시지를 읽어 적절한 테이블에 씁니다. ClickHouse Connect Sink는 데이터를 기존 테이블에 씁니다. 데이터 삽입을 시작하기 전에 적절한 스키마(schema)를 갖춘 대상 테이블(target table)이 ClickHouse에 생성되어 있는지 확인하십시오.

각 토픽에는 ClickHouse의 전용 대상 테이블이 필요합니다. 대상 테이블 이름은 원본 토픽 이름과 일치해야 합니다.

<div id="pre-processing">
  ### 사전 처리
</div>

메시지를 ClickHouse Kafka Connect
싱크로 보내기 전에 아웃바운드 메시지를 변환해야 한다면 [Kafka Connect Transformations](https://docs.confluent.io/platform/current/connect/transforms/overview.html)을 사용하십시오.

<div id="supported-data-types">
  ### 지원되는 데이터 타입
</div>

**스키마가 선언된 경우:**

| Kafka Connect 타입                        | ClickHouse 타입           | 지원 여부 | 기본형 |
| --------------------------------------- | ----------------------- | ----- | --- |
| STRING                                  | String                  | ✅     | 예   |
| STRING                                  | JSON. 아래 (1) 참조         | ✅     | 예   |
| INT8                                    | Int8                    | ✅     | 예   |
| INT16                                   | Int16                   | ✅     | 예   |
| INT32                                   | Int32                   | ✅     | 예   |
| INT64                                   | Int64                   | ✅     | 예   |
| FLOAT32                                 | Float32                 | ✅     | 예   |
| FLOAT64                                 | Float64                 | ✅     | 예   |
| BOOLEAN                                 | Boolean                 | ✅     | 예   |
| ARRAY                                   | Array(T)                | ✅     | 아니요 |
| MAP                                     | Map(Primitive, T)       | ✅     | 아니요 |
| STRUCT                                  | Variant(T1, T2, ...)    | ✅     | 아니요 |
| STRUCT                                  | Tuple(a T1, b T2, ...)  | ✅     | 아니요 |
| STRUCT                                  | Nested(a T1, b T2, ...) | ✅     | 아니요 |
| STRUCT                                  | JSON. 아래 (1), (2) 참조    | ✅     | 아니요 |
| BYTES                                   | String                  | ✅     | 아니요 |
| org.apache.kafka.connect.data.Time      | Int64 / DateTime64      | ✅     | 아니요 |
| org.apache.kafka.connect.data.Timestamp | Int32 / Date32          | ✅     | 아니요 |
| org.apache.kafka.connect.data.Decimal   | Decimal                 | ✅     | 아니요 |

* (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` 설정](https://docs.confluent.io/platform/current/schema-registry/connect.html#protobuf)입니다.

**스키마가 선언되지 않은 경우:**

레코드는 JSON으로 변환된 후 [JSONEachRow](/docs/ko/reference/formats/JSON/JSONEachRow) 포맷의 값으로 ClickHouse에 전송됩니다.

<div id="configuration-recipes">
  ### 구성 예시
</div>

빠르게 시작할 수 있도록 자주 사용하는 몇 가지 구성 예시를 소개합니다.

<div id="basic-configuration">
  #### 기본 구성
</div>

시작을 위한 가장 기본적인 구성입니다. Kafka Connect는 분산 모드로 실행되고, SSL이 활성화된 `localhost:8443`에서 ClickHouse 서버가 실행 중이며, 데이터는 스키마가 없는 JSON 형식이라고 가정합니다.

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    "tasks.max": "1",
    "consumer.override.max.poll.records": "5000",
    "consumer.override.max.partition.fetch.bytes": "5242880",
    "database": "default",
    "errors.retry.timeout": "60000",
    "exactlyOnce": "false",
    "hostname": "localhost",
    "port": "8443",
    "ssl": "true",
    "jdbcConnectionProperties": "?ssl=true&sslmode=strict",
    "username": "default",
    "password": "<PASSWORD>",
    "topics": "<TOPIC_NAME>",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    "clickhouseSettings": ""
  }
}
```

<Note>
  위의 커넥터 구성에서는 워커 구성에서 `connector.client.config.override.policy=All`을 설정해 클라이언트 재정의를 활성화해야 합니다. 자세한 내용은 [Kafka Connect 문서](https://docs.confluent.io/platform/current/connect/references/allconfigs.html#override-the-worker-configuration)를 참조하십시오.
</Note>

<div id="basic-configuration-with-multiple-topics">
  #### 여러 토픽을 사용하는 기본 구성
</div>

커넥터는 여러 토픽에서 데이터를 수집할 수 있습니다.

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "topics": "SAMPLE_TOPIC, ANOTHER_TOPIC, YET_ANOTHER_TOPIC",
    ...
  }
}
```

<div id="basic-configuration-with-dlq">
  #### DLQ를 포함한 기본 구성
</div>

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "<DLQ_TOPIC>",
    "errors.deadletterqueue.context.headers.enable": "true",
  }
}
```

<div id="avro-schema-support">
  ### Avro 스키마 지원
</div>

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "<SCHEMA_REGISTRY_HOST>:<PORT>",
    "value.converter.schemas.enable": "true",
  }
}
```

<div id="avro-type-mapping">
  #### Avro 타입 매핑
</div>

아래의 타입 매핑은 Kafka Connect의 공식 Avro 직렬화/역직렬화 구현인 `io.confluent.connect.avro.AvroConverter`에서 정의합니다. 변환 로직에 대한 자세한 내용은 Kafka Connect [문서](https://docs.confluent.io/platform/current/connect/userguide.html#avro)를 참조하십시오.

✅: 지원

❌: 미지원

️⚠️: 부분 지원

| Avro Type | Kafka Connect Type | Supported | Notes                                                                                                                                                                                                                                                    |
| --------- | ------------------ | --------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| null      | *N/A*              | ❌         | 독립형 타입으로는 지원되지 않지만 유니온에서는 사용할 수 있습니다                                                                                                                                                                                                                     |
| boolean   | BOOLEAN            | ✅         |                                                                                                                                                                                                                                                          |
| int       | INT8/INT16/INT32   | ✅         | 기본값은 INT32입니다. 스키마에 `connect.type=int8` 속성이 있으면 INT8로 해석됩니다 (`connect.type=int16`인 경우 INT16도 동일)                                                                                                                                                         |
| long      | INT64              | ✅         |                                                                                                                                                                                                                                                          |
| float     | FLOAT32            | ✅         |                                                                                                                                                                                                                                                          |
| double    | FLOAT64            | ✅         |                                                                                                                                                                                                                                                          |
| bytes     | BYTES              | ✅         |                                                                                                                                                                                                                                                          |
| string    | STRING             | ✅         |                                                                                                                                                                                                                                                          |
| record    | STRUCT             | ✅         |                                                                                                                                                                                                                                                          |
| enum      | STRING             | ✅         |                                                                                                                                                                                                                                                          |
| array     | ARRAY/MAP          | ✅         | 기본값은 ARRAY입니다. 필드가 원래 `AvroData.fromConnectSchema`를 통해 생성된 경우 MAP으로 해석됩니다 ([source](https://github.com/confluentinc/schema-registry/blob/174907bfc0d9424e8d02e788f450f4afcdda1750/avro-data/src/main/java/io/confluent/connect/avro/AvroData.java#L943)) |
| map       | MAP                | ✅         |                                                                                                                                                                                                                                                          |
| union     | STRUCT/`<T>`       | ⚠️        | 기본값은 STRUCT입니다. `flatten.singleton.unions=true`인 경우 유니온 정의의 singleton 타입 `T`로 해석됩니다 ([docs](https://docs.confluent.io/cloud/current/connectors/reference/connector-configuration.html#value-converter-flatten-singleton-unions) 참조)                      |
| fixed     | BYTES              | ⚠️        | Fixed `decimal` 논리 타입은 지원되지 않습니다 (아래 참조)                                                                                                                                                                                                                 |

Kafka Connect 타입과 ClickHouse 타입 간의 매핑은 [지원되는 데이터 타입](#supported-data-types)을 참조하십시오.

<div id="unsupported-avro-schemas">
  #### 지원되지 않는 Avro 스키마
</div>

커넥터는 다음 Avro 스키마를 지원하지 않습니다:

* `fixed` `decimal` 논리 타입

```json theme={null}
{"name": "decimal_18_4", "type": "fixed", "size": 8, "logicalType": "decimal", "precision": 18, "scale": 4}
```

* 널 허용 유니온

```json theme={null}
{"name": "mixed_union", "type": ["null", "string", "int"], "default": null}
```

* 레코드 유니온

```json theme={null}
{
  "name": "record_union",
  "type": [
    {
      "type": "record",
      "name": "TypeA",
      "fields": [
        {
          "name": "label",
          "type": "string"
        }
      ]
    },
    {
      "type": "record",
      "name": "TypeB",
      "fields": [
        {
          "name": "count",
          "type": "int"
        }
      ]
    }
  ]
}
```

<div id="protobuf-schema-support">
  ### Protobuf 스키마 지원
</div>

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "io.confluent.connect.protobuf.ProtobufConverter",
    "value.converter.schema.registry.url": "<SCHEMA_REGISTRY_HOST>:<PORT>",
    "value.converter.schemas.enable": "true",
  }
}
```

참고: 누락된 클래스 관련 문제가 발생하는 경우도 있습니다. 모든 환경에 protobuf 컨버터가 포함되는 것은 아니므로, 종속성이 함께 번들된 대체 jar 릴리스가 필요할 수 있습니다.

<div id="proto-type-mapping">
  #### Protobuf 타입 매핑
</div>

아래 타입 매핑은 Kafka Connect의 공식 Protobuf 직렬화/역직렬화 구현인 `io.confluent.connect.protobuf.ProtobufConverter`에서 정의합니다. 변환 로직에 관한 자세한 내용은 Kafka Connect [docs](https://docs.confluent.io/platform/current/connect/userguide.html#json-schema-and-protobuf)를 참조하십시오.

✅: 지원됨

❌: 지원되지 않음

️⚠️: 부분 지원

| Protobuf 타입                             | Kafka Connect 타입                        | ClickHouse 타입                                  | 지원 여부 | 참고                                                                                                                                                                |
| :-------------------------------------- | :-------------------------------------- | :--------------------------------------------- | :---- | :---------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| double                                  | FLOAT64                                 | Float64                                        | ✅     |                                                                                                                                                                   |
| float                                   | FLOAT32                                 | Float32                                        | ✅     |                                                                                                                                                                   |
| int32                                   | INT8/INT16/INT32                        | Int32                                          | ✅     | 기본값은 INT32입니다. 스키마에 `connect.type=int8` 옵션이 있으면 INT8로 해석됩니다(`connect.type=int16`이면 INT16도 동일한 방식으로 처리됨).                                                          |
| sint32                                  | INT8/INT16/INT32                        | Int32                                          | ✅     | 기본값은 INT32입니다. 스키마에 `connect.type=int8` 옵션이 있으면 INT8로 해석됩니다(`connect.type=int16`이면 INT16도 동일한 방식으로 처리됨).                                                          |
| sfixed32                                | INT8/INT16/INT32                        | Int32                                          | ✅     | 기본값은 INT32입니다. 스키마에 `connect.type=int8` 옵션이 있으면 INT8로 해석됩니다(`connect.type=int16`이면 INT16도 동일한 방식으로 처리됨).                                                          |
| uint32                                  | INT64                                   | UInt32                                         | ✅     |                                                                                                                                                                   |
| fixed32                                 | INT64                                   | UInt32                                         | ✅     |                                                                                                                                                                   |
| int64                                   | INT64                                   | Int64                                          | ✅     |                                                                                                                                                                   |
| uint64                                  | INT64                                   | UInt64                                         | ✅     |                                                                                                                                                                   |
| sint64                                  | INT64                                   | Int64                                          | ✅     |                                                                                                                                                                   |
| fixed64                                 | INT64                                   | UInt64                                         | ✅     |                                                                                                                                                                   |
| sfixed64                                | INT64                                   | Int64                                          | ✅     |                                                                                                                                                                   |
| bool                                    | BOOLEAN                                 | Bool                                           | ✅     |                                                                                                                                                                   |
| string                                  | STRING                                  | String                                         | ✅     |                                                                                                                                                                   |
| bytes                                   | BYTES                                   | String                                         | ✅     |                                                                                                                                                                   |
| enum                                    | INT32/STRING                            | Int32                                          | ✅     | 기본값은 STRING입니다. `int.for.enums=true`이면 INT32로 해석됩니다([schema registry docs](https://docs.confluent.io/platform/current/schema-registry/connect.html#protobuf) 참조). |
| message                                 | STRUCT                                  | Tuple / JSON                                   | ⚠️    | 아래의 지원되지 않는 스키마 섹션을 참조하십시오                                                                                                                                        |
| repeated T (where T is not a map entry) | ARRAY                                   | Array(T)                                       | ✅     |                                                                                                                                                                   |
| `map<K, V>`                             | MAP                                     | Map(K, V)                                      | ✅     |                                                                                                                                                                   |
| oneof                                   | STRUCT                                  | Tuple / Variant                                | ⚠️    | `oneof`를 ClickHouse 스키마로 변환하는 방법은 아래 섹션을 참조하십시오                                                                                                                   |
| google.protobuf.DoubleValue             | FLOAT64                                 | Nullable(Float64)                              | ✅     |                                                                                                                                                                   |
| google.protobuf.FloatValue              | FLOAT32                                 | Nullable(Float32)                              | ✅     |                                                                                                                                                                   |
| google.protobuf.Int64Value              | INT64                                   | Nullable(Int64)                                | ✅     |                                                                                                                                                                   |
| google.protobuf.UInt64Value             | INT64                                   | Nullable(UInt64)                               | ✅     |                                                                                                                                                                   |
| google.protobuf.UInt32Value             | INT64                                   | Nullable(UInt32)                               | ✅     |                                                                                                                                                                   |
| google.protobuf.Int32Value              | INT32                                   | Nullable(Int32)                                | ✅     |                                                                                                                                                                   |
| google.protobuf.BoolValue               | BOOLEAN                                 | Nullable(Bool)                                 | ✅     |                                                                                                                                                                   |
| google.protobuf.StringValue             | STRING                                  | Nullable(String)                               | ✅     |                                                                                                                                                                   |
| google.protobuf.BytesValue              | BYTES                                   | Nullable(String)                               | ✅     |                                                                                                                                                                   |
| google.protobuf.Timestamp               | org.apache.kafka.connect.data.Timestamp | DateTime64(3)                                  | ✅     |                                                                                                                                                                   |
| google.type.Date                        | org.apache.kafka.connect.data.Date      | Date                                           | ✅     |                                                                                                                                                                   |
| google.type.TimeOfDay                   | org.apache.kafka.connect.data.Time      | Int32 / Int64                                  | ✅     |                                                                                                                                                                   |
| google.protobuf.Duration                | STRUCT                                  | Tuple(`seconds` Int64, `nano` Nullable(Int32)) | ✅     |                                                                                                                                                                   |
| google.protobuf.Any                     | *N/A*                                   | *N/A*                                          | ❌     |                                                                                                                                                                   |
| google.protobuf.Empty                   | *N/A*                                   | *N/A*                                          | ❌     |                                                                                                                                                                   |

Kafka Connect 타입과 ClickHouse 타입의 매핑은 [지원되는 데이터 타입](#supported-data-types)을 참조하십시오.

<div id="oneof-translation">
  #### `oneof` 필드를 ClickHouse 컬럼으로 변환할 때 참고 사항
</div>

커넥터는 Protobuf 유니온(`oneof`)을 ClickHouse Variant 타입으로 변환하는 기능을 지원하지 않습니다. 대신 `oneof` 필드를 ClickHouse 테이블 스키마에 각각의 널 허용 필드로 나열하십시오.

예시:

```protobuf theme={null}
syntax = "proto3";

package com.clickhouse.kafka.connect.proto.test;

message StringIntUnion {
  oneof mixed {
    string mixed_string = 2;
    int32 mixed_int = 3;
  }
}

```

다음과 같은 ClickHouse 테이블 정의로 변환됩니다:

```sql theme={null}
CREATE TABLE IF NOT EXISTS `StringIntUnion`
(
    mixed_string Nullable(String),
    mixed_int Nullable(Int32)
) ENGINE = ...;
```

<div id="unsupported-proto-schemas">
  #### 지원되지 않는 Protobuf 스키마
</div>

다음 Protobuf 스키마는 커넥터에서 지원되지 않습니다:

* 다중 메시지 유니온 (**CH 버전 26.1 이전**)

```protobuf theme={null}
syntax = "proto3";

package com.clickhouse.kafka.connect.proto.test;

message TwoRecords {
  oneof payload {
    TypeA type_a = 2;
    TypeB type_b = 3;
  }

  // translates to Nullable(Tuple(label String)) in ClickHouse, which is unsupported
  message TypeA {
    string label = 1;
  }

  // translates to Nullable(Tuple(count Int32)) in ClickHouse, which is unsupported
  message TypeB {
    int32 count = 1;
  }
}
```

CH 버전 26.1부터는 `allow_experimental_nullable_tuple_type=1`로 설정하면 이 스키마가 지원됩니다([이 문서 페이지](/docs/ko/reference/settings/session-settings#allow_experimental_nullable_tuple_type) 참조).

<div id="json-schema-support">
  ### JSON 스키마 지원
</div>

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  }
}
```

<div id="string-support">
  ### String 컨버터 지원
</div>

커넥터는 [JSON](/docs/ko/reference/formats/JSON/JSONEachRow), [CSV](/docs/ko/reference/formats/CSV/CSV), [TSV](/docs/ko/reference/formats/TabSeparated/TabSeparated) 등 다양한 ClickHouse 포맷에서 String 컨버터를 지원합니다.

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "customInsertFormat": "true",
    "insertFormat": "CSV"
  }
}
```

<div id="internal-buffering">
  ### 내부 버퍼링
</div>

내부 버퍼링을 사용하면 싱크 작업이 여러 번의 `poll()` 호출에서 레코드를 누적한 뒤, 이를 더 큰 배치로 묶어 ClickHouse에 플러시할 수 있습니다. 이렇게 하면 각 `poll()`이 파티션별로 작은 배치를 많이 생성하는 워크로드에서 처리량을 높일 수 있습니다.

주요 동작:

* `bufferCount`는 플러시하기 전에 버퍼링할 레코드 수를 제어합니다.
* `bufferFlushTime`은 버퍼링된 레코드를 플러시하기 전까지의 최대 대기 시간(밀리초)을 설정합니다.
* `bufferFlushTime`은 `bufferCount > 0`일 때만 유효합니다.
* `bufferCount=0` 및 `bufferFlushTime=0`이면 버퍼링이 비활성화된 상태로 유지됩니다(기본 동작).
* `exactlyOnce=true`일 때는 버퍼링이 지원되지 않습니다.

버퍼링이 exactly-once 모드와 호환되지 않는 이유:
버퍼링은 배치 경계를 변경하므로 ClickHouse 블록 중복 제거와 커넥터 오프셋 상태 머신이 제대로 동작하지 않게 됩니다.
이 문제를 해결하려면 커넥터 설정에서 `exactlyOnce=false`로 exactly-once 모드를 비활성화하거나, `bufferCount=0`으로 버퍼링을 비활성화하십시오.

예시:

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "exactlyOnce": "false",
    "bufferCount": "5000",
    "bufferFlushTime": "2000"
  }
}
```

<div id="logging">
  ### 로깅
</div>

로깅은 Kafka Connect Platform에서 자동으로 제공됩니다.
로깅 대상과 포맷은 Kafka Connect [설정 파일](https://docs.confluent.io/platform/current/connect/logging.html#log4j-properties-file)을 통해 구성할 수 있습니다.

Confluent Platform을 사용하는 경우 CLI 명령을 실행하여 로그를 확인할 수 있습니다:

```bash theme={null}
confluent local services connect log
```

자세한 내용은 공식 [튜토리얼](https://docs.confluent.io/platform/current/connect/logging.html)을 참조하십시오.

<div id="monitoring">
  ### 모니터링
</div>

ClickHouse Kafka Connect는 [Java Management Extensions (JMX)](https://www.oracle.com/technical-resources/articles/javase/jmx.html)를 통해 런타임 메트릭을 제공합니다. JMX는 Kafka Connector에서 기본적으로 활성화되어 있습니다.

<div id="clickhouse-specific-metrics">
  #### ClickHouse 전용 메트릭
</div>

이 커넥터는 다음 MBean 이름으로 사용자 지정 메트릭을 노출합니다:

```java theme={null}
com.clickhouse:type=ClickHouseKafkaConnector,name=SinkTask{id}
```

| 메트릭 이름                 | 유형   | 설명                                                  |
| ---------------------- | ---- | --------------------------------------------------- |
| `receivedRecords`      | long | 수신된 레코드의 총 개수입니다.                                   |
| `recordProcessingTime` | long | 레코드를 그룹화하고 통합된 구조로 변환하는 데 소요된 총 시간이며, 나노초 단위입니다.    |
| `taskProcessingTime`   | long | 데이터를 처리하고 ClickHouse에 삽입하는 데 소요된 총 시간이며, 나노초 단위입니다. |

<div id="kafka-producer-consumer-metrics">
  #### Kafka 프로듀서/컨슈머 메트릭
</div>

커넥터는 데이터 흐름, 처리량, 성능을 파악할 수 있도록 표준 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`: 초당 평균 응답 수

이러한 메트릭은 다음 항목을 모니터링하는 데 도움이 됩니다.

* **처리량**: 데이터 수집 속도 추적
* **지연**: 병목과 처리 지연 식별
* **압축**: 데이터 압축 효율 측정
* **연결 상태**: 네트워크 연결 및 안정성 모니터링

<div id="kafka-connect-framework-metrics">
  #### Kafka Connect Framework 메트릭
</div>

이 커넥터는 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한 전체 레코드 수

<div id="monitoring-best-practices">
  #### 모니터링 모범 사례
</div>

1. **Consumer lag 모니터링**: 처리 병목을 식별할 수 있도록 파티션별 `records-lag`를 추적합니다
2. **오류율 추적**: 데이터 품질 문제를 감지할 수 있도록 `errors-total` 및 `records-skip-total`을 확인합니다
3. **작업 상태 관찰**: 작업이 정상적으로 실행되는지 확인할 수 있도록 작업 상태 메트릭을 모니터링합니다
4. **처리량 측정**: 수집 성능을 추적하기 위해 `records-send-rate` 및 `byte-rate`를 사용합니다
5. **연결 상태 모니터링**: 네트워크 문제를 확인하기 위해 노드 수준의 연결 메트릭을 점검합니다
6. **압축 효율 추적**: 데이터 전송을 최적화하기 위해 `compression-rate`를 사용합니다

자세한 JMX 메트릭 정의와 Prometheus 통합 방법은 [jmx-export-connector.yml](https://github.com/ClickHouse/clickhouse-kafka-connect/blob/main/jmx-export-connector.yml) 설정 파일을 참조하십시오.

<div id="limitations">
  ### 제한 사항
</div>

* 삭제는 지원되지 않습니다.
* 배치 크기는 Kafka Consumer 속성을 따릅니다.
* exactly-once에 KeeperMap을 사용하는 경우 오프셋이 변경되거나 되돌려지면 해당 토픽의 KeeperMap 내용을 삭제해야 합니다. (자세한 내용은 아래 문제 해결 가이드를 참조하십시오)

<div id="tuning-performance">
  ### 성능 튜닝 및 처리량 최적화
</div>

이 섹션에서는 ClickHouse Kafka Connect Sink의 성능 튜닝 전략을 설명합니다. 성능 튜닝은 처리량이 높은 환경을 다루거나 리소스 활용을 최적화하고 지연을 최소화해야 할 때 필수적입니다.

<div id="when-is-performance-tuning-needed">
  #### 성능 튜닝은 언제 필요합니까?
</div>

일반적으로 다음과 같은 상황에서는 성능 튜닝이 필요합니다.

* **고처리량 워크로드**: Kafka 토픽에서 초당 수백만 건의 이벤트를 처리하는 경우
* **Consumer lag**: 커넥터가 데이터 생성 속도를 따라가지 못해 lag가 계속 증가하는 경우
* **리소스 제약**: CPU, 메모리 또는 네트워크 사용량을 최적화해야 하는 경우
* **여러 토픽**: 대용량 토픽 여러 개를 동시에 소비하는 경우
* **작은 메시지 크기**: 서버 측 배칭의 이점을 얻을 수 있도록 작은 메시지를 대량으로 처리하는 경우

다음과 같은 경우에는 일반적으로 성능 튜닝이 **필요하지 않습니다**.

* 낮거나 중간 수준의 처리량(\< 10,000 messages/second)을 처리하는 경우
* 사용 사례에서 Consumer lag가 안정적이고 허용 가능한 수준인 경우
* 기본 커넥터 설정만으로도 필요한 처리량 요구 사항을 이미 충족하는 경우
* ClickHouse 클러스터가 유입되는 부하를 무리 없이 처리할 수 있는 경우

<div id="understanding-the-data-flow">
  #### 데이터 흐름 이해하기
</div>

튜닝에 앞서 데이터가 커넥터를 통해 어떻게 흐르는지 이해하는 것이 중요합니다.

1. **Kafka Connect Framework**가 백그라운드에서 Kafka 토픽의 메시지를 가져옵니다
2. **커넥터는 폴링을 통해** 프레임워크의 내부 버퍼에서 메시지를 가져옵니다
3. **커넥터는 메시지를 배치로 묶어** 폴링 크기에 따라 처리합니다
4. **ClickHouse는** HTTP/S를 통해 배치된 삽입을 수신합니다
5. **ClickHouse는** 삽입을 처리합니다(동기식 또는 비동기식)

이 각 단계에서 성능을 최적화할 수 있습니다.

<div id="connect-fetch-vs-connector-poll">
  #### Kafka Connect 배치 크기 튜닝
</div>

첫 번째 최적화 단계는 커넥터가 Kafka에서 배치당 수신하는 데이터 양을 제어하는 것입니다.

<div id="fetch-settings">
  ##### fetch 설정
</div>

Kafka Connect(프레임워크)는 커넥터와 무관하게 백그라운드에서 Kafka 토픽의 메시지를 fetch합니다.

* **`fetch.min.bytes`**: 프레임워크가 값을 커넥터에 전달하기 전에 필요한 최소 데이터 양(기본값: 1 byte)
* **`fetch.max.bytes`**: 단일 요청으로 fetch할 수 있는 최대 데이터 양(기본값: 52428800 / 50 MB)
* **`fetch.max.wait.ms`**: `fetch.min.bytes` 조건이 충족되지 않으면 데이터를 반환하기 전까지 대기하는 최대 시간(기본값: 500 ms)

<Note>
  Confluent Cloud에서는 이러한 설정을 조정하려면 Confluent Cloud 지원 케이스를 열어야 합니다.
</Note>

<div id="poll-settings">
  ##### 폴링 설정
</div>

커넥터는 프레임워크의 버퍼에서 메시지를 폴링합니다.

* **`max.poll.records`**: 한 번의 폴링으로 반환되는 최대 레코드 수(기본값: 500)
* **`max.partition.fetch.bytes`**: 파티션별 최대 데이터 크기(기본값: 1048576 / 1 MB)

<Note>
  Confluent Cloud에서는 이러한 설정을 조정하려면 Confluent Cloud에서 지원 케이스를 접수해야 합니다.
</Note>

<div id="recommended-batch-settings">
  ##### 고처리량을 위한 권장 설정
</div>

ClickHouse에서 최적의 성능을 얻으려면 더 큰 배치를 사용하십시오:

```properties theme={null}
# 폴링당 레코드 수 증가
consumer.override.max.poll.records=5000

# 파티션 fetch 크기 증가 (5 MB)
consumer.override.max.partition.fetch.bytes=5242880

# 선택 사항: 더 많은 데이터를 기다리기 위한 최소 fetch 크기 증가 (1 MB)
consumer.override.fetch.min.bytes=1048576

# 선택 사항: 지연 시간이 중요한 경우 대기 시간 단축
consumer.override.fetch.max.wait.ms=300
```

<Note>
  위 속성을 사용하려면 워커 구성에서 `connector.client.config.override.policy=All`을 통해 클라이언트 재정의를 활성화해야 합니다. 자세한 내용은 [Kafka Connect 문서](https://docs.confluent.io/platform/current/connect/references/allconfigs.html#override-the-worker-configuration)를 참조하십시오.
</Note>

**중요**: Kafka Connect의 fetch 설정은 압축된 데이터를 기준으로 하며, ClickHouse는 비압축 데이터를 수신합니다. 압축률을 고려해 이러한 설정의 균형을 맞추십시오.

**트레이드오프**:

* **더 큰 배치** = 더 나은 ClickHouse 수집 성능, 더 적은 파트, 더 낮은 오버헤드
* **더 큰 배치** = 더 높은 메모리 사용량, 종단 간 지연 시간 증가 가능성
* **너무 큰 배치** = timeout, OutOfMemory 오류 또는 `max.poll.interval.ms` 초과 위험

자세한 내용: [Confluent 문서](https://docs.confluent.io/platform/current/connect/references/allconfigs.html#override-the-worker-configuration) | [Kafka 문서](https://kafka.apache.org/documentation/#consumerconfigs)

<div id="asynchronous-inserts">
  #### 비동기 삽입
</div>

비동기 삽입은 커넥터가 비교적 작은 배치를 전송할 때, 또는 배칭 책임을 ClickHouse로 넘겨 수집을 한층 더 최적화하려는 경우에 매우 유용한 기능입니다.

<div id="when-to-use-async-inserts">
  ##### async 삽입을 사용해야 하는 경우
</div>

다음과 같은 경우 async 삽입 활성화를 고려하십시오.

* **작은 배치가 많은 경우**: 커넥터가 작은 배치(\< 배치당 1000행)를 자주 전송하는 경우
* **높은 동시성**: 여러 커넥터 작업이 동일한 테이블에 쓰기 작업을 수행하는 경우
* **분산 배포**: 서로 다른 호스트에서 많은 커넥터 인스턴스를 실행하는 경우
* **파트 생성 오버헤드**: "too many parts" 오류가 발생하는 경우
* **혼합 워크로드**: 실시간 수집과 쿼리 워크로드를 함께 실행하는 경우

다음과 같은 경우에는 async 삽입을 **사용하지 마십시오**.

* 이미 큰 배치(배치당 10,000행 초과)를 제어된 빈도로 전송하는 경우
* 데이터가 즉시 표시되어야 하는 경우(쿼리에서 데이터를 즉시 확인해야 함)
* `wait_for_async_insert=0`을 사용하는 정확히 한 번 처리 의미 체계가 요구 사항과 충돌하는 경우
* 대신 클라이언트 측 배칭 개선의 이점을 얻을 수 있는 사용 사례인 경우

<div id="how-async-inserts-work">
  ##### 비동기 삽입의 작동 방식
</div>

비동기 삽입이 활성화되면 ClickHouse는 다음과 같이 동작합니다:

1. 커넥터로부터 삽입 쿼리를 받습니다
2. 데이터를 메모리 버퍼에 기록합니다(즉시 디스크에 기록하지 않음)
3. 커넥터에 성공을 반환합니다 (`wait_for_async_insert=0`인 경우)
4. 다음 조건 중 하나가 충족되면 버퍼를 디스크로 플러시합니다:
   * 버퍼 크기가 `async_insert_max_data_size`에 도달함(기본값: 100 MB)
   * 첫 번째 삽입 후 `async_insert_busy_timeout_ms`밀리초가 경과함(기본값: 1000 ms)
   * 누적된 쿼리 수가 최대치에 도달함(`async_insert_max_query_number`, 기본값: 100)

이 방식은 생성되는 파트 수를 크게 줄이고 전체 처리량을 향상시킵니다.

<div id="enabling-async-inserts">
  ##### async 삽입 활성화
</div>

async 삽입 설정을 `clickhouseSettings` 구성 매개변수에 추가합니다:

```json theme={null}
{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "clickhouseSettings": "async_insert=1,wait_for_async_insert=1"
  }
}
```

**주요 설정**:

* **`async_insert=1`**: 비동기 삽입을 활성화합니다
* **`wait_for_async_insert=1`** (권장): 커넥터가 확인 응답을 보내기 전에 데이터가 ClickHouse 스토리지에 플러시될 때까지 기다립니다. 전송 보장을 제공합니다.
* **`wait_for_async_insert=0`**: 버퍼링 직후 커넥터가 즉시 확인 응답을 보냅니다. 성능은 더 좋지만 플러시되기 전에 서버가 충돌하면 데이터가 손실될 수 있습니다.

<div id="tuning-async-inserts">
  ##### async insert 동작 조정
</div>

async insert의 플러시 동작을 세부적으로 조정할 수 있습니다:

```json theme={null}
"clickhouseSettings": "async_insert=1,wait_for_async_insert=1,async_insert_max_data_size=104857600,async_insert_busy_timeout_ms=1000"
```

일반적인 튜닝 매개변수:

* **`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`인 경우 서버 장애 시 데이터가 손실될 수 있으며, 버퍼가 크면 메모리 사용량 압박이 발생할 수 있습니다

<div id="async-inserts-with-exactly-once">
  ##### 정확히 한 번 처리 의미 체계의 async 삽입
</div>

비동기 삽입에서 `exactlyOnce=true`를 사용하는 경우:

```json theme={null}
{
  "config": {
    "exactlyOnce": "true",
    "clickhouseSettings": "async_insert=1,wait_for_async_insert=1"
  }
}
```

**중요**: 오프셋 커밋이 데이터가 영구 저장된 후에만 수행되도록 하려면 exactly-once와 함께 항상 `wait_for_async_insert=1`을 사용하십시오.

async 삽입에 대한 자세한 내용은 [ClickHouse async 삽입 문서](/docs/ko/concepts/best-practices/selecting-an-insert-strategy#asynchronous-inserts)를 참조하십시오.

<div id="connector-parallelism">
  #### 커넥터 병렬성
</div>

처리량을 높이려면 병렬성을 늘리십시오:

<div id="tasks-per-connector">
  ##### 커넥터별 작업 수
</div>

```json theme={null}
"tasks.max": "4"
```

각 작업은 토픽 파티션의 일부를 처리합니다. 작업 수가 많을수록 병렬성은 높아지지만, 다음 사항에 유의해야 합니다:

* 효과적인 최대 작업 수 = 토픽 파티션 수
* 각 작업은 ClickHouse와의 자체 연결을 유지합니다
* 작업 수가 많을수록 오버헤드가 커지고 리소스 경합이 발생할 가능성이 높아집니다

**권장 사항**: `tasks.max`를 토픽 파티션 수와 같게 설정하여 시작한 다음, CPU 및 처리량 메트릭을 기준으로 조정하세요.

<div id="ignoring-partitions">
  ##### 배칭 시 파티션 무시
</div>

기본적으로 커넥터는 파티션별로 메시지를 배칭합니다. 처리량을 높이려면 파티션 간에 걸쳐 배칭할 수 있습니다.

```json theme={null}
"ignorePartitionsWhenBatching": "true"
```

\*\* 경고\*\*: `exactlyOnce=false`일 때만 사용하십시오. 이 설정을 사용하면 더 큰 배치를 생성해 처리량을 높일 수 있지만, 파티션별 순서 보장은 유지되지 않습니다.

<div id="multiple-high-throughput-topics">
  #### 여러 고처리량 토픽
</div>

커넥터가 여러 토픽을 구독하도록 구성되어 있고, `topic2TableMap`을 사용해 토픽을 테이블에 매핑하며, 삽입 병목으로 인해 consumer lag이 발생하는 경우에는 토픽별로 커넥터를 하나씩 생성하는 방안을 고려하십시오.

이 문제가 발생하는 주된 이유는 현재 배치가 각 테이블에 [순차적으로](https://github.com/ClickHouse/clickhouse-kafka-connect/blob/578ac07e8be1a920aaa3b26e49183595c3edd04b/src/main/java/com/clickhouse/kafka/connect/sink/ProxySinkTask.java#L95-L100) 삽입되기 때문입니다.

**권장 사항**: 여러 고처리량 토픽의 경우, 병렬 삽입 처리량을 극대화하려면 토픽별로 커넥터 인스턴스를 하나씩 배포하십시오.

<div id="table-engine-considerations">
  #### ClickHouse 테이블 엔진 고려 사항
</div>

사용 사례에 맞는 ClickHouse 테이블 엔진을 선택하십시오:

* **`MergeTree`**: 대부분의 사용 사례에 가장 적합하며, 쿼리와 삽입 성능의 균형이 좋습니다
* **`ReplicatedMergeTree`**: 고가용성을 위해 필요하며, 복제 오버헤드가 추가됩니다
* **적절한 `ORDER BY`를 사용하는 `*MergeTree`**: 쿼리 패턴에 맞게 최적화하십시오

**고려할 설정**:

```sql theme={null}
CREATE TABLE my_table (...)
ENGINE = MergeTree()
ORDER BY (timestamp, id)
SETTINGS 
    -- 병렬 파트 쓰기를 위한 최대 삽입 스레드 수 증가
    max_insert_threads = 4,
    -- 안정성을 위해 쿼럼(quorum) 삽입 허용 (ReplicatedMergeTree)
    insert_quorum = 2
```

커넥터 수준 삽입 설정:

```json theme={null}
"clickhouseSettings": "insert_quorum=2,insert_quorum_timeout=60000"
```

<div id="connection-pooling">
  #### 연결 풀링 및 타임아웃
</div>

이 커넥터는 ClickHouse와의 HTTP 연결을 유지합니다. 지연 시간이 큰 네트워크에서는 타임아웃을 조정하십시오:

```json theme={null}
"clickhouseSettings": "socket_timeout=300000,connection_timeout=30000"
```

* **`socket_timeout`** (기본값: 30000 ms): 읽기 작업의 최대 대기 시간
* **`connection_timeout`** (기본값: 10000 ms): 연결을 설정하는 최대 대기 시간

대용량 배치에서 timeout 오류가 발생하면 이 값을 늘리십시오.

<div id="monitoring-performance">
  #### 성능 모니터링 및 문제 해결
</div>

다음 핵심 메트릭을 모니터링하세요:

1. **Consumer lag**: Kafka 모니터링 도구를 사용해 파티션별 지연을 추적합니다
2. **커넥터 메트릭**: JMX를 통해 `receivedRecords`, `recordProcessingTime`, `taskProcessingTime`를 모니터링합니다([모니터링](#monitoring) 참조)
3. **ClickHouse 메트릭**:
   * `system.asynchronous_inserts`: async insert 버퍼 사용량을 모니터링합니다
   * `system.parts`: 머지 문제를 감지할 수 있도록 파트 수를 모니터링합니다
   * `system.merges`: 진행 중인 머지를 모니터링합니다
   * `system.events`: `InsertedRows`, `InsertedBytes`, `FailedInsertQuery`를 추적합니다

**일반적인 성능 문제**:

| 증상                  | 가능한 원인             | 해결 방법                                                 |
| ------------------- | ------------------ | ----------------------------------------------------- |
| Consumer lag 증가     | 배치가 너무 작음          | `max.poll.records`를 늘리고 async 삽입을 활성화합니다              |
| "Too many parts" 오류 | 작고 빈번한 삽입          | async 삽입을 활성화하고 배치 크기를 늘립니다                           |
| timeout 오류          | 배치 크기가 큼, 네트워크가 느림 | 배치 크기를 줄이고 `socket_timeout`을 늘린 다음 네트워크를 확인합니다        |
| 높은 CPU 사용량          | 작은 파트가 너무 많음       | async 삽입을 활성화하고 머지 설정을 높입니다                           |
| OutOfMemory 오류      | 배치 크기가 너무 큼        | `max.poll.records`, `max.partition.fetch.bytes`를 줄입니다 |
| 작업 부하 불균형           | 파티션 분배가 고르지 않음     | 파티션을 리밸런싱하거나 `tasks.max`를 조정합니다                       |

<div id="performance-best-practices">
  #### 모범 사례 요약
</div>

1. **기본값으로 시작**한 다음 실제 성능을 측정하고 그 결과에 따라 조정하세요
2. **더 큰 배치를 우선하세요**: 가능하면 한 번의 삽입당 10,000\~100,000개 행을 목표로 하세요
3. **작은 배치를 많이 보내거나 동시성(Concurrency)이 높을 때는 async 삽입을 사용하세요**
4. **정확히 한 번 처리 의미 체계를 위해 `wait_for_async_insert=1`을 항상 사용하세요**
5. **수평 확장하세요**: `tasks.max`를 파티션 수만큼 늘리세요
6. **처리량을 최대화하려면 대용량 토픽마다 커넥터를 하나씩 사용하세요**
7. **지속적으로 모니터링하세요**: consumer lag, part 개수, 머지 활동을 추적하세요
8. **충분히 테스트하세요**: 프로덕션 배포 전에 반드시 실제와 유사한 부하에서 구성 변경을 테스트하세요

<div id="example-high-throughput">
  #### 예시: 고처리량 구성
</div>

다음은 고처리량에 맞게 최적화한 전체 예시입니다:

```json theme={null}
{
  "name": "clickhouse-high-throughput",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    "tasks.max": "8",
    
    "topics": "high_volume_topic",
    "hostname": "my-clickhouse-host.cloud",
    "port": "8443",
    "database": "default",
    "username": "default",
    "password": "<PASSWORD>",
    "ssl": "true",
    
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    
    "exactlyOnce": "false",
    "ignorePartitionsWhenBatching": "true",
    
    "consumer.override.max.poll.records": "10000",
    "consumer.override.max.partition.fetch.bytes": "5242880",
    "consumer.override.fetch.min.bytes": "1048576",
    "consumer.override.fetch.max.wait.ms": "500",
    
    "clickhouseSettings": "async_insert=1,wait_for_async_insert=1,async_insert_max_data_size=16777216,async_insert_busy_timeout_ms=1000,socket_timeout=300000"
  }
}
```

<Note>
  위 커넥터 구성에서는 worker 구성에서 `connector.client.config.override.policy=All`을 설정해 클라이언트 재정의를 활성화해야 합니다. 자세한 내용은 [Kafka Connect 문서](https://docs.confluent.io/platform/current/connect/references/allconfigs.html#override-the-worker-configuration)를 참조하십시오.
</Note>

**이 구성은**:

* 폴링 한 번에 최대 10,000개의 레코드를 처리합니다
* 더 큰 삽입을 위해 여러 파티션의 데이터를 배치로 묶습니다
* 16 MB 버퍼를 사용하는 async 삽입을 사용합니다
* 8개의 병렬 작업을 실행합니다(파티션 수에 맞춰 설정)
* 엄격한 순서 보장보다 처리량에 중점을 두고 최적화되었습니다

<div id="troubleshooting">
  ### 문제 해결
</div>

<div id="state-mismatch-for-topic-sometopic-partition-0">
  #### "토픽 `[someTopic]` 파티션 `[0]`의 상태가 일치하지 않음"
</div>

이 오류는 KeeperMap에 저장된 오프셋과 Kafka에 저장된 오프셋이 서로 다를 때 발생하며, 보통 토픽이 삭제되었거나
오프셋이 수동으로 조정된 경우에 나타납니다.
이 문제를 해결하려면 해당 토픽 + 파티션에 저장된 기존 값을 삭제해야 합니다:

```sql theme={null}
-- 먼저, 데이터를 저장하는 데 사용된 데이터베이스를 확인합니다.
SELECT * FROM [database].connect_state

-- 해당 토픽과 파티션에 해당하는 키를 확인합니다.
ALTER TABLE [database].connect_state DELETE WHERE key = [keyname]
```

<Note>
  이 조정은 exactly-once 동작에 영향을 줄 수 있습니다.
</Note>

<div id="what-errors-will-the-connector-retry">
  #### "커넥터는 어떤 오류를 재시도합니까?"
</div>

현재는 일시적이어서 재시도할 수 있는 오류를 식별하는 데 중점을 두고 있으며, 여기에는 다음이 포함됩니다:

* `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` - 네트워크에 문제가 있을 때 이 예외가 발생합니다.

<div id="all-my-data-is-blankzeroes">
  #### "모든 데이터가 비어 있거나 0입니다"
</div>

데이터의 필드가 테이블의 필드와 일치하지 않을 가능성이 높습니다. 이는 특히 CDC(및 Debezium 포맷)에서 자주 발생합니다.
일반적인 해결 방법 중 하나는 connector 구성에 flatten 변환을 추가하는 것입니다:

```properties theme={null}
transforms=flatten
transforms.flatten.type=org.apache.kafka.connect.transforms.Flatten$Value
transforms.flatten.delimiter=_
```

이렇게 하면 데이터가 중첩된 JSON에서 평탄화된 JSON으로 변환됩니다(`_`를 구분자(Delimiter)로 사용). 그러면 테이블(table)의 필드는 "field1\_field2\_field3" 포맷을 따르게 됩니다(예: "before\_id", "after\_id" 등).

<div id="i-want-to-use-my-kafka-keys-in-clickhouse">
  #### "ClickHouse에서 Kafka 키를 사용하고 싶습니다"
</div>

기본적으로 Kafka 키는 value 필드에 저장되지 않지만, `KeyToValue` 변환을 사용하면 키를 value 필드의 새 `_key` 필드로 옮길 수 있습니다:

```properties theme={null}
transforms=keyToValue
transforms.keyToValue.type=com.clickhouse.kafka.connect.transforms.KeyToValue
transforms.keyToValue.field=_key
```
