> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-mintlify-fbfa8bee.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

> Apache Flink 与 ClickHouse 入门

# Flink 连接器

export const ClickHouseSupportedBadge = () => {
  return <div className="ClickHouseSupportedBadge">
            <div className="ClickHouseSupportedIcon">
                <svg width="16" height="16" viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
                    <path d="M1.30762 1.39073C1.30762 1.3103 1.37465 1.22986 1.46849 1.22986H2.64824C2.72868 1.22986 2.80912 1.29689 2.80912 1.39073V14.4886C2.80912 14.5691 2.74209 14.6495 2.64824 14.6495H1.46849C1.38805 14.6495 1.30762 14.5825 1.30762 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M4.2832 1.39073C4.2832 1.3103 4.35023 1.22986 4.44408 1.22986H5.62383C5.70427 1.22986 5.7847 1.29689 5.7847 1.39073V14.4886C5.7847 14.5691 5.71767 14.6495 5.62383 14.6495H4.44408C4.36364 14.6495 4.2832 14.5825 4.2832 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M7.25977 1.39073C7.25977 1.3103 7.3268 1.22986 7.42064 1.22986H8.60039C8.68083 1.22986 8.76127 1.29689 8.76127 1.39073V14.4886C8.76127 14.5691 8.69423 14.6495 8.60039 14.6495H7.42064C7.3402 14.6495 7.25977 14.5825 7.25977 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M10.2354 1.39073C10.2354 1.3103 10.3024 1.22986 10.3962 1.22986H11.576C11.6564 1.22986 11.7369 1.29689 11.7369 1.39073V14.4886C11.7369 14.5691 11.6698 14.6495 11.576 14.6495H10.3962C10.3158 14.6495 10.2354 14.5825 10.2354 14.4886V1.39073Z" fill="currentColor" />
                    <path d="M13.2256 6.6057C13.2256 6.52526 13.2926 6.44482 13.3865 6.44482H14.5662C14.6466 6.44482 14.7271 6.51186 14.7271 6.6057V9.27354C14.7271 9.35398 14.6601 9.43442 14.5662 9.43442H13.3865C13.306 9.43442 13.2256 9.36739 13.2256 9.27354V6.6057Z" fill="currentColor" />
                </svg>
            </div>
            支持 ClickHouse
        </div>;
};

