> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-detect-table-modification.mintlify.site/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>;
};

**Apache Beam** は、開発者がバッチ処理とストリーム (継続的) データ処理パイプラインの両方を定義・実行できる、オープンソースの統一的なプログラミングモデルです。Apache Beam の柔軟性は、ETL (Extract, Transform, Load) 処理から複雑なイベント処理、リアルタイム分析まで、幅広いデータ処理シナリオに対応できる点にあります。
このインテグレーションでは、データ挿入の基盤レイヤーとして、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 バージョン**

  `ClickHouseIO` コネクタは、Apache Beam バージョン `2.59.0` 以降で使用することを推奨します。
  それ以前のバージョンでは、コネクタの機能が完全にサポートされない可能性があります。
</Warning>

アーティファクトは[公式 Maven リポジトリ](https://mvnrepository.com/artifact/org.apache.beam/beam-sdks-java-io-clickhouse)で確認できます。

<div id="code-example">
  ### コード例
</div>

次の例では、`input.csv` という名前のCSVファイルを `PCollection` として読み込み、定義したスキーマを使って `Row` オブジェクトに変換し、`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` は固定長の<br />バイト配列を表す `LogicalType` で、<br />`org.apache.beam.sdk.schemas.logicaltypes` にあります |
|                                    | `Schema.TypeName#DECIMAL`  | ❌      |                                                                                                         |
|                                    | `Schema.TypeName#MAP`      | ❌      |                                                                                                         |

<div id="clickhouseiowrite-parameters">
  ## ClickHouseIO.Write パラメーター
</div>

以下のセッター関数を使用して、`ClickHouseIO.Write` の設定を調整できます。

| パラメーター設定関数                  | 引数の型                        | デフォルト値                        | 説明                                 |
| --------------------------- | --------------------------- | ----------------------------- | ---------------------------------- |
| `withMaxInsertBlockSize`    | `(long maxInsertBlockSize)` | `1000000`                     | 挿入する行ブロックの最大サイズ。                   |
| `withMaxRetries`            | `(int maxRetries)`          | `5`                           | 失敗した insert の最大再試行回数。              |
| `withMaxCumulativeBackoff`  | `(Duration maxBackoff)`     | `Duration.standardDays(1000)` | 再試行における累積バックオフ時間の最大値。              |
| `withInitialBackoff`        | `(Duration initialBackoff)` | `Duration.standardSeconds(5)` | 最初の再試行前の初期バックオフ時間。                 |
| `withInsertDistributedSync` | `(Boolean sync)`            | `true`                        | true の場合、分散テーブルへの insert 操作を同期します。 |
| `withInsertQuorum`          | `(Long quorum)`             | `null`                        | insert 操作の確認に必要なレプリカ数。             |
| `withInsertDeduplicate`     | `(Boolean deduplicate)`     | `true`                        | true の場合、insert 操作の重複排除が有効になります。   |
| `withTableSchema`           | `(TableSchema schema)`      | `null`                        | ClickHouse のターゲットテーブルのスキーマ。        |

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

コネクタを使用する際は、以下の制限事項に注意してください。

* 現時点でサポートされているのは Sink 操作のみで、コネクタは Source 操作には対応していません。
* ClickHouse は、`ReplicatedMergeTree` または `ReplicatedMergeTree` を基盤とする `Distributed` テーブルへの挿入時に重複排除を実行します。レプリケーションがない場合、通常の MergeTree への挿入では、挿入が失敗したあと再試行で成功すると、重複が発生する可能性があります。ただし、各ブロックはアトミックに挿入され、ブロックサイズは `ClickHouseIO.Write.withMaxInsertBlockSize(long)` を使用して設定できます。重複排除は、挿入されたブロックのチェックサムを用いて行われます。重複排除の詳細については、[Deduplication](/ja/concepts/features/operations/insert/deduplication) および [Deduplicate insertion config](/ja/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)。
