Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Spark-коннектор

Поддерживается в ClickHouse

Этот коннектор использует оптимизации 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.

Скачайте библиотеку

Имя бинарного 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 (или в команды CLI spark-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()

Пример записи с использованием 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()

Возможности 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()

Параметры подключения 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()

Настройка параметров 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()

Использование Catalog API

При использовании Catalog API укажите в конфигурации Spark формат spark.sql.catalog.<catalog_name>.option.<key>:

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

Или задайте их при создании таблиц через Spark SQL:

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

Настройки ClickHouse Cloud

При подключении к ClickHouse Cloud обязательно включите SSL и укажите подходящий режим SSL. Например:

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

Чтение данных

public static void main(String[] args) {
        // Создание сеанса Spark
        SparkSession spark = SparkSession.builder()
                .appName("example")
                .master("local[*]")
                .config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
                .config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
                .config("spark.sql.catalog.clickhouse.protocol", "http")
                .config("spark.sql.catalog.clickhouse.http_port", "8123")
                .config("spark.sql.catalog.clickhouse.user", "default")
                .config("spark.sql.catalog.clickhouse.password", "123456")
                .config("spark.sql.catalog.clickhouse.database", "default")
                .config("spark.clickhouse.write.format", "json")
                .getOrCreate();

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

        df.show();

        spark.stop();
    }

Запись данных

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

        // Создание сеанса Spark
        SparkSession spark = SparkSession.builder()
                .appName("example")
                .master("local[*]")
                .config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
                .config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
                .config("spark.sql.catalog.clickhouse.protocol", "http")
                .config("spark.sql.catalog.clickhouse.http_port", "8123")
                .config("spark.sql.catalog.clickhouse.user", "default")
                .config("spark.sql.catalog.clickhouse.password", "123456")
                .config("spark.sql.catalog.clickhouse.database", "default")
                .config("spark.clickhouse.write.format", "json")
                .getOrCreate();

        // Определение схемы для DataFrame
        StructType schema = new StructType(new StructField[]{
                DataTypes.createStructField("id", DataTypes.IntegerType, false),
                DataTypes.createStructField("name", DataTypes.StringType, false),
        });

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

        // Создание DataFrame
        Dataset<Row> df = spark.createDataFrame(data, schema);

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

        spark.stop();
    }

Операции DDL

Вы можете выполнять 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()

Запись данных типа 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'
  )
""")

Создание таблиц 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-строки

Поддержка форматов записи

Поддержка записи 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

Лучшие практики

  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. Подсказки для столбцов для повышения производительности: При использовании полей 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)

Конфигурации

Ниже перечислены настраиваемые конфигурации, доступные в коннекторе.


Параметр По умолчанию Описание С версии
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!

Navigation