Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Spark コネクタ

ClickHouse対応

このコネクタは、高度なパーティション化や述語プッシュダウンなどの ClickHouse 固有の最適化を活用して、 クエリ性能とデータ処理を向上させます。 このコネクタは ClickHouse 公式の JDBC コネクタ をベースとしており、 独自のカタログを管理します。

Spark 3.0 より前は、Spark には組み込みのカタログという概念がなかったため、ユーザーは通常、 Hive Metastore や AWS Glue などの外部カタログシステムを利用していました。 これらの外部ソリューションでは、Spark からアクセスする前に、データソースのテーブルを手動で登録する必要がありました。 しかし、Spark 3.0 でカタログの概念が導入されたことで、Spark は カタログプラグインを登録するだけでテーブルを自動的に検出できるようになりました。

Spark のデフォルトカタログは spark_catalog であり、テーブルは {catalog name}.{database}.{table} で識別されます。新しい カタログ機能により、1 つの Spark アプリケーションで複数のカタログを追加して利用できるようになりました。

Catalog API と TableProvider API の選び方

ClickHouse Spark コネクタは、Catalog APITableProvider API (フォーマットベースのアクセス) という 2 つのアクセスパターンをサポートしています。その違いを理解することで、ユースケースに適したアプローチを選べます。

Catalog API と TableProvider API の比較

機能 Catalog API TableProvider API
設定 Spark の設定で一元管理 操作ごとにオプションで指定
テーブルの検出 カタログ経由で自動的に検出 テーブルを手動で指定
DDL 操作 完全対応 (CREATE、DROP、ALTER) 制限あり (自動テーブル作成のみ)
Spark SQL インテグレーション ネイティブ (clickhouse.database.table) フォーマットの指定が必要
ユースケース 設定を一元管理した長期的で安定した接続 アドホック、動的、または一時的なアクセス

要件

  • Java 8 または 17 (Spark 4.0 では Java 17 以降が必要)
  • Scala 2.12 または 2.13 (Spark 4.0 は Scala 2.13 のみをサポート)
  • Apache Spark 3.3、3.4、3.5、または 4.0

互換性マトリクス

バージョン 対応する Spark バージョン ClickHouse JDBC バージョン
main Spark 3.3, 3.4, 3.5, 4.0 0.9.4
0.10.0 Spark 3.3, 3.4, 3.5, 4.0 0.9.5
0.9.0 Spark 3.3, 3.4, 3.5, 4.0 0.9.4
0.8.1 Spark 3.3, 3.4, 3.5 0.6.3
0.7.3 Spark 3.3, 3.4 0.4.6
0.6.0 Spark 3.3 0.3.2-patch11
0.5.0 Spark 3.2, 3.3 0.3.2-patch11
0.4.0 Spark 3.2, 3.3 依存なし
0.3.0 Spark 3.2, 3.3 依存なし
0.2.1 Spark 3.2 依存なし
0.1.2 Spark 3.2 依存なし

インストールとセットアップ

ClickHouse を Spark と統合するには、プロジェクト構成に応じて複数のインストール方法を選択できます。 ClickHouse Spark コネクタは、プロジェクトのビルドファイル (Maven の pom.xml や SBT の build.sbt など) に依存関係として直接追加できます。 また、必要な JAR ファイルを $SPARK_HOME/jars/ フォルダーに配置するか、spark-submit コマンドで --jars フラグを使って Spark のオプションとして直接指定することもできます。 どちらの方法でも、ClickHouse Spark コネクタを Spark 環境で利用できるようになります。

依存関係として追加

<dependency>
  <groupId>com.clickhouse.spark</groupId>
  <artifactId>clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}</artifactId>
  <version>{{ stable_version }}</version>
</dependency>
<dependency>
  <groupId>com.clickhouse</groupId>
  <artifactId>clickhouse-jdbc</artifactId>
  <classifier>all</classifier>
  <version>{{ clickhouse_jdbc_version }}</version>
  <exclusions>
    <exclusion>
      <groupId>*</groupId>
      <artifactId>*</artifactId>
    </exclusion>
  </exclusions>
</dependency>

