此连接器利用 ClickHouse 特有的优化能力,例如高级分区和谓词下推,以 提升查询性能和数据处理效率。 该连接器基于 ClickHouse 官方 JDBC connector,并 自行管理其目录。
在 Spark 3.0 之前,Spark 缺少内置的目录概念,因此用户通常依赖 Hive Metastore 或 AWS Glue 等外部目录系统。 使用这些外部方案时,用户必须先手动注册数据源表,然后才能在 Spark 中访问它们。 不过,自 Spark 3.0 引入目录概念后,Spark 现在可以通过注册 目录插件来自动发现表。
Spark 的默认目录是 spark_catalog,表通过 {catalog name}.{database}.{table} 标识。借助这一新的
目录功能,现在可以在单个 Spark 应用中添加并使用多个目录。
在 Catalog API 和 TableProvider API 之间进行选择
ClickHouse Spark connector 支持两种访问方式:Catalog API 和 TableProvider API (基于 format 的访问) 。了解它们之间的差异,有助于你根据具体用例选择更合适的方式。
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 connector 直接作为依赖项添加到项目的构建文件中 (例如 Maven 的 pom.xml
或 SBT 的 build.sbt) 。
或者,也可以将所需的 JAR 文件放入 $SPARK_HOME/jars/ 目录中,或在 spark-submit 命令中
使用 --jars 参数将其直接作为 Spark 选项传入。
这两种方式都能确保 ClickHouse 连接器在您的 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 版本,请参阅 Sonatype 的通过 Maven 使用 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 版本,请参阅 Sonatype 的通过 Gradle 使用 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 的 shell 选项 (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.jar如果你不想将 JAR 文件复制到 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您可以在 Maven Central Repository 中找到所有已发布的 JAR 文件。 每日构建的 SNAPSHOT JAR 文件可通过上方配置的 Sonatype snapshots 仓库获取。
注册 目录 (必需)
要访问您的 ClickHouse 表,必须使用以下配置新增一个 Spark 目录:
| Property | Value | Default Value | Required |
|---|---|---|---|
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 命令) 。 - 在初始化 Context 时添加配置。
使用 TableProvider API (基于格式的访问)
除了基于 目录 的方法外,ClickHouse Spark connector 还支持通过 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 功能
TableProvider API 具备多项强大功能:
自动创建表
当写入一个不存在的表时,连接器会自动使用合适的 schema 创建该表。连接器提供了以下智能默认值:
- Engine:如果未指定,默认使用
MergeTree()。你可以使用engine选项指定其他引擎 (例如ReplacingMergeTree()、SummingMergeTree()等) - ORDER BY:必需 - 创建新表时,必须显式指定
order_by选项。连接器会验证所有指定的列都存在于 schema 中。 - Nullable Key 支持:如果 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 |
表名称 | N/A | 是 |
user |
用于身份验证的用户名 | default |
否 |
password |
用于身份验证的密码 | (空字符串) | 否 |
ssl |
启用 SSL 连接 | false |
否 |
ssl_mode |
SSL 模式 (NONE、STRICT 等) |
STRICT |
否 |
timezone |
用于日期/时间操作的时区 | server |
否 |
表创建选项
当表不存在且需要创建时,可使用以下选项:
| Option | Description | Default Value | Required |
|---|---|---|---|
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选项。所有指定列都必须存在于 schema 中。 ** 如果 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 选项,你可以配置 ClickHouse 特有的设置,例如 allow_nullable_key、index_granularity 以及其他表级或查询级设置。这些不同于连接器选项 (如 host、database、table) ;后者用于控制连接器如何连接到 ClickHouse。
使用 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 的 schema
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 的 schema
val rows = Seq(Row(1, "John"), Row(2, "Doe"))
val schema = List(
StructField("id", DataTypes.IntegerType, nullable = false),
StructField("name", StringType, nullable = true)
)
// 创建 df
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 查询,你可以在应用程序中使用任意 API——Java、Scala、 PySpark 或 shell——运行这些查询。
使用 VariantType
该连接器 支持 Spark 的 VariantType,可用于处理半结构化数据。VariantType 会映射到 ClickHouse 的 JSON 和 Variant 类型,让您能够高效存储和查询 schema 灵活的数据。
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 对象)
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(仅限 JSON 对象)
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 对象)
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 DDL 创建 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 类型
要支持基本类型、数组和 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 类型
以下 ClickHouse 类型可用于 Variant():
- 基本类型:
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 数据。 |
| Arrow | ⚠️ 部分支持 | 支持写入 ClickHouse JSON 类型,不支持 ClickHouse Variant 类型。完整支持仍需等待 https://github.com/ClickHouse/ClickHouse/issues/92752 解决。 |
配置写入格式:
spark.conf.set("spark.clickhouse.write.format", "json") // 推荐用于 Variant 类型最佳实践
- 对纯 JSON 数据使用 JSON 类型:如果你只存储 JSON 对象,请使用默认的 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 中启用 Experimental 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 中启用 Experimental 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 ...,会将查询 schema 的所有 field 都标记为 Nullable。请注意,此配置需要 SPARK-43390 (在 Spark 3.5 中可用) ;如果没有这个补丁,则始终等同于 true。 |
0.8.0 |
| spark.clickhouse.write.batchSize | 10000 | 写入 ClickHouse 时每个批次包含的记录数。 | 0.1.0 |
| spark.clickhouse.write.compression.codec | lz4 | 用于压缩写入数据的编解码器。支持的编解码器:none、lz4。 | 0.3.0 |
| spark.clickhouse.write.distributed.convertLocal | false | 写入 Distributed 表时,会改为写入其本地表,而不是 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 | 写入 Distributed 表时,会写入集群中的所有节点。 | 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 返回的可重试 error 代码。 | 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 |
✅ | 否 | 数组元素类型也会一并转换 |
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 |
✅ | 否 | 同时支持命名元组和非命名元组。命名元组按名称映射为结构体字段,非命名元组使用 _1、_2 等字段名。支持嵌套结构体和可空字段 |
Point |
❌ | |||
Polygon |
❌ | |||
MultiPolygon |
❌ | |||
Ring |
❌ | |||
IntervalQuarter |
❌ | |||
IntervalWeek |
❌ | |||
Decimal256 |
❌ | |||
AggregateFunction |
❌ | |||
SimpleAggregateFunction |
❌ |
将数据从 Spark 插入到 ClickHouse
| Spark 数据类型 | ClickHouse 数据类型 | 支持 | 是否为基本类型 | 说明 |
|---|---|---|---|---|
BooleanType |
Bool |
✅ | 是 | 自 0.9.0 版本起映射为 Bool 类型 (而不是 UInt8) |
ByteType |
Int8 |
✅ | 是 | |
ShortType |
Int16 |
✅ | 是 | |
IntegerType |
Int32 |
✅ | 是 | |
LongType |
Int64 |
✅ | 是 | |
FloatType |
Float32 |
✅ | 是 | |
DoubleType |
Float64 |
✅ | 是 | |
StringType |
String |
✅ | 是 | |
VarcharType |
String |
✅ | 是 | |
CharType |
String |
✅ | 是 | |
DecimalType |
Decimal(p, s) |
✅ | 是 | 精度和标度最高支持 Decimal128 |
DateType |
Date |
✅ | 是 | |
TimestampType |
DateTime |
✅ | 是 | |
ArrayType (list, tuple, or array) |
Array |
✅ | 否 | Array 的元素类型也会被转换 |
MapType |
Map |
✅ | 否 | 键仅限为 StringType |
StructType |
Tuple |
✅ | 否 | 会转换为带字段名的命名 Tuple。 |
VariantType |
JSON or Variant |
✅ | 否 | 需要 Spark 4.0+ 和 ClickHouse 25.3+。默认使用 JSON 类型。使用 clickhouse.column.<name>.variant_types 属性可指定包含多种类型的 Variant。 |
Object |
❌ | |||
Nested |
❌ |
贡献与支持
如果您想为该项目作出贡献或报告问题,欢迎向我们反馈! 请访问我们的 GitHub 代码仓库 创建 issue、提出改进建议, 或提交拉取请求。 欢迎贡献!开始之前,请先查看代码仓库中的贡献指南。 感谢您帮助改进我们的 ClickHouse Spark 连接器!