Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Spark Connector

支持 ClickHouse

此连接器利用 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 APITableProvider 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 发行版说明

下载库

二进制 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-sql CLI 命令) 。
  • 在初始化 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()

基于格式的写入示例

# 使用 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 功能

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

TableProvider 连接选项

使用基于格式的 API 时,可用以下连接选项:

连接选项

选项 说明 默认值 必填
host ClickHouse server 主机名 localhost
protocol 连接协议 (httphttps) http
http_port HTTP/HTTPS 端口 8123
database 数据库名称 default
table 表名称 N/A
user 用于身份验证的用户名 default
password 用于身份验证的密码 (空字符串)
ssl 启用 SSL 连接 false
ssl_mode SSL 模式 (NONESTRICT 等) 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()

配置 ClickHouse 选项

Catalog API 和 TableProvider API 都支持配置 ClickHouse 专用选项 (而不是连接器选项) 。创建表或执行查询时,这些选项会传递给 ClickHouse。

通过 ClickHouse 选项,你可以配置 ClickHouse 特有的设置,例如 allow_nullable_keyindex_granularity 以及其他表级或查询级设置。这些不同于连接器选项 (如 hostdatabasetable) ;后者用于控制连接器如何连接到 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()

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

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 的 JSONVariant 类型,让您能够高效存储和查询 schema 灵活的数据。

ClickHouse 类型映射

ClickHouse 类型 Spark 类型 描述
JSON VariantType 仅存储 JSON 对象 (必须以 { 起始)
Variant(T1, T2, ...) VariantType 可存储多种类型,包括基本类型、数组和 JSON

读取 VariantType 数据

从 ClickHouse 读取数据时,JSONVariant 列会自动映射到 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 对象)
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()

  • 基本类型StringInt8Int16Int32Int64UInt8UInt16UInt32UInt64Float32Float64Bool
  • 数组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 ✅ 完全支持 同时支持 JSONVariant 类型。推荐用于 VariantType 数据。
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 对象,请使用默认的 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 ...,会将查询 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=trueignoreUnsupportedTransform=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 连接器!

Navigation