SNAPSHOT バージョンを使用するには、Maven で Sonatype の SNAPSHOT リリースを利用するための手順 に従ってください。

ライブラリをダウンロードする

バイナリJARの命名パターンは次のとおりです。

clickhouse-spark-runtime-${spark_binary_version}_${scala_binary_version}-${version}.jar

利用可能なリリース済み JAR ファイルはすべて Maven Central Repository で入手できます。 日次ビルドの SNAPSHOT JAR ファイルは、上記で設定した Sonatype snapshots repository を通じて利用できます。

カタログを登録する (必須)

ClickHouse のテーブルにアクセスするには、以下の設定で新しい Spark カタログを構成する必要があります。

プロパティ デフォルト値 必須
spark.sql.catalog.<catalog_name> com.clickhouse.spark.ClickHouseCatalog N/A はい
spark.sql.catalog.<catalog_name>.host <clickhouse_host> localhost いいえ
spark.sql.catalog.<catalog_name>.protocol http http いいえ
spark.sql.catalog.<catalog_name>.http_port <clickhouse_port> 8123 いいえ
spark.sql.catalog.<catalog_name>.user <clickhouse_username> default いいえ
spark.sql.catalog.<catalog_name>.password <clickhouse_password> (空文字列) いいえ
spark.sql.catalog.<catalog_name>.database <database> default いいえ
spark.<catalog_name>.write.format json arrow いいえ

これらの設定は、次のいずれかの方法で指定できます。

  • spark-defaults.conf を編集または作成する。
  • 設定を spark-submit コマンド (または spark-shell / spark-sql CLI コマンド) に渡す。
  • コンテキストの初期化時に設定を追加する。

TableProvider API の使用 (フォーマットベースのアクセス)

カタログベースの方法に加え、ClickHouse Spark コネクタは、TableProvider API を通じた フォーマットベースのアクセス方式 もサポートしています。

フォーマットベースの読み取りの例

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# format API を使用して ClickHouse から読み込む
df = spark.read \
    .format("clickhouse") \
    .option("host", "your-clickhouse-host") \
    .option("protocol", "https") \
    .option("http_port", "8443") \
    .option("database", "default") \
    .option("table", "your_table") \
    .option("user", "default") \
    .option("password", "your_password") \
    .option("ssl", "true") \
    .load()

df.show()

フォーマットベースの書き込み例

# format APIを使用してClickHouseに書き込む
df.write \
    .format("clickhouse") \
    .option("host", "your-clickhouse-host") \
    .option("protocol", "https") \
    .option("http_port", "8443") \
    .option("database", "default") \
    .option("table", "your_table") \
    .option("user", "default") \
    .option("password", "your_password") \
    .option("ssl", "true") \
    .mode("append") \
    .save()

TableProvider API の機能

TableProvider API には、便利な機能がいくつか用意されています。

自動テーブル作成

存在しないテーブルに書き込むと、コネクタが適切なスキーマでそのテーブルを自動的に作成します。コネクタには、適切なデフォルト設定が用意されています。

  • Engine: 指定しない場合は MergeTree() がデフォルトで使用されます。engine オプションを使って別のエンジンを指定できます (例: ReplacingMergeTree(), SummingMergeTree() など)
  • ORDER BY: 必須 - 新しいテーブルを作成する際は、order_by オプションを明示的に指定する必要があります。コネクタは、指定されたすべてのカラムがスキーマ内に存在することを検証します。
  • Nullable Key Support: ORDER BY に Nullable カラムが含まれる場合は、自動的に settings.allow_nullable_key=1 を追加します
# テーブルは明示的な ORDER BY を指定して自動的に作成されます(必須)
df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "new_table") \
    .option("order_by", "id") \
    .mode("append") \
    .save()

# カスタムエンジンを使用してテーブル作成オプションを指定
df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "new_table") \
    .option("order_by", "id, timestamp") \
    .option("engine", "ReplacingMergeTree()") \
    .option("settings.allow_nullable_key", "1") \
    .mode("append") \
    .save()

TableProvider の接続オプション

フォーマットベースの API を使用する場合は、次の接続オプションを利用できます。

