このコネクタは、高度なパーティション化や述語プッシュダウンなどの 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 API と TableProvider 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 リリースを利用するための手順 に従ってください。
dependencies {
implementation("com.clickhouse.spark:clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }}")
implementation("com.clickhouse:clickhouse-jdbc:{{ clickhouse_jdbc_version }}:all") { transitive = false }
}SNAPSHOT バージョンを使用するには、Gradle で Sonatype の SNAPSHOT リリースを利用するための手順 に従ってください。
libraryDependencies += "com.clickhouse" % "clickhouse-jdbc" % {{ clickhouse_jdbc_version }} classifier "all"
libraryDependencies += "com.clickhouse.spark" %% clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }} % {{ stable_version }}Spark のシェルオプション (Spark SQL CLI、Spark Shell CLI、Spark Submit コマンド) を使用する場合は、必要な JAR を渡して依存関係を登録できます。
$SPARK_HOME/bin/spark-sql \
--jars /path/clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }}.jar,/path/clickhouse-jdbc-{{ clickhouse_jdbc_version }}-all.jarJAR ファイルを Spark クライアントノードにコピーしたくない場合は、代わりに次を使用できます。
--repositories https://{maven-central-mirror or private-nexus-repo} \
--packages com.clickhouse.spark:clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }},com.clickhouse:clickhouse-jdbc:{{ clickhouse_jdbc_version }}注: SQL のみのユースケースでは、本番環境では Apache Kyuubi の使用を推奨します。
ライブラリをダウンロードする
バイナリ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-sqlCLI コマンド) に渡す。 - コンテキストの初期化時に設定を追加する。
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()val 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()Dataset<Row> 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()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()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()// テーブルは明示的な 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()// テーブルは明示的な 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 モード (NONE、STRICT など) |
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()// 上書きモード(最初にテーブルを空にする)
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.mode("overwrite")
.save()// 上書きモード(最初にテーブルを空にする)
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_key、index_granularity、その他のテーブルレベルまたはクエリレベルの設定など、ClickHouse 固有の設定を構成できます。これらは、コネクタが ClickHouse に接続する方法を制御するコネクタ オプション (host、database、table など) とは異なります。
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()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()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();
}object NativeSparkRead extends App {
val 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
val df = spark.sql("select * from clickhouse.default.example_table")
df.show()
spark.stop()
}from pyspark.sql import SparkSession
packages = [
"com.clickhouse.spark:clickhouse-spark-runtime-3.4_2.12:0.8.0",
"com.clickhouse:clickhouse-client:0.7.0",
"com.clickhouse:clickhouse-http-client:0.7.0",
"org.apache.httpcomponents.client5:httpclient5:5.2.1"
]
spark = (SparkSession.builder
.config("spark.jars.packages", ",".join(packages))
.getOrCreate())
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "127.0.0.1")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "http")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "8123")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "123456")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
df = spark.sql("select * from clickhouse.default.example_table")
df.show() CREATE TEMPORARY VIEW jdbcTable
USING org.apache.spark.sql.jdbc
OPTIONS (
url "jdbc:ch://localhost:8123/default",
dbtable "schema.tablename",
user "username",
password "password",
driver "com.clickhouse.jdbc.ClickHouseDriver"
);
SELECT * FROM jdbcTable;データの書き込み
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();
}object NativeSparkWrite extends App {
// Spark セッションを作成
val spark: SparkSession = 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 のスキーマを定義
val rows = Seq(Row(1, "John"), Row(2, "Doe"))
val schema = List(
StructField("id", DataTypes.IntegerType, nullable = false),
StructField("name", StringType, nullable = true)
)
// DataFrame を作成
val df: DataFrame = spark.createDataFrame(
spark.sparkContext.parallelize(rows),
StructType(schema)
)
df.writeTo("clickhouse.default.example_table").append()
spark.stop()
}from pyspark.sql import SparkSession
from pyspark.sql import Row
# 上記の互換性マトリクスを満たす任意のパッケージの組み合わせを使用できます。
packages = [
"com.clickhouse.spark:clickhouse-spark-runtime-3.4_2.12:0.8.0",
"com.clickhouse:clickhouse-client:0.7.0",
"com.clickhouse:clickhouse-http-client:0.7.0",
"org.apache.httpcomponents.client5:httpclient5:5.2.1"
]
spark = (SparkSession.builder
.config("spark.jars.packages", ",".join(packages))
.getOrCreate())
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "127.0.0.1")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "http")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "8123")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "123456")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
# DataFrame を作成
data = [Row(id=11, name="John"), Row(id=12, name="Doe")]
df = spark.createDataFrame(data)
# DataFrame を ClickHouse に書き込む
df.writeTo("clickhouse.default.example_table").append() -- resultTable は、clickhouse.default.example_table に挿入する Spark の中間 DataFrame です
INSERT INTO TABLE clickhouse.default.example_table
SELECT * FROM resultTable;
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()# JSONカラムをVariantTypeとして読み込む
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# Variantデータにアクセス
df.show()
# 確認のため、VariantをJSON文字列に変換
from pyspark.sql.functions import to_json
df.select(
"id",
to_json("data").alias("data_json")
).show()// JSONカラムをVariantTypeとして読み込む
Dataset<Row> df = spark.sql("SELECT id, data FROM clickhouse.default.json_table");
// Variantデータにアクセス
df.show();
// 確認のため、VariantをJSON文字列に変換
import static 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'
)
""")from pyspark.sql.functions import parse_json
# JSON データを含む DataFrame を作成
json_data = [
(1, '{"name": "Alice", "age": 30}'),
(2, '{"name": "Bob", "age": 25}'),
(3, '{"name": "Charlie", "city": "NYC"}')
]
df = spark.createDataFrame(json_data, ["id", "json_string"])
# JSON 文字列を VariantType にパース
variant_df = df.select(
"id",
parse_json("json_string").alias("data")
)
# JSON 型を使って ClickHouse に書き込む
variant_df.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'
)
""")import static org.apache.spark.sql.functions.*;
// JSON データを含む DataFrame を作成
List<Row> jsonData = Arrays.asList(
RowFactory.create(1, "{\"name\": \"Alice\", \"age\": 30}"),
RowFactory.create(2, "{\"name\": \"Bob\", \"age\": 25}"),
RowFactory.create(3, "{\"name\": \"Charlie\", \"city\": \"NYC\"}")
);
StructType schema = new StructType(new StructField[]{
DataTypes.createStructField("id", DataTypes.IntegerType, false),
DataTypes.createStructField("json_string", DataTypes.StringType, false)
});
Dataset<Row> jsonDF = spark.createDataFrame(jsonData, schema);
// JSON 文字列を VariantType にパース
Dataset<Row> variantDF = jsonDF.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 になる# JSON/Variant を VariantType ではなく文字列として読み取る
spark.conf.set("spark.clickhouse.read.jsonAs", "string")
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# data カラムは JSON 文字列を含む StringType になる// JSON/Variant を VariantType ではなく文字列として読み取る
spark.conf().set("spark.clickhouse.read.jsonAs", "string");
Dataset<Row> 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 型に推奨ベストプラクティス
- JSON 専用のデータには JSON 型を使用する: JSON object だけを保存する場合は、デフォルトの JSON 型 (
variant_typesプロパティなし) を使用します - 型を明示的に指定する:
Variant()を使用する場合は、保存する予定のすべての型を明示的に列挙します - 実験的機能を有効にする: ClickHouse で
allow_experimental_json_type = 1が有効になっていることを確認します - 書き込みには JSON フォーマットを使用する: 互換性を高めるため、VariantType データの書き込みには JSON フォーマットを使用することを推奨します
- クエリパターンを考慮する: JSON/Variant 型は、効率的にフィルタリングするための ClickHouse の JSON パスクエリをサポートしています
- パフォーマンス向上のためのカラムヒント: 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)from pyspark.sql.functions import parse_json, to_timestamp
# 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'
)
""")
# 混合型のデータを準備
events = [
(1, "2024-01-01 10:00:00", '{"action": "login", "user_id": 123}'),
(2, "2024-01-01 10:05:00", '{"action": "purchase", "amount": 99.99}'),
(3, "2024-01-01 10:10:00", '{"action": "logout", "duration": 600}')
]
df = spark.createDataFrame(events, ["event_id", "event_time", "json_data"])
# VariantTypeに変換して書き込み
variant_events = df.select(
"event_id",
to_timestamp("event_time").alias("event_time"),
parse_json("json_data").alias("event_data")
)
variant_events.writeTo("clickhouse.default.events").append()
# 読み取りとクエリ
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(truncate=False)import static 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'" +
")");
// 複数の型が混在するデータを準備
List<Row> events = Arrays.asList(
RowFactory.create(1L, "2024-01-01 10:00:00", "{\"action\": \"login\", \"user_id\": 123}"),
RowFactory.create(2L, "2024-01-01 10:05:00", "{\"action\": \"purchase\", \"amount\": 99.99}"),
RowFactory.create(3L, "2024-01-01 10:10:00", "{\"action\": \"logout\", \"duration\": 600}")
);
StructType eventSchema = new StructType(new StructField[]{
DataTypes.createStructField("event_id", DataTypes.LongType, false),
DataTypes.createStructField("event_time", DataTypes.StringType, false),
DataTypes.createStructField("json_data", DataTypes.StringType, false)
});
Dataset<Row> eventsDF = spark.createDataFrame(events, eventSchema);
// VariantTypeに変換して書き込み
Dataset<Row> variantEvents = eventsDF.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();
// 読み取りとクエリ
Dataset<Row> 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.ignoreUnsupportedTransform を false に設定してください。 |
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 コネクタの改善にご協力いただき、ありがとうございます!