Этот коннектор использует оптимизации ClickHouse, такие как расширенное партиционирование и pushdown предикатов, чтобы повысить производительность запросов и эффективность обработки данных. Коннектор основан на официальном коннекторе JDBC для ClickHouse и управляет собственным каталогом.
До 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 поддерживает два варианта доступа: Catalog API и TableProvider API (доступ на основе формата). Понимание различий между ними поможет выбрать подходящий вариант для вашего сценария использования.
Catalog API vs 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 как зависимость напрямую в файл сборки проекта (например, в pom.xml
для Maven или build.sbt для SBT).
Либо можно поместить необходимые JAR-файлы в каталог $SPARK_HOME/jars/ или передать их напрямую как параметр Spark
с помощью флага --jars в команде spark-submit.
Оба подхода позволяют сделать коннектор 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 по использованию SNAPSHOT-релизов для Maven.
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 по использованию SNAPSHOT-релизов для Gradle.
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.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Все доступные выпущенные JAR-файлы можно найти в репозитории Maven Central. JAR-файлы ежедневных SNAPSHOT-сборок доступны через указанный выше репозиторий снимков Sonatype.
Зарегистрируйте каталог (обязательно)
Чтобы получить доступ к таблицам 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(или в команды CLIspark-shell/spark-sql). - Добавить конфигурацию при инициализации контекста.
Использование TableProvider API (доступ на основе формата)
Помимо подхода на основе каталога, коннектор ClickHouse Spark поддерживает доступ на основе формата через TableProvider API.
Пример чтения через format API
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
# Чтение из ClickHouse через format API
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
# Запись в ClickHouse с использованием API format
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-ключей: Автоматически добавляет
settings.allow_nullable_key=1, если ORDER BY содержит столбец с типом Nullable
# Таблица будет создана автоматически с явно указанным 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 | localhost |
Да |
protocol |
Протокол подключения (http или https) |
http |
Нет |
http_port |
Порт HTTP/HTTPS | 8123 |
Нет |
database |
Имя базы данных | default |
Да |
table |
Имя таблицы | Н/Д | Да |
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 |
Разрешить nullable-ключи в ORDER BY (для ClickHouse Cloud) | Определяется автоматически** | Нет |
settings.<key> |
Любая настройка таблицы ClickHouse | N/A | Нет |
cluster |
Имя кластера для Distributed tables | N/A | Нет |
clickhouse.column.<name>.variant_types |
Список типов ClickHouse для Variant column, разделённых запятыми (например, String, Int64, Bool, JSON). Имена типов чувствительны к регистру. Пробелы после запятых необязательны. |
N/A | Нет |
- Параметр
order_byобязателен при создании новой таблицы. Все указанные столбцы должны существовать в схеме. ** Автоматически устанавливается в1, если ORDER BY содержит столбец с типом Nullable и параметр не задан явно.
Режимы записи
Spark-коннектор (и TableProvider API, и Catalog API) поддерживает следующие режимы записи в Spark:
append: Добавляет данные в существующую таблицуoverwrite: Заменяет все данные в таблице (предварительно очищает таблицу)
# Режим overwrite (сначала очищает таблицу)
df.write \
.format("clickhouse") \
.option("host", "your-host") \
.option("database", "default") \
.option("table", "my_table") \
.mode("overwrite") \
.save()// Режим overwrite (сначала очищает таблицу)
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.mode("overwrite")
.save()// Режим 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_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
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)
)
// Создание 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 — это промежуточный df Spark, который нужно вставить в clickhouse.default.example_table
INSERT INTO TABLE clickhouse.default.example_table
SELECT * FROM resultTable;
Операции DDL
Вы можете выполнять DDL-операции в своём экземпляре ClickHouse с помощью Spark SQL, при этом все изменения сразу сохраняются в 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
Коннектор поддерживает тип VariantType в Spark для работы с полуструктурированными данными. VariantType сопоставляется с типами ClickHouse JSON и Variant, что позволяет эффективно хранить данные с гибкой схемой и выполнять по ним запросы.
Сопоставление типов ClickHouse
| Тип ClickHouse | Тип Spark | Описание |
|---|---|---|
JSON |
VariantType |
Хранит только объекты JSON (должны начинаться с {) |
Variant(T1, T2, ...) |
VariantType |
Хранит значения нескольких типов, включая примитивные типы, массивы и JSON |
Чтение данных типа VariantType
При чтении из ClickHouse столбцы JSON и Variant автоматически преобразуются в VariantType Spark:
// Чтение 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
Вы можете записывать данные типа VariantType в ClickHouse, используя типы столбцов JSON или Variant:
import org.apache.spark.sql.functions._
// Создайте DataFrame с данными JSON
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")
)
// Запишите в ClickHouse с типом JSON (только объекты 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
# Создайте DataFrame с данными JSON
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")
)
# Запишите в 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.*;
// Создайте DataFrame с данными JSON
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")
);
// Запишите в ClickHouse с типом JSON (только объекты 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'" +
")");Создание таблиц VariantType в Spark SQL
Таблицы VariantType можно создавать с помощью DDL Spark SQL:
-- Создать таблицу с типом 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 не указано, для столбца по умолчанию используется тип JSON в ClickHouse, который принимает только объекты 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
В 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 будет иметь тип StringType и содержать JSON-строки# Считывать JSON/Variant как строки вместо VariantType
spark.conf.set("spark.clickhouse.read.jsonAs", "string")
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# столбец data будет иметь тип StringType и содержать JSON-строки// Считывать 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 будет иметь тип StringType и содержать JSON-строкиПоддержка форматов записи
Поддержка записи VariantType зависит от формата:
| Format | Support | Notes |
|---|---|---|
| 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 для эффективной фильтрации
- Подсказки для столбцов для повышения производительности: При использовании полей JSON в ClickHouse добавление подсказок для столбцов повышает производительность запросов. В настоящее время добавление подсказок для столбцов через Spark не поддерживается. Отслеживать эту возможность можно в GitHub issue #497.
Пример: полный процесс
import org.apache.spark.sql.functions._
// Включить экспериментальный тип JSON в ClickHouse
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
# Включить экспериментальный тип JSON в ClickHouse
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.*;
// Включить экспериментальный тип JSON в ClickHouse
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 | При чтении distributed таблицы вместо неё используется локальная таблица. Если true, параметр spark.clickhouse.read.distributed.useClusterNodes игнорируется. |
0.1.0 |
| spark.clickhouse.read.fixedStringAs | binary | Читает тип FixedString в ClickHouse как указанный тип данных Spark. Поддерживаемые типы: binary, string | 0.8.0 |
| spark.clickhouse.read.format | json | Формат сериализации при чтении. Поддерживаемые форматы: json, binary | 0.6.0 |
| spark.clickhouse.read.runtimeFilter.enabled | false | Включить runtime filter для чтения. | 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. | 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 | Разрешает запись в таблицы Distributed с 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 возвращает при сбое записи. | 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 |
✅ | Да | Точность и scale — до 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 |
✅ | Нет | Используется соответствующий тип interval |
JSON, Variant |
VariantType |
✅ | Нет | Требуются Spark 4.0+ и ClickHouse 25.3+. Можно читать как StringType, используя spark.clickhouse.read.jsonAs=string |
Object |
❌ | |||
Nested |
❌ | |||
Tuple |
StructType |
✅ | Нет | Поддерживаются как именованные, так и неименованные Tuple. Именованные Tuple сопоставляются с полями struct по имени, а для неименованных используются _1, _2 и т. д. Также поддерживаются вложенные struct и поля Nullable |
Point |
❌ | |||
Polygon |
❌ | |||
MultiPolygon |
❌ | |||
Ring |
❌ | |||
IntervalQuarter |
❌ | |||
IntervalWeek |
❌ | |||
Decimal256 |
❌ | |||
AggregateFunction |
❌ | |||
SimpleAggregateFunction |
❌ |
Вставка данных из Spark в ClickHouse
| Тип данных Spark | Тип данных ClickHouse | Поддерживается | Примитивный | Примечания |
|---|---|---|---|---|
BooleanType |
Bool |
✅ | Да | Сопоставляется с типом Bool (а не UInt8), начиная с версии 0.9.0 |
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 |
✅ | Нет | Тип элементов массива также преобразуется |
MapType |
Map |
✅ | Нет | Ключи ограничены типом StringType |
StructType |
Tuple |
✅ | Нет | Преобразуется в именованный Tuple с именами полей. |
VariantType |
JSON или Variant |
✅ | Нет | Требует Spark 4.0+ и ClickHouse 25.3+. По умолчанию используется тип JSON. Используйте свойство clickhouse.column.<name>.variant_types, чтобы указать Variant с несколькими типами. |
Object |
❌ | |||
Nested |
❌ |
Участие в проекте и поддержка
Если вы хотите внести вклад в проект или сообщить о проблеме, мы будем рады вашей помощи! Посетите наш репозиторий GitHub, чтобы создать issue, предложить улучшения или отправить pull request. Мы приветствуем ваш вклад! Перед началом работы ознакомьтесь с рекомендациями по участию в репозитории. Спасибо, что помогаете улучшать наш коннектор ClickHouse Spark!