接続オプション

オプション 説明 デフォルト値 必須
host ClickHouse server のホスト名 localhost はい
protocol 接続プロトコル (http または https) http いいえ
http_port HTTP/HTTPS のポート 8123 いいえ
database データベース名 default はい
table テーブル名 該当なし はい
user authentication 用のユーザー名 default いいえ
password authentication 用のパスワード (空文字列) いいえ
ssl SSL 接続を有効にする false いいえ
ssl_mode SSL モード (NONESTRICT など) STRICT いいえ
timezone 日付/時刻の操作に使用するタイムゾーン server いいえ

テーブル作成オプション

これらのオプションは、テーブルが存在せず、新たに作成する必要がある場合に使用します。

オプション 説明 デフォルト値 必須
order_by ORDER BY 句に使用するカラム。複数カラムを指定する場合はカンマ区切り N/A はい
engine ClickHouse のテーブルエンジン (例: MergeTree(), ReplacingMergeTree(), SummingMergeTree() など) MergeTree() いいえ
settings.allow_nullable_key ORDER BY で Nullable なキーを有効にします (ClickHouse Cloud 向け) 自動検出** いいえ
settings.<key> 任意の ClickHouse テーブル設定 N/A いいえ
cluster 分散テーブル用のクラスター名 N/A いいえ
clickhouse.column.<name>.variant_types Variant カラム用の ClickHouse 型のカンマ区切りリスト (例: String, Int64, Bool, JSON) 。型名では大文字と小文字が区別されます。カンマの後のスペースは省略可能です。 N/A いいえ
  • 新しいテーブルを作成する場合、order_by オプションは必須です。指定するすべてのカラムはスキーマ内に存在している必要があります。 ** ORDER BY に Nullable カラムが含まれており、明示的に指定されていない場合は、自動的に 1 に設定されます。

書き込みモード

Spark コネクタ (TableProvider API と Catalog API の両方) は、以下の Spark 書き込みモードをサポートしています。

  • append: 既存のテーブルにデータを追加
  • overwrite: テーブル内のすべてのデータを置き換える (テーブルを空にする)
# 上書きモード(最初にテーブルを空にする)
df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "my_table") \
    .mode("overwrite") \
    .save()

ClickHouse オプションの設定

Catalog API と TableProvider API はどちらも、ClickHouse 固有のオプション (コネクタのオプションではなく) を設定できます。これらのオプションは、テーブルの作成時やクエリの実行時に ClickHouse に渡されます。

ClickHouse オプションを使用すると、allow_nullable_keyindex_granularity、その他のテーブルレベルまたはクエリレベルの設定など、ClickHouse 固有の設定を構成できます。これらは、コネクタが ClickHouse に接続する方法を制御するコネクタ オプション (hostdatabasetable など) とは異なります。

TableProvider API を使用する

TableProvider API では、settings.<key> 形式のオプションを使用します。

df.write \
    .format("clickhouse") \
    .option("host", "your-host") \
    .option("database", "default") \
    .option("table", "my_table") \
    .option("order_by", "id") \
    .option("settings.allow_nullable_key", "1") \
    .option("settings.index_granularity", "8192") \
    .mode("append") \
    .save()

Catalog API を使用する

Catalog API を使用する場合は、Spark の設定で spark.sql.catalog.<catalog_name>.option.<key> の形式を使用します。

spark.sql.catalog.clickhouse.option.allow_nullable_key 1
spark.sql.catalog.clickhouse.option.index_granularity 8192

または、Spark SQL でテーブルを作成する際に設定することもできます:

CREATE TABLE clickhouse.default.my_table (
  id INT,
  name STRING
) USING ClickHouse
TBLPROPERTIES (
  engine = 'MergeTree()',
  order_by = 'id',
  'settings.allow_nullable_key' = '1',
  'settings.index_granularity' = '8192'
)

ClickHouse Cloud の設定

ClickHouse Cloud に接続する際は、SSL を有効にし、適切な SSL モードを設定してください。たとえば、次のようにします。

spark.sql.catalog.clickhouse.option.ssl        true
spark.sql.catalog.clickhouse.option.ssl_mode   NONE

