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 ステートメントを一切実行しないため、ターゲットテーブルは挿入前にあらかじめ存在している必要があります。
ClickHouseIOクラスのドキュメント。- サンプルを含む
Githubリポジトリ clickhouse-beam-connector。