Skip to main content
Apache Beam — это открытая унифицированная модель программирования, которая позволяет разработчикам определять и выполнять как батч-, так и потоковые (непрерывные) конвейеры обработки данных. Гибкость Apache Beam заключается в поддержке широкого спектра сценариев обработки данных — от операций ETL (извлечение, преобразование, загрузка) до сложной обработки событий и Real-time аналитики. Эта интеграция использует официальный коннектор JDBC ClickHouse в качестве базового слоя для вставки.

Пакет интеграции

Пакет интеграции, необходимый для работы Apache Beam с ClickHouse, поддерживается и разрабатывается в рамках Apache Beam I/O Connectors — набора интеграций для множества популярных систем хранения данных и баз данных. Реализация org.apache.beam.sdk.io.clickhouse.ClickHouseIO находится в репозитории Apache Beam.

Настройка пакета ClickHouse для Apache Beam

Установка пакета

Добавьте следующую зависимость в систему управления пакетами:
Рекомендуемая версия BeamКоннектор ClickHouseIO рекомендуется использовать с Apache Beam версии 2.59.0 и выше. Более ранние версии могут не полностью поддерживать функциональность коннектора.
Артефакты доступны в официальном репозитории Maven.

Пример кода

В следующем примере CSV-файл input.csv считывается в виде PCollection, преобразуется в объект Row (с использованием заданной схемы) и вставляется в локальный экземпляр ClickHouse с помощью ClickHouseIO:

Поддерживаемые типы данных

Параметры ClickHouseIO.Write

Конфигурацию ClickHouseIO.Write можно настроить с помощью следующих функций-сеттеров:

Ограничения

Учитывайте следующие ограничения при использовании коннектора:
  • На данный момент поддерживается только операция Sink. Коннектор не поддерживает операцию Source.
  • ClickHouse выполняет дедупликацию при вставке в таблицу ReplicatedMergeTree или в таблицу Distributed, построенную поверх ReplicatedMergeTree. Без репликации вставка в обычную таблицу MergeTree может приводить к появлению дубликатов, если вставка завершается ошибкой, а затем успешно повторяется. Однако каждый блок вставляется атомарно, а размер блока можно настроить с помощью ClickHouseIO.Write.withMaxInsertBlockSize(long). Дедупликация достигается за счёт использования контрольных сумм вставленных блоков. Подробнее о дедупликации см. в разделах Deduplication и Deduplicate insertion config.
  • Коннектор не выполняет никаких DDL-операторов; поэтому целевая таблица должна существовать до вставки.
Последнее изменение 3 июля 2026 г.