Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Apache Beam と ClickHouse のインテグレーション

ClickHouse対応

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

インテグレーションパッケージ

Apache Beam と ClickHouse をインテグレーションするために必要なパッケージは、Apache Beam I/O Connectors で保守・開発されています。これは、多くの一般的なデータストレージシステムやデータベース向けインテグレーションをまとめたバンドルです。 org.apache.beam.sdk.io.clickhouse.ClickHouseIO の実装は、Apache Beam repo 内にあります。

Apache Beam ClickHouse パッケージの設定

パッケージのインストール

以下の依存関係をパッケージ管理フレームワークに追加します。

<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-clickhouse</artifactId>
    <version>${beam.version}</version>
</dependency>

アーティファクトは公式 Maven リポジトリで確認できます。

コード例

次の例では、input.csv という名前のCSVファイルを PCollection として読み込み、定義したスキーマを使って Row オブジェクトに変換し、ClickHouseIO を使用してローカルの ClickHouse インスタンスに挿入します。


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

サポートされているデータ型

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 で、
org.apache.beam.sdk.schemas.logicaltypes にあります
Schema.TypeName#DECIMAL
Schema.TypeName#MAP

ClickHouseIO.Write パラメーター

以下のセッター関数を使用して、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 のターゲットテーブルのスキーマ。

制限事項

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

  • 現時点でサポートされているのは Sink 操作のみで、コネクタは Source 操作には対応していません。
  • ClickHouse は、ReplicatedMergeTree または ReplicatedMergeTree を基盤とする Distributed テーブルへの挿入時に重複排除を実行します。レプリケーションがない場合、通常の MergeTree への挿入では、挿入が失敗したあと再試行で成功すると、重複が発生する可能性があります。ただし、各ブロックはアトミックに挿入され、ブロックサイズは ClickHouseIO.Write.withMaxInsertBlockSize(long) を使用して設定できます。重複排除は、挿入されたブロックのチェックサムを用いて行われます。重複排除の詳細については、Deduplication および Deduplicate insertion config を参照してください。
  • コネクタは DDL ステートメントを一切実行しないため、ターゲットテーブルは挿入前にあらかじめ存在している必要があります。
Navigation