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 г.