データの読み取り

public static void main(String[] args) {
        // Spark セッションを作成
        SparkSession spark = SparkSession.builder()
                .appName("example")
                .master("local[*]")
                .config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
                .config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
                .config("spark.sql.catalog.clickhouse.protocol", "http")
                .config("spark.sql.catalog.clickhouse.http_port", "8123")
                .config("spark.sql.catalog.clickhouse.user", "default")
                .config("spark.sql.catalog.clickhouse.password", "123456")
                .config("spark.sql.catalog.clickhouse.database", "default")
                .config("spark.clickhouse.write.format", "json")
                .getOrCreate();

        Dataset<Row> df = spark.sql("select * from clickhouse.default.example_table");

        df.show();

        spark.stop();
    }

データの書き込み

 public static void main(String[] args) throws AnalysisException {

        // Spark セッションを作成
        SparkSession spark = SparkSession.builder()
                .appName("example")
                .master("local[*]")
                .config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
                .config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
                .config("spark.sql.catalog.clickhouse.protocol", "http")
                .config("spark.sql.catalog.clickhouse.http_port", "8123")
                .config("spark.sql.catalog.clickhouse.user", "default")
                .config("spark.sql.catalog.clickhouse.password", "123456")
                .config("spark.sql.catalog.clickhouse.database", "default")
                .config("spark.clickhouse.write.format", "json")
                .getOrCreate();

        // DataFrame のスキーマを定義
        StructType schema = new StructType(new StructField[]{
                DataTypes.createStructField("id", DataTypes.IntegerType, false),
                DataTypes.createStructField("name", DataTypes.StringType, false),
        });

        List<Row> data = Arrays.asList(
                RowFactory.create(1, "Alice"),
                RowFactory.create(2, "Bob")
        );

        // DataFrame を作成
        Dataset<Row> df = spark.createDataFrame(data, schema);

        df.writeTo("clickhouse.default.example_table").append();

        spark.stop();
    }

DDL 操作

Spark SQL を使用すると、ClickHouse インスタンスに対して DDL 操作を実行でき、すべての変更は即座に ClickHouse に永続化されます。 Spark SQL では、ClickHouse とまったく同じようにクエリを記述できるため、 たとえば CREATE TABLE や TRUNCATE などのコマンドも、変更を加えることなく直接実行できます。

USE clickhouse; 

CREATE TABLE test_db.tbl_sql (
  create_time TIMESTAMP NOT NULL,
  m           INT       NOT NULL COMMENT 'part key',
  id          BIGINT    NOT NULL COMMENT 'sort key',
  value       STRING
) USING ClickHouse
PARTITIONED BY (m)
TBLPROPERTIES (
  engine = 'MergeTree()',
  order_by = 'id',
  settings.index_granularity = 8192
);

上記の例は Spark SQL クエリを示しています。これらのクエリは、Java、Scala、 PySpark、またはシェルなど、任意の API を使ってアプリケーション内で実行できます。

VariantType を扱う

このコネクタは、半構造化データを扱うための Spark の VariantType をサポートしています。VariantType は ClickHouse の JSON 型および Variant 型にマッピングされるため、柔軟なスキーマを持つデータを効率的に保存およびクエリできます。

ClickHouse 型マッピング