这是由 ClickHouse 官方支持的 [Apache Flink Sink Connector](https://github.com/ClickHouse/flink-connector-clickhouse)。它基于 Flink 的 [AsyncSinkBase](https://cwiki.apache.org/confluence/display/FLINK/FLIP-171%3A+Async+Sink) 和官方 ClickHouse [Java 客户端](https://github.com/ClickHouse/clickhouse-java) 构建。

该 连接器 支持 Apache Flink 的 DataStream API。对 Table API 的支持[计划在后续 release 中提供](https://github.com/ClickHouse/flink-connector-clickhouse/issues/42)。

<div id="requirements">
  ## 要求
</div>

* Java 11+ (用于 Flink 1.17+) 或 17+ (用于 Flink 2.0+)
* Apache Flink 1.17+

<div id="flink-compatibility-matrix">
  ## Flink 版本兼容性矩阵
</div>

该连接器分为两个制品，以同时支持 Flink 1.17+ 和 Flink 2.0+。请选择与所用 Flink 版本对应的制品：

| Flink 版本 | 制品                               | ClickHouse Java 客户端 版本 | 所需 Java  |
| -------- | -------------------------------- | ---------------------- | -------- |
| latest   | flink-connector-clickhouse-2.0.0 | 0.9.5                  | Java 17+ |
| 2.0.1    | flink-connector-clickhouse-2.0.0 | 0.9.5                  | Java 17+ |
| 2.0.0    | flink-connector-clickhouse-2.0.0 | 0.9.5                  | Java 17+ |
| 1.20.2   | flink-connector-clickhouse-1.17  | 0.9.5                  | Java 11+ |
| 1.19.3   | flink-connector-clickhouse-1.17  | 0.9.5                  | Java 11+ |
| 1.18.1   | flink-connector-clickhouse-1.17  | 0.9.5                  | Java 11+ |
| 1.17.2   | flink-connector-clickhouse-1.17  | 0.9.5                  | Java 11+ |

<Note>
  该连接器尚未针对早于 Flink 1.17.2 的版本进行测试。
</Note>

<div id="installation--setup">
  ## 安装与设置
</div>

<div id="import-as-a-dependency">
  ### 作为依赖引入
</div>

<div id="flink-2">
  #### 对于 Flink 2.0+
</div>

<Tabs>
  <Tab title="Maven">
    ```maven theme={null}
    <dependency>
        <groupId>com.clickhouse.flink</groupId>
        <artifactId>flink-connector-clickhouse-2.0.0</artifactId>
        <version>{{ stable_version }}</version>
        <classifier>all</classifier>
    </dependency>
    ```
  </Tab>

  <Tab title="Gradle">
    ```gradle theme={null}
    dependencies {
        implementation("com.clickhouse.flink:flink-connector-clickhouse-2.0.0:{{ stable_version }}")
    }
    ```
  </Tab>

  <Tab title="SBT">
    ```sbt theme={null}
    libraryDependencies += "com.clickhouse.flink" % "flink-connector-clickhouse-2.0.0" % {{ stable_version }} classifier "all"
    ```
  </Tab>
</Tabs>

<div id="flink-117">
  #### 适用于 Flink 1.17 及以上版本
</div>

<Tabs>
  <Tab title="Maven">
    ```maven theme={null}
    <dependency>
        <groupId>com.clickhouse.flink</groupId>
        <artifactId>flink-connector-clickhouse-1.17</artifactId>
        <version>{{ stable_version }}</version>
        <classifier>all</classifier>
    </dependency>
    ```
  </Tab>

  <Tab title="Gradle">
    ```gradle theme={null}
    dependencies {
        implementation("com.clickhouse.flink:flink-connector-clickhouse-1.17:{{ stable_version }}")
    }
    ```
  </Tab>

  <Tab title="SBT">
    ```sbt theme={null}
    libraryDependencies += "com.clickhouse.flink" % "flink-connector-clickhouse-1.17" % {{ stable_version }} classifier "all"
    ```
  </Tab>
</Tabs>

<div id="download-the-binary">
  ### 下载二进制包
</div>

二进制 JAR 的命名规则如下：

```bash theme={null}
flink-connector-clickhouse-${flink_version}-${stable_version}-all.jar
```

其中：

* `flink_version` 为 `2.0.0` 或 `1.17`
* `stable_version` 是[稳定版本制品的发布版本](https://github.com/ClickHouse/flink-connector-clickhouse/releases)

你可以在 [Maven Central 仓库](https://repo1.maven.org/maven2/com/clickhouse/flink/) 中找到所有已发布的 JAR 文件。

<div id="using-the-datastream-api">
  ## 使用 DataStream API
</div>

<div id="datastream-snippet">
  ### 代码示例
</div>

假设你想将原始 CSV 数据插入到 ClickHouse：

<Tabs>
  <Tab title="Java">
    ```java theme={null}
    public static void main(String[] args) {
        // 配置 ClickHouseClient
        ClickHouseClientConfig clickHouseClientConfig = new ClickHouseClientConfig(url, username, password, database, tableName);

        // 创建一个 ElementConverter
        ElementConverter<String, ClickHousePayload> convertorString = new ClickHouseConvertor<>(String.class);

        // 创建 sink，并使用 `setClickHouseFormat` 设置格式
        ClickHouseAsyncSink<String> csvSink = new ClickHouseAsyncSink<>(
                convertorString,
                MAX_BATCH_SIZE,
                MAX_IN_FLIGHT_REQUESTS,
                MAX_BUFFERED_REQUESTS,
                MAX_BATCH_SIZE_IN_BYTES,
                MAX_TIME_IN_BUFFER_MS,
                MAX_RECORD_SIZE_IN_BYTES,
                clickHouseClientConfig
        );

        csvSink.setClickHouseFormat(ClickHouseFormat.CSV);

        // 最后，将 DataStream 连接到 sink。
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        Path csvFilePath = new Path(fileFullName);
        FileSource<String> csvSource = FileSource
                .forRecordStreamFormat(new TextLineInputFormat(), csvFilePath)
                .build();

        env.fromSource(
                csvSource,
                WatermarkStrategy.noWatermarks(),
                "GzipCsvSource"
        ).sinkTo(csvSink);
    }
    ```
  </Tab>
</Tabs>

更多示例和代码片段可在我们的测试中找到：

* [flink-connector-clickhouse-1.17](https://github.com/ClickHouse/flink-connector-clickhouse/tree/main/flink-connector-clickhouse-1.17/src/test/java/org/apache/flink/connector/clickhouse/sink)
* [flink-connector-clickhouse-2.0.0](https://github.com/ClickHouse/flink-connector-clickhouse/tree/main/flink-connector-clickhouse-2.0.0/src/test/java/org/apache/flink/connector/clickhouse/sink)

<div id="datastream-quick-start">
  ### 快速入门示例
</div>

我们提供了基于 Maven 的示例，方便您快速开始使用 ClickHouse Sink：

* [Flink 1.17+](https://github.com/ClickHouse/flink-connector-clickhouse/tree/main/examples/maven/flink-v1.7/covid)
* [Flink 2.0.0+](https://github.com/ClickHouse/flink-connector-clickhouse/tree/main/examples/maven/flink-v2/covid)

如需更详细的说明，请参阅 [示例指南](https://github.com/ClickHouse/flink-connector-clickhouse/blob/main/examples/README.md)

<div id="datastream-api-connection-options">
  ### DataStream API 连接选项
</div>

<div id="client-options">
  #### ClickHouse 客户端选项
</div>

| Parameters                  | Description                                                                                   | Default Value | Required |
| --------------------------- | --------------------------------------------------------------------------------------------- | ------------- | -------- |
| `url`                       | 完整的 ClickHouse URL                                                                            | 不适用           | 是        |
| `username`                  | ClickHouse 数据库用户名                                                                             | 不适用           | 是        |
| `password`                  | ClickHouse 数据库密码                                                                              | 不适用           | 是        |
| `database`                  | ClickHouse 数据库名称                                                                              | 不适用           | 是        |
| `table`                     | ClickHouse 表名                                                                                 | 不适用           | 是        |
| `options`                   | Java 客户端配置选项的映射                                                                               | 空映射           | 否        |
| `serverSettings`            | ClickHouse 服务端会话设置的映射                                                                         | 空映射           | 否        |
| `enableJsonSupportAsString` | 用于让 ClickHouse 服务端期望 [JSON 数据类型](/zh/reference/data-types/newjson) 采用 JSON 格式 `String` 的服务端设置 | true          | 否        |

`options` 和 `serverSettings` 应以 `Map<String, String>` 的形式传递给客户端。如果其中任一项为空映射，则分别使用客户端或服务端的默认值。

<Note>
  所有可用的 Java 客户端选项均列在 [ClientConfigProperties.java](https://github.com/ClickHouse/clickhouse-java/blob/main/client-v2/src/main/java/com/clickhouse/client/api/ClientConfigProperties.java) 和[此文档页面](/zh/integrations/language-clients/java/client#configuration)中。

  所有可用的服务端会话设置均列在[此文档页面](/zh/reference/settings/session-settings)中。
</Note>

例如：

<Tabs>
  <Tab title="Java">
    ```java theme={null}
    Map<String, String> javaClientOptions = Map.of(
        ClientConfigProperties.CA_CERTIFICATE.getKey(), "<my_CA_cert>",
        ClientConfigProperties.SSL_CERTIFICATE.getKey(), "<my_SSL_cert>",
        ClientConfigProperties.CLIENT_NETWORK_BUFFER_SIZE.getKey(), "30000",
        ClientConfigProperties.HTTP_MAX_OPEN_CONNECTIONS.getKey(), "5"
    );

    Map<String, String> serverSettings = Map.of(
        "insert_deduplicate", "1"
    );

    ClickHouseClientConfig clickHouseClientConfig = new ClickHouseClientConfig(
        url,
        username,
        password,
        database,
        tableName,
        javaClientOptions,
        serverSettings,
        false // enableJsonSupportAsString
    );
    ```
  </Tab>
</Tabs>

<div id="sink-options">
  #### Sink 选项
</div>

以下选项直接来自 Flink 的 `AsyncSinkBase`：

| 参数                     | 描述                                        | 默认值 | 必填 |
| ---------------------- | ----------------------------------------- | --- | -- |
| `maxBatchSize`         | 单个批次中可插入的最大记录数                            | N/A | 是  |
| `maxInFlightRequests`  | 在 sink 开始施加背压之前，允许的最大进行中请求数               | N/A | 是  |
| `maxBufferedRequests`  | 在 sink 开始施加背压之前，可在其中缓冲的最大记录数              | N/A | 是  |
| `maxBatchSizeInBytes`  | 一个批次允许达到的最大大小 (以字节为单位) 。所有发送的批次都将小于或等于该大小 | N/A | 是  |
| `maxTimeInBufferMS`    | 记录在被刷新前可在 sink 中停留的最长时间                   | N/A | 是  |
| `maxRecordSizeInBytes` | sink 可接受的最大记录大小，超过该大小的记录会被自动拒绝            | N/A | 是  |

<div id="supported-data-types">
  ## 支持的数据类型
</div>

下表快速列出了将数据从 Flink 插入 ClickHouse 时的数据类型转换对应关系。

<div id="inserting-data-from-flink-into-clickhouse">
  ### 将 Flink 中的数据插入 ClickHouse
</div>

[//]: # "TODO: 添加表 API 支持后，增加一个“Flink SQL 类型”列 "

| Java 类型             | ClickHouse 类型     | 是否支持 | 序列化方法                         |
| ------------------- | ----------------- | ---- | ----------------------------- |
| `byte`/`Byte`       | `Int8`            | ✅    | `DataWriter.writeInt8`        |
| `short`/`Short`     | `Int16`           | ✅    | `DataWriter.writeInt16`       |
| `int`/`Integer`     | `Int32`           | ✅    | `DataWriter.writeInt32`       |
| `long`/`Long`       | `Int64`           | ✅    | `DataWriter.writeInt64`       |
| `BigInteger`        | `Int128`          | ✅    | `DataWriter.writeInt128`      |
| `BigInteger`        | `Int256`          | ✅    | `DataWriter.writeInt256`      |
| `short`/`Short`     | `UInt8`           | ✅    | `DataWriter.writeUInt8`       |
| `int`/`Integer`     | `UInt8`           | ✅    | `DataWriter.writeUInt8 `      |
| `int`/`Integer`     | `UInt16`          | ✅    | `DataWriter.writeUInt16`      |
| `long`/`Long`       | `UInt32`          | ✅    | `DataWriter.writeUInt32`      |
| `long`/`Long`       | `UInt64`          | ✅    | `DataWriter.writeUInt64`      |
| `BigInteger`        | `UInt64`          | ✅    | `DataWriter.writeUInt64`      |
| `BigInteger`        | `UInt128`         | ✅    | `DataWriter.writeUInt128`     |
| `BigInteger`        | `UInt256`         | ✅    | `DataWriter.writeUInt256`     |
| `BigDecimal`        | `Decimal`         | ✅    | `DataWriter.writeDecimal`     |
| `BigDecimal`        | `Decimal32`       | ✅    | `DataWriter.writeDecimal`     |
| `BigDecimal`        | `Decimal64`       | ✅    | `DataWriter.writeDecimal`     |
| `BigDecimal`        | `Decimal128`      | ✅    | `DataWriter.writeDecimal`     |
| `BigDecimal`        | `Decimal256`      | ✅    | `DataWriter.writeDecimal`     |
| `float`/`Float`     | `Float`           | ✅    | `DataWriter.writeFloat32`     |
| `double`/`Double`   | `Double`          | ✅    | `DataWriter.writeFloat64`     |
| `boolean`/`Boolean` | `Boolean`         | ✅    | `DataWriter.writeBoolean`     |
| `String`            | `String`          | ✅    | `DataWriter.writeString`      |
| `String`            | `FixedString`     | ✅    | `DataWriter.writeFixedString` |
| `LocalDate`         | `Date`            | ✅    | `DataWriter.writeDate`        |
| `LocalDate`         | `Date32`          | ✅    | `DataWriter.writeDate32`      |
| `LocalDateTime`     | `DateTime`        | ✅    | `DataWriter.writeDateTime`    |
| `ZonedDateTime`     | `DateTime`        | ✅    | `DataWriter.writeDateTime`    |
| `LocalDateTime`     | `DateTime64`      | ✅    | `DataWriter.writeDateTime64`  |
| `ZonedDateTime`     | `DateTime64`      | ✅    | `DataWriter.writeDateTime64`  |
| `int`/`Integer`     | `Time`            | ❌    | N/A                           |
| `long`/`Long`       | `Time64`          | ❌    | N/A                           |
| `byte`/`Byte`       | `Enum8`           | ✅    | `DataWriter.writeInt8`        |
| `int`/`Integer`     | `Enum16`          | ✅    | `DataWriter.writeInt16`       |
| `java.util.UUID`    | `UUID`            | ✅    | `DataWriter.writeIntUUID`     |
| `String`            | `JSON`            | ✅    | `DataWriter.writeJSON`        |
| `Array<Type>`       | `Array<Type>`     | ✅    | `DataWriter.writeArray`       |
| `Map<K,V>`          | `Map<K,V>`        | ✅    | `DataWriter.writeMap`         |
| `Tuple<Type,..>`    | `Tuple<T1,T2,..>` | ✅    | `DataWriter.writeTuple`       |
| `Object`            | `Variant`         | ❌    | N/A                           |

注意：

* 执行日期操作时，必须提供 `ZoneId`。
* 执行 decimal 操作时，必须提供[精度和标度](/zh/reference/data-types/decimal#decimal-value-ranges)。
* 要让 ClickHouse 将 Java String 解析为 JSON，需要在 `ClickHouseClientConfig` 中启用 `enableJsonSupportAsString`。
* 该 连接器 需要一个 `ElementConvertor`，用于将输入 DataStream 中的元素映射为 ClickHouse 载荷。为此，连接器 提供了 `ClickHouseConvertor` 和 `POJOConvertor`，你可以结合上述 `DataWriter` 序列化方法使用它们来实现这种映射。

<div id="supported-input-formats">
  ## 支持的输入格式
</div>

你可以在[此文档页面](/zh/reference/formats/index#formats-overview)和 [ClickHouseFormat.java](https://github.com/ClickHouse/clickhouse-java/blob/main/clickhouse-data/src/main/java/com/clickhouse/data/ClickHouseFormat.java) 中查看可用的 ClickHouse 输入格式列表。

要指定连接器用于将 `DataStream` 序列化为 ClickHouse 载荷的格式，请使用 `setClickHouseFormat` 函数。例如：

```java theme={null}
ClickHouseAsyncSink<String> csvSink = new ClickHouseAsyncSink<>(
        convertorString,
        MAX_BATCH_SIZE,
        MAX_IN_FLIGHT_REQUESTS,
        MAX_BUFFERED_REQUESTS,
        MAX_BATCH_SIZE_IN_BYTES,
        MAX_TIME_IN_BUFFER_MS,
        MAX_RECORD_SIZE_IN_BYTES,
        clickHouseClientConfig
);
csvSink.setClickHouseFormat(ClickHouseFormat.CSV);
```

<Note>
  默认情况下，如果在 `ClickHouseClientConfig` 中显式将 `setSupportDefault` 设置为 true 或 false，连接器将分别使用 [RowBinaryWithDefaults](/zh/reference/formats/RowBinary/RowBinaryWithDefaults) 或 [RowBinary](/zh/reference/formats/RowBinary/RowBinary)。
</Note>

<div id="metrics">
  ## 指标
</div>

该连接器在 Flink 现有指标的基础上，还额外暴露了以下指标：

| 指标                                      | 描述                                                                                                                                     | 类型  | 状态 |
| --------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------- | --- | -- |
| `numBytesSend`                          | 请求载荷中发送到 ClickHouse 的总字节数。*注意：该指标衡量的是通过网络发送的序列化数据大小，可能与 ClickHouse 的 `system.query_log` 中的 `written_bytes` 不同；后者反映的是数据经过处理后实际写入存储的字节数* | 计数器 | ✅  |
| `numRecordSend`                         | 发送到 ClickHouse 的记录总数                                                                                                                   | 计数器 | ✅  |
| `numRequestSubmitted`                   | 已发送的请求总数 (即实际执行的 flush 次数)                                                                                                             | 计数器 | ✅  |
| `numOfDroppedBatches`                   | 因不可重试失败而丢弃的批次总数                                                                                                                        | 计数器 | ✅  |
| `numOfDroppedRecords`                   | 因不可重试失败而丢弃的记录总数                                                                                                                        | 计数器 | ✅  |
| `totalBatchRetries`                     | 因可重试失败而进行的批次重试总数                                                                                                                       | 计数器 | ✅  |
| `writeLatencyHistogram`                 | 成功写入延迟分布的直方图 (毫秒)                                                                                                                      | 直方图 | ✅  |
| `writeFailureLatencyHistogram`          | 写入失败延迟分布的直方图 (毫秒)                                                                                                                      | 直方图 | ✅  |
| `triggeredByMaxBatchSizeCounter`        | 因达到 `maxBatchSize` 而触发的 flush 总数                                                                                                       | 计数器 | ✅  |
| `triggeredByMaxBatchSizeInBytesCounter` | 因达到 `maxBatchSizeInBytes` 而触发的 flush 总数                                                                                                | 计数器 | ✅  |
| `triggeredByMaxTimeInBufferMSCounter`   | 因达到 `maxTimeInBufferMS` 而触发的 flush 总数                                                                                                  | 计数器 | ✅  |
| `actualRecordsPerBatch`                 | 实际批次大小分布的直方图                                                                                                                           | 直方图 | ✅  |
| `actualBytesPerBatch`                   | 每个批次实际字节数分布的直方图                                                                                                                        | 直方图 | ✅  |

[//]: # "| actualTimeInBuffer           | flush 前在 buffer 中实际停留时间分布的直方图 | 直方图 | ❌      |"

<div id="limitations">
  ## 局限性
</div>

* 该 sink 当前提供至少一次交付保证。实现精确一次语义的相关工作可在[此处](https://github.com/ClickHouse/flink-connector-clickhouse/issues/106)跟踪。
* 该 sink 目前尚不支持用于缓冲无法处理记录的死信队列 (DLQ) 。在此之前，连接器 会尝试重新 insert 失败的记录；如果仍然失败，则会将其丢弃。此功能可在[此处](https://github.com/ClickHouse/flink-connector-clickhouse/issues/105)跟踪。
* 该 sink 目前尚不支持通过 Flink 的 Table API 或 Flink SQL 创建。此功能可在[此处](https://github.com/ClickHouse/flink-connector-clickhouse/issues/42)跟踪。

<div id="compatibility-and-security">
  ## ClickHouse 版本兼容性与安全
</div>

* 该连接器通过每日 CI 工作流针对一系列较新的 ClickHouse 版本进行测试，包括 latest 和 head。随着新的 ClickHouse 发行版进入活跃状态，测试版本也会定期更新。有关该连接器每日测试的版本，请参见[此处](https://github.com/ClickHouse/flink-connector-clickhouse/blob/main/.github/workflows/tests-nightly.yaml#L15)。
* 有关已知安全漏洞以及如何报告漏洞，请参见 [ClickHouse 安全策略](https://github.com/ClickHouse/ClickHouse/blob/master/SECURITY.md#security-change-log-and-support)。
* 我们建议持续升级该连接器，以免错过安全修复和新改进。
* 如果你在迁移过程中遇到问题，请创建 GitHub [issue](https://github.com/ClickHouse/flink-connector-clickhouse/issues)，我们会回复！

<div id="advanced-and-recommended-usage">
  ## 高级和推荐用法
</div>

* 为获得最佳性能，请确保你的 DataStream 元素类型**不是** Generic 类型——请参阅[这篇关于 Flink 类型区分的说明](https://nightlies.apache.org/flink/flink-docs-release-2.2/docs/dev/datastream/fault-tolerance/serialization/types_serialization/#flinks-typeinformation-class)。非 Generic 元素可避免 Kryo 引入的序列化开销，并提升写入 ClickHouse 的吞吐量。
* 我们建议将 `maxBatchSize` 设置为至少 1000，理想范围是 10,000 到 100,000。更多信息请参阅[这篇关于批量 insert 的指南](/zh/concepts/features/operations/insert/bulkinserts)。
* 如需在 ClickHouse 中执行 OLTP 风格的去重或 upsert，请参阅[此文档页面](/zh/concepts/features/operations/insert/deduplication#options-for-deduplication)。*注意：不要将其与发生在 retries 时的批次去重混淆。*

<div id="troubleshooting">
  ## 故障排查
</div>

<div id="cannot_read_all_data">
  ### CANNOT\_READ\_ALL\_DATA
</div>

可能会发生以下错误：

```text theme={null}
com.clickhouse.client.api.ServerException: Code: 33. DB::Exception: Cannot read all data. Bytes read: 9205. Bytes expected: 1100022.: (at row 9) : While executing BinaryRowInputFormat. (CANNOT_READ_ALL_DATA)
```

**原因**：最常见的情况是，CANNOT\_READ\_ALL\_DATA 错误表示你的 ClickHouse 表 schema 与 Flink 记录 schema 不一致。这通常发生在其中一方以不向后兼容的方式发生变更时。

**解决方案**：更新 ClickHouse 表 schema 或 连接器 输入数据类型中的 schema (或同时更新两者) ，使其相互兼容。如有需要，请参阅[type mapping](#inserting-data-from-flink-into-clickhouse)，了解如何将 Java 类型映射为 ClickHouse 类型。*注意：如果仍有记录在传输过程中，则在重启 连接器 时需要重置 Flink 状态。*

<div id="low_throughput">
  ### 吞吐量低
</div>

将数据写入 ClickHouse 时，你可能会发现连接器的吞吐量不会随着作业并行度 (Flink task 数量) 提升而线性扩展。

**原因**：ClickHouse 的后台[分片合并过程](/zh/concepts/core-concepts/merges)可能会拖慢 insert。在配置的批次大小过小、连接器 flush 过于频繁，或两者同时存在时，都可能出现这种情况。

**解决方案**：监控 `numRequestSubmitted` 和 `actualRecordsPerBatch` 指标，以帮助判断如何调整批次大小 (`maxBatchSize`) 以及 flush 频率。另请参阅[高级和推荐用法](#advanced-and-recommended-usage)中的批次大小建议。

[//]: # "TODO: 一旦 https://github.com/ClickHouse/flink-connector-clickhouse/issues/121 关闭，就取消注释此部分"

[//]: # "### 我在 ClickHouse 表中看到了重复的行批次 {#duplicate_batches}"

[//]: #

[//]: # "**原因**：如果 Flink 某个批次中的一条或多条记录因可重试故障而无法插入 ClickHouse，连接器会重试**整个批次**。如果未启用 [insert deduplication](https://clickhouse.com/docs/guides/developer/deduplicating-inserts-on-retries#query-level-insert-deduplication)，就可能导致重复记录写入 ClickHouse 表。否则，也有可能是去重窗口或窗口耗时过短，导致这些块在连接器重试之前就已过期。"

[//]: #

[//]: # "**解决方案**:"

[//]: # "- 如果你的表使用的是 `Replicated*MergeTree` table engine："

[//]: # "  1. 确保 server session setting `insert_deduplicate=1` (如有需要，可参见上面的[示例](#client-options)了解如何设置)。请注意，复制表默认启用 `insert_deduplicate`。"

[//]: # "  2. 如有必要，增大 `MergeTree` table settings [`replicated_deduplication_window`](https://clickhouse.com/docs/operations/settings/merge-tree-settings#replicated_deduplication_window) 或 [`replicated_deduplication_window_seconds`](https://clickhouse.com/docs/operations/settings/merge-tree-settings#replicated_deduplication_window_seconds)。"

[//]: # "- 如果你的表使用的是非复制表 `*MergeTree` table engine，请增大 `MergeTree` table setting [`non_replicated_deduplication_window`](https://clickhouse.com/docs/operations/settings/merge-tree-settings#non_replicated_deduplication_window)。"

[//]: #

[//]: # "_注 1：此解决方案依赖于 [synchronous inserts](https://clickhouse.com/docs/best-practices/selecting-an-insert-strategy#synchronous-inserts-by-default)，这也是 Flink 连接器推荐使用的方式。请确保 server session setting `async_insert=0`。_"

[//]: #

[//]: # "_注 2：`(non_)replicated_deduplication_window` 的值过大可能会拖慢 insert，因为需要比较更多 entries。_"

<div id="missing_rows">
  ### 我的 ClickHouse 表中有行缺失
</div>

**原因**：这些批次被丢弃了，原因可能是发生了不可重试的故障，或者在配置的重试次数内仍无法插入 (可通过 `ClickHouseClientConfig.setNumberOfRetries()` 设置) 。*注意：默认情况下，连接器会在丢弃某个批次之前最多尝试重新插入 3 次。*

**解决方案**：检查 TaskManager 日志和/或堆栈跟踪，以定位根本原因。

<div id="contributing-and-support">
  ## 贡献与支持
</div>

如果您想为该项目贡献力量或报告任何问题，欢迎向我们反馈！
请访问我们的 [GitHub 仓库](https://github.com/ClickHouse/flink-connector-clickhouse) 提交 issue、提出
改进建议，或提交拉取请求。

欢迎贡献！开始之前，请先查阅仓库中的[贡献指南](https://github.com/ClickHouse/flink-connector-clickhouse/blob/main/CONTRIBUTING.md)。
感谢您帮助改进 ClickHouse Flink 连接器！
