> ## 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.

> 您可以使用 Apache Beam 将数据摄取到 ClickHouse

# 将 Apache Beam 与 ClickHouse 集成

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>;
};

<ClickHouseSupportedBadge />

**Apache Beam** 是一种开源的统一编程模型，使开发者能够定义并执行批次和 stream (连续) 数据处理管道。Apache Beam 的灵活性在于，它支持广泛的数据处理场景，从 ETL (提取、转换、加载) 操作到复杂事件处理和实时分析。
此集成使用 ClickHouse 官方的 [JDBC 连接器](https://github.com/ClickHouse/clickhouse-java) 作为底层插入机制。

<div id="integration-package">
  ## 集成包
</div>

用于集成 Apache Beam 和 ClickHouse 的集成包由 [Apache Beam I/O Connectors](https://beam.apache.org/documentation/io/connectors/) 维护和开发；这是一个包含多种流行数据存储系统和数据库的集成组件集合。
`org.apache.beam.sdk.io.clickhouse.ClickHouseIO` 的实现位于 [Apache Beam repo](https://github.com/apache/beam/tree/0bf43078130d7a258a0f1638a921d6d5287ca01e/sdks/java/io/clickhouse/src/main/java/org/apache/beam/sdk/io/clickhouse) 中。

<div id="setup-of-the-apache-beam-clickhouse-package">
  ## Apache Beam ClickHouse 软件包设置
</div>

<div id="package-installation">
  ### 软件包安装
</div>

将以下依赖项添加到包管理框架中：

```xml theme={null}
<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-clickhouse</artifactId>
    <version>${beam.version}</version>
</dependency>
```

<Warning>
  **推荐的 Beam 版本**

  建议从 Apache Beam `2.59.0` 版本起使用 `ClickHouseIO` 连接器。
  更早的版本可能无法完全支持该连接器的功能。
</Warning>

这些制品可在[官方 Maven 仓库](https://mvnrepository.com/artifact/org.apache.beam/beam-sdks-java-io-clickhouse)中获取。

<div id="code-example">
  ### 代码示例
</div>

以下示例将名为 `input.csv` 的 CSV 文件读取为 `PCollection`，再将其转换为 Row 对象 (使用已定义的 schema) ，并通过 `ClickHouseIO` 将其插入本地 ClickHouse 实例中：

```java theme={null}

package org.example;

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.io.clickhouse.ClickHouseIO;
import org.apache.beam.sdk.schemas.Schema;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
import org.joda.time.DateTime;

public class Main {

    public static void main(String[] args) {
        // 创建一个 Pipeline 对象。
        Pipeline p = Pipeline.create();

        Schema SCHEMA =
                Schema.builder()
                        .addField(Schema.Field.of("name", Schema.FieldType.STRING).withNullable(true))
                        .addField(Schema.Field.of("age", Schema.FieldType.INT16).withNullable(true))
                        .addField(Schema.Field.of("insertion_time", Schema.FieldType.DATETIME).withNullable(false))
                        .build();

        // 对管道应用转换操作。
        PCollection<String> lines = p.apply("ReadLines", TextIO.read().from("src/main/resources/input.csv"));

        PCollection<Row> rows = lines.apply("ConvertToRow", ParDo.of(new DoFn<String, Row>() {
            @ProcessElement
            public void processElement(@Element String line, OutputReceiver<Row> out) {

                String[] values = line.split(",");
                Row row = Row.withSchema(SCHEMA)
                        .addValues(values[0], Short.parseShort(values[1]), DateTime.now())
                        .build();
                out.output(row);
            }
        })).setRowSchema(SCHEMA);

        rows.apply("Write to ClickHouse",
                        ClickHouseIO.write("jdbc:clickhouse://localhost:8123/default?user=default&password=******", "test_table"));

        // 运行管道。
        p.run().waitUntilFinish();
    }
}

```

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

| ClickHouse                         | Apache Beam                | 是否支持 | 说明                                                                                                      |
| ---------------------------------- | -------------------------- | ---- | ------------------------------------------------------------------------------------------------------- |
| `TableSchema.TypeName.FLOAT32`     | `Schema.TypeName#FLOAT`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.FLOAT64`     | `Schema.TypeName#DOUBLE`   | ✅    |                                                                                                         |
| `TableSchema.TypeName.INT8`        | `Schema.TypeName#BYTE`     | ✅    |                                                                                                         |
| `TableSchema.TypeName.INT16`       | `Schema.TypeName#INT16`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.INT32`       | `Schema.TypeName#INT32`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.INT64`       | `Schema.TypeName#INT64`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.STRING`      | `Schema.TypeName#STRING`   | ✅    |                                                                                                         |
| `TableSchema.TypeName.UINT8`       | `Schema.TypeName#INT16`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.UINT16`      | `Schema.TypeName#INT32`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.UINT32`      | `Schema.TypeName#INT64`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.UINT64`      | `Schema.TypeName#INT64`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.DATE`        | `Schema.TypeName#DATETIME` | ✅    |                                                                                                         |
| `TableSchema.TypeName.DATETIME`    | `Schema.TypeName#DATETIME` | ✅    |                                                                                                         |
| `TableSchema.TypeName.ARRAY`       | `Schema.TypeName#ARRAY`    | ✅    |                                                                                                         |
| `TableSchema.TypeName.ENUM8`       | `Schema.TypeName#STRING`   | ✅    |                                                                                                         |
| `TableSchema.TypeName.ENUM16`      | `Schema.TypeName#STRING`   | ✅    |                                                                                                         |
| `TableSchema.TypeName.BOOL`        | `Schema.TypeName#BOOLEAN`  | ✅    |                                                                                                         |
| `TableSchema.TypeName.TUPLE`       | `Schema.TypeName#ROW`      | ✅    |                                                                                                         |
| `TableSchema.TypeName.FIXEDSTRING` | `FixedBytes`               | ✅    | `FixedBytes` 是一种 `LogicalType`，表示固定长度的 <br /> 字节数组，位于 <br /> `org.apache.beam.sdk.schemas.logicaltypes` |
|                                    | `Schema.TypeName#DECIMAL`  | ❌    |                                                                                                         |
|                                    | `Schema.TypeName#MAP`      | ❌    |                                                                                                         |

<div id="clickhouseiowrite-parameters">
  ## ClickHouseIO.Write 参数
</div>

你可以使用以下 setter 函数调整 `ClickHouseIO.Write` 配置：

| 参数设置函数                      | 参数类型                        | 默认值                           | 说明                       |
| --------------------------- | --------------------------- | ----------------------------- | ------------------------ |
| `withMaxInsertBlockSize`    | `(long maxInsertBlockSize)` | `1000000`                     | 单次插入的行块最大大小。             |
| `withMaxRetries`            | `(int maxRetries)`          | `5`                           | 插入失败时的最大重试次数。            |
| `withMaxCumulativeBackoff`  | `(Duration maxBackoff)`     | `Duration.standardDays(1000)` | 重试的最大累计退避耗时。             |
| `withInitialBackoff`        | `(Duration initialBackoff)` | `Duration.standardSeconds(5)` | 首次重试前的初始退避耗时。            |
| `withInsertDistributedSync` | `(Boolean sync)`            | `true`                        | 如果为 true，则会同步分布式表的插入操作。  |
| `withInsertQuorum`          | `(Long quorum)`             | `null`                        | 确认一次插入操作所需的副本数。          |
| `withInsertDeduplicate`     | `(Boolean deduplicate)`     | `true`                        | 如果为 true，则会对插入操作启用去重。    |
| `withTableSchema`           | `(TableSchema schema)`      | `null`                        | 目标 ClickHouse 表的 schema。 |

<div id="limitations">
  ## 限制
</div>

使用该连接器时，请注意以下限制：

* 截至目前，仅支持 Sink 操作，不支持 Source 操作。
* 向 `ReplicatedMergeTree` 或基于 `ReplicatedMergeTree` 构建的 `Distributed` 表插入数据时，ClickHouse 会执行去重。未启用复制时，如果插入失败后重试成功，向普通 MergeTree 表插入可能会产生重复数据。不过，每个块的插入都是原子的，并且可以使用 `ClickHouseIO.Write.withMaxInsertBlockSize(long)` 配置块大小。去重是通过对已插入块的校验和进行比对来实现的。有关去重的更多信息，请参阅 [去重](/docs/zh/concepts/features/operations/insert/deduplication) 和 [插入去重配置](/docs/zh/reference/settings/session-settings#insert_deduplicate)。
* 该连接器不会执行任何 DDL 语句；因此，目标表必须在插入前已存在。

<div id="related-content">
  ## 相关内容
</div>

* `ClickHouseIO` 类的[文档](https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/clickhouse/ClickHouseIO.html)。
* 示例 `Github` 仓库：[clickhouse-beam-connector](https://github.com/ClickHouse/clickhouse-beam-connector)。