ClickHouse 型 Spark 型 説明
JSON VariantType JSONオブジェクトのみを格納します ({ で始まる必要があります)
Variant(T1, T2, ...) VariantType プリミティブ、配列、JSON など、複数の型を格納します

VariantType データの読み込み

ClickHouse から読み込む際、JSON カラムと Variant カラムは自動的に Spark の VariantType にマッピングされます。

// JSONカラムをVariantTypeとして読み込む
val df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")

// Variantデータにアクセス
df.show()

// 確認のため、VariantをJSON文字列に変換
import org.apache.spark.sql.functions._
df.select(
  col("id"),
  to_json(col("data")).as("data_json")
).show()

VariantType データの書き込み

JSON または Variant カラム型を使用して、VariantType データを ClickHouse に書き込めます。

import org.apache.spark.sql.functions._

// JSON データを含む DataFrame を作成
val jsonData = Seq(
  (1, """{"name": "Alice", "age": 30}"""),
  (2, """{"name": "Bob", "age": 25}"""),
  (3, """{"name": "Charlie", "city": "NYC"}""")
).toDF("id", "json_string")

// JSON 文字列を VariantType にパース
val variantDF = jsonData.select(
  col("id"),
  parse_json(col("json_string")).as("data")
)

// JSON 型を使って ClickHouse に書き込む(JSON object のみ)
variantDF.writeTo("clickhouse.default.user_data").create()

// または、複数の型を持つ Variant を指定
spark.sql("""
  CREATE TABLE clickhouse.default.mixed_data (
    id INT,
    data VARIANT
  ) USING clickhouse
  TBLPROPERTIES (
    'clickhouse.column.data.variant_types' = 'String, Int64, Bool, JSON',
    'engine' = 'MergeTree()',
    'order_by' = 'id'
  )
""")

Spark SQLでのVariantType テーブルの作成

Spark SQL の DDL を使用して、VariantType テーブルを作成できます。

-- JSON 型(デフォルト)でテーブルを作成
CREATE TABLE clickhouse.default.json_table (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)
-- 複数の型をサポートする Variant 型のテーブルを作成する
CREATE TABLE clickhouse.default.flexible_data (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'clickhouse.column.data.variant_types' = 'String, Int64, Float64, Bool, Array(String), JSON',
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)

Variant 型の設定

VariantType カラムを含むテーブルを作成する際は、使用する ClickHouse の型を指定できます。

JSON 型 (デフォルト)

variant_types プロパティが指定されていない場合、このカラムはデフォルトで ClickHouse の JSON 型となり、JSON オブジェクトのみを受け付けます:

CREATE TABLE clickhouse.default.json_table (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)

これにより、次の ClickHouse クエリが生成されます。

CREATE TABLE json_table (id Int32, data JSON) ENGINE = MergeTree() ORDER BY id

複数の型をサポートする Variant 型

プリミティブ、Array、JSON オブジェクトをサポートするには、variant_types プロパティに型を指定します。

CREATE TABLE clickhouse.default.flexible_data (
  id INT,
  data VARIANT
) USING clickhouse
TBLPROPERTIES (
  'clickhouse.column.data.variant_types' = 'String, Int64, Float64, Bool, Array(String), JSON',
  'engine' = 'MergeTree()',
  'order_by' = 'id'
)

これにより、次の ClickHouse クエリが生成されます:

CREATE TABLE flexible_data (
  id Int32, 
  data Variant(String, Int64, Float64, Bool, Array(String), JSON)
) ENGINE = MergeTree() ORDER BY id

サポートされている Variant 型

Variant() では、以下の ClickHouse 型を使用できます。

  • プリミティブ: String, Int8, Int16, Int32, Int64, UInt8, UInt16, UInt32, UInt64, Float32, Float64, Bool
  • 配列: Array(T) (T は、ネストした配列を含む任意のサポート対象の型)
  • JSON: JSON オブジェクトを格納するための JSON

読み取りフォーマットの設定

デフォルトでは、JSON および Variant カラムは VariantType として読み取られます。これを上書きして、文字列として読み取ることもできます。

// JSON/Variant を VariantType ではなく文字列として読み取る
spark.conf.set("spark.clickhouse.read.jsonAs", "string")

val df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
// data カラムは JSON 文字列を含む StringType になる

書き込みフォーマットのサポート

VariantType の書き込みサポートは、フォーマットによって異なります。

フォーマット サポート 注記
JSON ✅ 完全 JSON 型と Variant 型の両方をサポートします。VariantType データには JSON を推奨します
Arrow ⚠️ 部分的 ClickHouse の JSON 型への書き込みをサポートします。ClickHouse の Variant 型はサポートしていません。完全対応は https://github.com/ClickHouse/ClickHouse/issues/92752 の解決待ちです

書き込みフォーマットを設定します。

spark.conf.set("spark.clickhouse.write.format", "json")  // Variant 型に推奨

ベストプラクティス

  1. JSON 専用のデータには JSON 型を使用する: JSON object だけを保存する場合は、デフォルトの JSON 型 (variant_types プロパティなし) を使用します
  2. 型を明示的に指定する: Variant() を使用する場合は、保存する予定のすべての型を明示的に列挙します
  3. 実験的機能を有効にする: ClickHouse で allow_experimental_json_type = 1 が有効になっていることを確認します
  4. 書き込みには JSON フォーマットを使用する: 互換性を高めるため、VariantType データの書き込みには JSON フォーマットを使用することを推奨します
  5. クエリパターンを考慮する: JSON/Variant 型は、効率的にフィルタリングするための ClickHouse の JSON パスクエリをサポートしています
  6. パフォーマンス向上のためのカラムヒント: ClickHouse で JSON フィールドを使用する場合、カラムヒントを追加するとクエリのパフォーマンスが向上します。現在、Spark 経由でのカラムヒントの追加はサポートされていません。この機能の進捗状況は GitHub issue #497 を参照してください。

例: ワークフロー全体

import org.apache.spark.sql.functions._

// ClickHouseで実験的なJSON型を有効にする
spark.sql("SET allow_experimental_json_type = 1")

// Variant カラムを持つテーブルを作成する
spark.sql("""
  CREATE TABLE clickhouse.default.events (
    event_id BIGINT,
    event_time TIMESTAMP,
    event_data VARIANT
  ) USING clickhouse
  TBLPROPERTIES (
    'clickhouse.column.event_data.variant_types' = 'String, Int64, Bool, JSON',
    'engine' = 'MergeTree()',
    'order_by' = 'event_time'
  )
""")

// 混合型のデータを準備する
val events = Seq(
  (1L, "2024-01-01 10:00:00", """{"action": "login", "user_id": 123}"""),
  (2L, "2024-01-01 10:05:00", """{"action": "purchase", "amount": 99.99}"""),
  (3L, "2024-01-01 10:10:00", """{"action": "logout", "duration": 600}""")
).toDF("event_id", "event_time", "json_data")

// VariantTypeに変換して書き込む
val variantEvents = events.select(
  col("event_id"),
  to_timestamp(col("event_time")).as("event_time"),
  parse_json(col("json_data")).as("event_data")
)

variantEvents.writeTo("clickhouse.default.events").append()

// 読み取りとクエリ
val result = spark.sql("""
  SELECT event_id, event_time, event_data
  FROM clickhouse.default.events
  WHERE event_time >= '2024-01-01'
  ORDER BY event_time
""")

result.show(false)

設定

以下は、コネクタで利用可能な調整可能な設定です。


設定キー デフォルト 説明 導入バージョン
spark.clickhouse.ignoreUnsupportedTransform true ClickHouse では、cityHash64(col_1, col_2) のような複雑な式を分片キーやパーティション値として使用できますが、現時点では Spark ではそれらをサポートしていません。true の場合はサポートされていない式を無視して警告をログに記録し、そうでない場合は例外を発生させて直ちに失敗します。警告: spark.clickhouse.write.distributed.convertLocal=true の場合、サポートされていない分片キーを無視するとデータが破損するおそれがあります。コネクタはこれを検証し、デフォルトではエラーをスローします。これを許可するには、spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding=true を明示的に設定してください。 0.4.0
spark.clickhouse.read.compression.codec lz4 読み取り時にデータを伸張するために使用するコーデックです。対応するコーデック: none, lz4. 0.5.0
spark.clickhouse.read.distributed.convertLocal true 分散テーブルの読み取り時には、分散テーブル自体ではなくローカルテーブルを読み取ります。true の場合、spark.clickhouse.read.distributed.useClusterNodes は無視されます。 0.1.0
spark.clickhouse.read.fixedStringAs binary ClickHouse の FixedString 型を、指定した Spark のデータ型として読み込みます。対応する型: binary、string 0.8.0
spark.clickhouse.read.format json 読み取り時のシリアライズフォーマット。対応フォーマット: json、binary 0.6.0
spark.clickhouse.read.runtimeFilter.enabled false 読み取り時のランタイムフィルタを有効にします。 0.8.0
spark.clickhouse.read.splitByPartitionId true true の場合、パーティション値ではなく、仮想カラム _partition_id を使って入力パーティションのフィルターを構築します。パーティション値に基づいて SQL 述語を組み立てる場合には、既知の問題があります。この機能を利用するには ClickHouse Server v21.6+ が必要です 0.4.0
spark.clickhouse.useNullableQuerySchema false true の場合、テーブル作成時に CREATE/REPLACE TABLE ... AS SELECT ... を実行すると、クエリスキーマ内のすべてのフィールドが Nullable としてマークされます。なお、この設定には SPARK-43390 (Spark 3.5 で利用可能) が必要です。このパッチがない場合は、常に true として動作します。 0.8.0
spark.clickhouse.write.batchSize 10000 ClickHouse への書き込み時の1バッチあたりのレコード数。 0.1.0
spark.clickhouse.write.compression.codec lz4 書き込み時のデータ圧縮に使用するコーデック。対応コーデック: none、lz4。 0.3.0
spark.clickhouse.write.distributed.convertLocal false Distributed テーブルに書き込む際、テーブル自体ではなくローカルテーブルに書き込みます。true の場合、spark.clickhouse.write.distributed.useClusterNodes は無視されます。これにより ClickHouse のネイティブルーティングがバイパスされるため、Spark がシャーディングキーを評価する必要があります。サポートされていないシャーディング式を使用する場合は、データ分散エラーが表面化しないまま発生するのを防ぐため、spark.clickhouse.ignoreUnsupportedTransformfalse に設定してください。 0.1.0
spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding false 分片キーがサポートされていない場合でも、convertLocal=true かつ ignoreUnsupportedTransform=true で分散テーブルへの書き込みを許可します。これは危険であり、分片が正しく行われないことでデータ破損を引き起こすおそれがあります。true に設定する場合、Spark はサポートされていない分片式を評価できないため、書き込み前にデータが適切にソートおよび分片されていることを必ず確認してください。true に設定するのは、リスクを理解し、データの分布を検証済みの場合に限ってください。デフォルトでは、この組み合わせは、気付かないうちにデータが破損するのを防ぐためにエラーになります。 0.10.0
spark.clickhouse.write.distributed.useClusterNodes true 分散テーブルへの書き込み時に、クラスター内のすべてのノードに書き込みます。 0.1.0
spark.clickhouse.write.format arrow 書き込み時のシリアライズフォーマット。対応フォーマット: json, arrow 0.4.0
spark.clickhouse.write.localSortByKey true true の場合、書き込み前にソートキーに基づいてローカルでソートします。 0.3.0
spark.clickhouse.write.localSortByPartition spark.clickhouse.write.repartitionByPartition の値 true の場合、書き込み前にパーティションごとにローカルソートを行います。未設定の場合は、spark.clickhouse.write.repartitionByPartition と同じ値になります。 0.3.0
spark.clickhouse.write.maxRetry 3 再試行可能なコードにより単一のバッチ書き込みが失敗した場合に、その書き込みを再試行する最大回数です。 0.1.0
spark.clickhouse.write.repartitionByPartition true 書き込み前に、ClickHouseテーブルのデータ分布に合わせて、ClickHouseのパーティションキーでデータを再パーティション化するかどうか。 0.3.0
spark.clickhouse.write.repartitionNum 0 書き込み前に、ClickHouseテーブルの分散に合わせてデータを再パーティション化する必要があります。再パーティション数を指定するには、この設定を使用します。値が1未満の場合は、要件なしを意味します。 0.1.0
spark.clickhouse.write.repartitionStrictly false true の場合、Spark は書き込み時にレコードをログソーステーブルへ渡す前に、必要な分散要件を満たすよう、受信したレコードを各パーティションに厳密に分散します。そうでない場合、Spark はクエリを高速化するために一部の最適化を適用できますが、その結果、分散要件が満たされなくなる可能性があります。なお、この設定には SPARK-37523 (Spark 3.4 で利用可能) が必要で、このパッチがない場合は常に true として動作します。 0.3.0
spark.clickhouse.write.retryInterval 10s 書き込み再試行の間隔 (秒) 。 0.1.0
spark.clickhouse.write.retryableErrorCodes 241 書き込みの失敗時に ClickHouse server から返される、再試行可能なエラーコード。 0.1.0

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

このセクションでは、Spark と ClickHouse 間のデータ型の対応関係を概説します。以下の表は、ClickHouse から Spark にデータを読み込む場合と、Spark から ClickHouse にデータを挿入する場合のデータ型変換をすばやく参照できる一覧です。

ClickHouse から Spark へのデータの読み込み

ClickHouse データ型 Spark データ型 サポート プリミティブ型か 注記
Nothing NullType はい
Bool BooleanType はい
UInt8, Int16 ShortType はい
Int8 ByteType はい
UInt16,Int32 IntegerType はい
UInt32,Int64, UInt64 LongType はい
Int128,UInt128, Int256, UInt256 DecimalType(38, 0) はい
Float32 FloatType はい
Float64 DoubleType はい
String, UUID, Enum8, Enum16, IPv4, IPv6 StringType はい
FixedString BinaryType, StringType はい 設定 READ_FIXED_STRING_AS によって制御されます
Decimal DecimalType はい 精度とスケールは Decimal128 まで
Decimal32 DecimalType(9, scale) はい
Decimal64 DecimalType(18, scale) はい
Decimal128 DecimalType(38, scale) はい
Date, Date32 DateType はい
DateTime, DateTime32, DateTime64 TimestampType はい
Array ArrayType いいえ Array の要素型も変換されます
Map MapType いいえ キーは StringType に限定されます
IntervalYear YearMonthIntervalType(Year) はい
IntervalMonth YearMonthIntervalType(Month) はい
IntervalDay, IntervalHour, IntervalMinute, IntervalSecond DayTimeIntervalType いいえ 該当するインターバル型が使用されます
JSON, Variant VariantType いいえ Spark 4.0+ および ClickHouse 25.3+ が必要です。spark.clickhouse.read.jsonAs=string を使うと StringType として読み取れます
Object
Nested
Tuple StructType いいえ 名前付き/名前なしの両方の Tuple をサポートします。名前付き Tuple は名前で struct フィールドにマッピングされ、名前なし Tuple では _1_2 などが使われます。ネストされた struct と NULL 許容フィールドもサポートします
Point
Polygon
MultiPolygon
Ring
IntervalQuarter
IntervalWeek
Decimal256
AggregateFunction
SimpleAggregateFunction

Spark から ClickHouse へのデータ挿入

Spark データ型 ClickHouse データ型 サポート プリミティブ型か 注記
BooleanType Bool はい バージョン 0.9.0 以降は、UInt8 ではなく Bool 型にマッピングされます
ByteType Int8 はい
ShortType Int16 はい
IntegerType Int32 はい
LongType Int64 はい
FloatType Float32 はい
DoubleType Float64 はい
StringType String はい
VarcharType String はい
CharType String はい
DecimalType Decimal(p, s) はい 精度と scale は Decimal128 まで対応
DateType Date はい
TimestampType DateTime はい
ArrayType (list, tuple, or array) Array いいえ Array の element type も変換されます
MapType Map いいえ キーは StringType のみに制限されます
StructType Tuple いいえ フィールド名を持つ named Tuple に変換されます。
VariantType JSON or Variant いいえ Spark 4.0+ および ClickHouse 25.3+ が必要です。デフォルトでは JSON 型が使用されます。複数の型を持つ Variant を指定するには、clickhouse.column.<name>.variant_types プロパティを使用します。
Object
Nested

コントリビューションとサポート

プロジェクトへのコントリビューションや問題の報告をご希望の場合は、ぜひご協力ください。 issue の起票、改善の提案、またはプルリクエストの送信は、GitHub リポジトリから行えます。 コントリビューションを歓迎します。作業を始める前に、リポジトリ内のコントリビューションガイドラインをご確認ください。 ClickHouse Spark コネクタの改善にご協力いただき、ありがとうございます!

Navigation