Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 12.11 · 30 мин
Продвинутый
ClickPipesCDCPeerDBingestionClickHouse CloudTerraform

ClickPipes: managed ingestion для ClickHouse Cloud

ClickPipes — нативный managed-сервис непрерывной загрузки данных в ClickHouse Cloud. Вместо самостоятельной настройки Kafka Connect, Debezium или MaterializedPostgreSQL пользователь подключает источник через UI/Terraform, а ClickHouse Cloud управляет всем: репликационными слотами, backpressure, schema evolution, retry, мониторингом. Это рекомендуемый production-путь для CH Cloud, особенно для Postgres CDC, который GA с мая 2025 на базе PeerDB (приобретён ClickHouse в июле 2024).


Что такое ClickPipes и зачем он нужен

ClickPipes решает три исторические боли ingestion в ClickHouse:

  1. Operational complexity. Самостоятельный CDC-стек (Debezium + Kafka Connect + Schema Registry + sink connector) — это минимум 4 движущиеся части, каждая со своим SLA, мониторингом и операционной нагрузкой.
  2. Latency vs throughput tradeoff. В ручном Kafka-pipeline балансирование batch size, consumer lag, partition allocation и MV-INSERT throughput ложится на инженера.
  3. Schema evolution. ALTER TABLE в источнике (Postgres/MySQL) обычно ломает MaterializedPostgreSQL и требует ручного вмешательства; ClickPipes детектит DDL и применяет совместимые изменения автоматически.

ClickPipes выступает как control plane: пользователь декларирует источник, целевую таблицу и политику schema evolution, а ClickHouse Cloud разворачивает изолированный compute pool для каждого pipeline. Под капотом — PeerDB-движок (Go) для Postgres/MySQL/MongoDB CDC и собственные Kafka/Kinesis/S3 коннекторы.

INFO

ClickPipes доступен только в ClickHouse Cloud. Self-hosted ClickHouse продолжает использовать MaterializedPostgreSQL (урок 06), Debezium + Kafka Connect (урок 03) или Kafka Engine (урок 02). PeerDB Enterprise при этом open-source под ELv2 — его можно развернуть самостоятельно для Postgres CDC, но это уже отдельная инфраструктура, а не managed-сервис.


Поддерживаемые источники

На апрель 2026 ClickPipes покрывает следующие источники (статусы периодически меняются — смотрите официальный changelog):

ИсточникСтатусПод капотом
Postgres CDCGA с мая 2025PeerDB (logical replication slot, pgoutput)
MySQL CDCGAPeerDB (binlog, row-based)
MongoDB CDCPublic BetaPeerDB (change streams, oplog)
Kafka (Confluent, MSK, Redpanda, Azure Event Hubs, WarpStream)GANative ClickPipes consumer
Amazon Kinesis Data StreamsGANative ClickPipes consumer
S3 / GCS / Azure Blob (object storage)GANative batch/streaming reader
BigQueryGANative exporter

Postgres CDC — флагман. За год после private preview ClickPipes Postgres вырос с горстки клиентов до 400+ компаний и 200+ ТБ ежемесячной репликации, что делает ClickHouse самым быстрорастущим CDC-таргетом среди OLAP-баз (обгоняя Snowflake и BigQuery по росту).


Архитектура: что происходит под капотом

Архитектура ClickPipes разная для каждого источника, но базовый паттерн один: исходный CDC-поток декодируется в строки, затем батчится и пишется в целевую ReplacingMergeTree (или иную) таблицу.

ClickPipes Postgres CDC: путь данных от WAL до ClickHouse
PostgreSQL primary (wal_level=logical, publication, replication slot)PostgreSQL с wal_level=logical и созданной публикацией (CREATE PUBLICATION). ClickPipes создаёт logical replication slot с failover-флагом — slot переживает HA-переключение primary, что критично для production.
pgoutput logical decoding
ClickPipes Reader (PeerDB engine)ClickPipes Reader (PeerDB-движок, Go): держит постоянное подключение к replication slot — никогда не отпускает, чтобы избежать дорогих restart-циклов. Декодирует WAL в логические события (INSERT/UPDATE/DELETE с full row values), батчит, отправляет на writer.
Batched events + schema metadata
ClickPipes Writer (ClickHouse Cloud)ClickPipes Writer: применяет батчи к целевой ReplacingMergeTree. Для UPDATE/DELETE пишет новую версию строки с увеличенным _peerdb_version. ReplacingMergeTree(_peerdb_version) при merge оставляет максимальную версию = актуальное состояние.
ReplacingMergeTree(_peerdb_version)
_peerdb_synced_at_peerdb_synced_at: timestamp применения события. Полезно для мониторинга lag: max(_peerdb_synced_at) vs now() показывает текущее отставание репликации.
_peerdb_is_deleted_peerdb_is_deleted: флаг 0/1 для tombstone-семантики DELETE. ReplacingMergeTree с argMax по версии возвращает последнее состояние, FINAL фильтрует _peerdb_is_deleted=1.
_peerdb_version_peerdb_version: монотонный счётчик. ReplacingMergeTree(_peerdb_version) при slow merge оставляет строку с максимальной версией. Запросы должны использовать FINAL или argMax для корректного представления.

Ключевые архитектурные решения, которые отличают ClickPipes от наивной самостоятельной реализации:

  • Persistent replication connection. Reader не отпускает соединение к Postgres между батчами — это устраняет дорогие restart-циклы и резко снижает probability of slot drainage у крупных клиентов.
  • Failover replication slots. ClickPipes поддерживает Postgres 17 failover slots — slot переживает promotion реплики и продолжает читать с того же LSN, не теряя событий.
  • Parallel snapshot. Initial load использует партицирование по primary key и параллельные потоки — заявленное ускорение до 10x против sequential pg_dump, ТБ переливаются за часы вместо дней.
  • 50+ pre-flight checks. Перед запуском pipeline проверяет конфигурацию Postgres: wal_level=logical, права репликационного юзера, наличие primary key или REPLICA IDENTITY FULL у таблиц, доступность сети, размер max_wal_size. Большинство производственных аварий ловятся до старта репликации.

Для MySQL CDC под капотом — чтение row-based binlog с GTID-tracking; для MongoDB CDC — change streams через oplog (поддерживаются sharded clusters и Amazon DocumentDB). Все три CDC-коннектора используют одну общую PeerDB-кодовую базу, что объясняет схожесть служебных колонок (_peerdb_*).


Сравнение: ClickPipes vs Debezium vs MaterializedPostgreSQL

КритерийClickPipes Postgres CDCMaterializedPostgreSQLDebezium + Kafka Connect Sink
ПлатформаТолько ClickHouse CloudЛюбой self-hosted ClickHouseЛюбой ClickHouse + Kafka-стек
Operational overheadManaged (несколько кликов или Terraform)Самоуправляемое, встроено в ClickHouseKafka Connect + Schema Registry + connectors
Latency (steady state)СекундыСекунды-десятки секундСекунды-минуты (зависит от Kafka topic flush)
Initial snapshot speedПараллельный, до 10x быстрееSequential, медленнееПараллельный (Debezium incremental snapshot)
Schema evolutionАвто (ADD/DROP COLUMN), для type changes — resyncЧасто прерывается, требует ручного фиксаЗависит от connector + SMT
DDL handlingПарсит DDL, применяет совместимые измененияDROP TABLE прерывает репликациюЗависит от connector
Failover slotsДа (Postgres 17+)НетЗависит от Debezium config
Pre-flight validation50+ checksНетМинимальная
Стоимость~5x дешевле external ETLБесплатно (compute included)Стоимость Kafka + Connect compute
Гибкость трансформацийОграничена (column mapping, exclusion)Гибкая (через MV в ClickHouse)Максимальная (SMT, custom sinks, multi-target)
Multi-target (один источник, несколько sinks)НетНетДа (это сильная сторона Kafka)

Когда что выбирать (decision tree):

  1. ClickHouse Cloud + Postgres/MySQL/MongoDB источник -> ClickPipes. Без вариантов: дешевле, надёжнее, меньше operational load.
  2. Self-hosted ClickHouse + прямая сеть до источника + нет существующего Kafka-стека -> MaterializedPostgreSQL/MaterializedMySQL (урок 06).
  3. Self-hosted + уже есть Kafka + нужны SMT, multi-sink или enterprise governance -> Debezium + ClickHouse Sink Connector (урок 03).
  4. Self-hosted + хочется PeerDB-движка -> PeerDB Enterprise OSS (ELv2) — компромисс: тот же движок, что в ClickPipes, но ставите и оперируете сами.
TIP

Если вы на ClickHouse Cloud и тратите более $1–2k/месяц на Fivetran/Airbyte/Estuary для Postgres -> ClickHouse, миграция на ClickPipes почти всегда окупается: сам ClickHouse оценивает разницу примерно в 5x в свою пользу, плюс убирается один SaaS-вендор из критического пути данных.


Setup: UI flow

Базовый UI-flow для Postgres CDC (примерно одинаков для других CDC-коннекторов):

  1. Подготовка Postgres.
    • wal_level=logical (требует перезапуска, не reload).
    • Создать replication user с REPLICATION атрибутом.
    • Создать публикацию: CREATE PUBLICATION clickpipes_pub FOR ALL TABLES; (или явный список).
    • Убедиться, что max_replication_slots и max_wal_senders не исчерпаны.
  2. В ClickHouse Cloud Console -> Data Sources -> ClickPipes -> New Pipe -> Postgres CDC.
  3. Configure source. Хост, порт, креды replication user. Опционально — SSH tunnel или AWS PrivateLink для приватных Postgres.
  4. Pre-flight checks автоматически проверяют 50+ условий и предлагают исправления.
  5. Select tables. Выбираете таблицы из публикации, для каждой задаёте mapping: целевая БД, имя таблицы, ORDER BY.
  6. Schema evolution policy. ADD COLUMN — применять авто; DROP COLUMN — применять или игнорировать; type changes — pause или resync.
  7. Launch. ClickPipes выполняет initial snapshot параллельно, затем переходит в CDC-режим.

Setup: Terraform / OpenAPI

С 2025 ClickPipes доступен через clickhouse_clickpipe ресурс в официальном Terraform-провайдере (стабильно с v3.14.0). OpenAPI-эндпоинты также вышли из беты. Это критично для production: pipeline становится частью IaC, ревьюится в PR и воссоздаётся в любом environment.

# Минимальный пример Postgres CDC pipeline через Terraform
resource "clickhouse_clickpipe" "postgres_orders" {
  service_id = clickhouse_service.analytics.id
  name       = "postgres-orders-cdc"

  source = {
    postgres = {
      host     = "postgres.internal:5432"
      database = "production"
      username = "clickpipes_replicator"
      password = var.pg_replication_password
      slot_name      = "clickpipes_orders"
      publication    = "clickpipes_pub"
      ssl_mode       = "require"
    }
  }

  destination = {
    database = "raw"
    tables = [
      {
        source_name      = "public.orders"
        destination_name = "orders_cdc"
        order_by         = ["id"]
      },
      {
        source_name      = "public.order_items"
        destination_name = "order_items_cdc"
        order_by         = ["order_id", "item_id"]
      }
    ]
  }

  schema_evolution = {
    add_column_policy  = "apply"
    drop_column_policy = "ignore"
    type_change_policy = "pause"
  }
}
WARNING

Поле password в clickhouse_clickpipe хранится в Terraform state в открытом виде. В production используйте remote state с шифрованием (S3 + KMS, Terraform Cloud) и читайте секрет через data-source (Vault, AWS Secrets Manager) — никогда не коммитьте terraform.tfvars с паролями.

OpenAPI endpoint позволяет управлять pipeline из CI/CD-скриптов без Terraform — например, чтобы автоматически создавать ClickPipe для новой preview-среды Postgres в PR pipeline. Все ресурсы поддерживают full CRUD и importable в существующий state.


Цены и метеринг

Pre-GA ClickPipes Postgres CDC был бесплатен в трейле. С 1 сентября 2025 метеринг включён для всех клиентов (existing и new).

Модель оплаты Postgres CDC:

  • Ingested data — биллится за гигабайт несжатых данных, прилетающих из Postgres (initial load + CDC через slot).
  • Compute — фиксированная плата (0.5 или 1 compute unit, в зависимости от тира организации) даже когда pipeline на паузе. Это compute pool, выделенный под reader/writer.

Object storage коннекторы (S3/GCS/Azure Blob) обычно дешевле — биллится только compute, без per-GB ingest charge, что делает их экономичными для bulk-load сценариев.

Внутренние расчёты ClickHouse показывают, что ClickPipes Postgres CDC примерно в 5x дешевле external ETL-инструментов (Fivetran, Airbyte, Estuary) при сравнимой нагрузке. Для большинства существующих клиентов миграция увеличивает общий счёт ClickHouse Cloud на 0–15% post-trial — это намного меньше, чем счёт за external ETL, который удаляется из стека.

NOTE

Точные числа меняются — смотрите clickhouse.com/pricing и billing changelog. На момент написания (апрель 2026) per-GB цена не публична как фиксированный список, требуется calculator или contact sales для крупных volumes.


Limitations

Реальные ограничения, о которых стоит знать до выбора ClickPipes:

  1. Только ClickHouse Cloud. Self-hosted, BYOC (bring your own cluster) и dedicated не поддерживаются. PeerDB OSS — отдельный продукт.
  2. Регион ограничения. ClickPipes доступен в большинстве, но не во всех регионах CH Cloud — проверяйте matrix перед миграцией. Cross-region pipeline с источником в одном AWS-регионе и целевым в другом может стоить дороже из-за egress.
  3. Schema evolution: только совместимые изменения. ADD COLUMN — авто; DROP COLUMN — авто или ignore; rename column / type change — требуется pause + resync, потому что ReplacingMergeTree не умеет атомарно эволюционировать схему без полного refill.
  4. Нет multi-target. Один Postgres-источник = один ClickPipe = один CH-сервис. Если нужна fan-out (один Postgres -> CH + Snowflake + S3), используйте Debezium.
  5. Soft limits на parallelism. По умолчанию parallel snapshot имеет лимит на количество одновременных таблиц/партиций; для очень крупных БД (>10 ТБ initial) обращаются в support для bumping.
  6. Tables without PK / REPLICA IDENTITY FULL. Pre-flight check блокирует pipeline для таблиц, где PG не может однозначно идентифицировать строку для UPDATE/DELETE — нужно добавить PK или включить REPLICA IDENTITY FULL.
  7. Custom transformations ограничены. Можно делать column mapping, exclusion, простые type casts; сложная бизнес-логика трансформации делается уже в ClickHouse через MV или dbt поверх _cdc-таблицы.

Operations: мониторинг, retry, conflicts

Production-эксплуатация ClickPipes опирается на три источника наблюдаемости.

1. ClickPipes UI. В Cloud Console для каждого pipeline видны: ingestion rate (rows/s), ingestion volume (bytes), error count, replication lag, статус slot-а в Postgres, recent errors с DDL-парсингом.

2. Centralized system table. ClickHouse Cloud экспонирует системную таблицу с агрегированными логами всех ClickPipes — Kafka, Kinesis, S3, CDC. Можно мониторить из SQL так же, как system.merges или system.query_log:

-- Последние ошибки по всем ClickPipes за час
SELECT
    pipe_name,
    source_type,
    event_time,
    error_code,
    error_message
FROM system.clickpipes_log
WHERE event_time > now() - INTERVAL 1 HOUR
  AND error_code != ''
ORDER BY event_time DESC
LIMIT 50;

-- Throughput по pipeline за день
SELECT
    pipe_name,
    sum(rows_ingested) AS rows,
    sum(bytes_ingested) AS bytes,
    avg(latency_ms) AS avg_latency_ms
FROM system.clickpipes_log
WHERE event_time > now() - INTERVAL 1 DAY
GROUP BY pipe_name
ORDER BY rows DESC;

3. Prometheus metrics. ClickPipes публикует метрики через ClickHouse Cloud Prometheus integration: clickpipes_rows_ingested_total, clickpipes_bytes_ingested_total, clickpipes_errors_total, clickpipes_replication_lag_seconds. Это позволяет строить алерты в Grafana/Datadog рядом с остальной инфраструктурой.

Retry semantics. Постоянные сетевые ошибки -> exponential backoff с автоматическим возобновлением. Permanent errors (DDL, который ClickPipes не умеет применить, или нарушение schema policy) -> pipeline ставится на pause, операторы получают notification в UI/Slack/email.

DDL handling. ClickPipes парсит DDL-события из WAL/binlog. Совместимые изменения (ADD COLUMN со совпадающим типом, DROP COLUMN при policy=apply) применяются автоматически. Несовместимые (rename, type narrowing, RENAME TABLE) переводят pipeline в paused — оператор в UI делает resync для соответствующих таблиц.

Conflicts handling. ReplacingMergeTree разрешает дубли через _peerdb_version — последняя запись с максимальной версией становится actual. Для корректных запросов используйте FINAL или paramaterized view с argMax(...) GROUP BY pk. Для больших таблиц включают OPTIMIZE TABLE ... FINAL периодически или используют ReplacingMergeTree(version, is_deleted) с clean_deleted_rows=true.

-- Корректный SELECT поверх ClickPipes-таблицы
SELECT *
FROM raw.orders_cdc FINAL
WHERE _peerdb_is_deleted = 0
  AND order_date >= today() - 7;

Ключевые выводы

  1. ClickPipes — managed integration platform для ClickHouse Cloud; Postgres CDC GA с мая 2025 на базе PeerDB (приобретение ClickHouse 2024).
  2. Поддерживаются Postgres/MySQL CDC (GA), MongoDB CDC (Public Beta), Kafka, Kinesis, object storage (S3/GCS/Azure), BigQuery.
  3. Архитектурные ключи: persistent replication connection, failover slots, параллельный initial snapshot до 10x, 50+ pre-flight checks, служебные колонки _peerdb_version/_peerdb_is_deleted/_peerdb_synced_at.
  4. Decision tree: CH Cloud -> ClickPipes; self-hosted с прямым доступом -> MaterializedPostgreSQL; Kafka-стек или multi-target -> Debezium.
  5. Управление через UI, OpenAPI или Terraform (clickhouse_clickpipe, GA с v3.14.0). IaC-flow рекомендуется для production.
  6. Метеринг с 1 сентября 2025: ingest per GB + compute per hour, ~5x дешевле external ETL.
  7. Limitations: только CH Cloud, ограничения регионов, schema evolution только для совместимых изменений, нет native multi-target, нужен PK/REPLICA IDENTITY FULL.
  8. Observability на трёх уровнях: UI, system.clickpipes_log, Prometheus. Алерты на replication_lag_seconds и errors_total обязательны для production.
CDC-паттерны: WAL replication, Debezium и change streams Logical replication в PostgreSQL: publication, subscription и decoding

Проверь себя

Вопрос 1. На self-hosted ClickHouse 26.3 LTS нужно настроить CDC из Postgres 16. Архитектурно доступны три варианта: ClickPipes, MaterializedPostgreSQL и Debezium + ClickHouse Sink Connector. Какой вариант невозможен и почему?

Ответ

ClickPipes невозможен — он доступен только в ClickHouse Cloud (managed). На self-hosted остаются два варианта: MaterializedPostgreSQL (если есть прямая сеть PG -> CH и не нужен fan-out) или Debezium + Kafka Connect + ClickHouse Sink Connector (если уже есть Kafka-стек или нужны multi-sink/SMT). Третий, менее популярный путь — самостоятельно развернуть PeerDB Enterprise OSS (ELv2), это даст тот же движок, что в ClickPipes, но без managed-сервиса.

Вопрос 2. ClickPipes Postgres CDC использует ReplacingMergeTree со служебной колонкой _peerdb_version. Запрос SELECT count() FROM orders_cdc WHERE order_date = today() иногда возвращает завышенное значение. В чём причина и как починить?

Ответ

ReplacingMergeTree схлопывает дубли только при background merge — между UPDATE-событиями в исходной БД и моментом merge в ClickHouse в таблице может одновременно лежать несколько версий одной строки. SELECT без FINAL видит все версии, поэтому count() завышен. Решения: добавить FINAL (медленнее, но точно), или использовать argMax(version_field, _peerdb_version) GROUP BY pk для актуальных значений, или периодически OPTIMIZE TABLE ... FINAL. Также фильтр _peerdb_is_deleted = 0 исключает tombstone-строки удалённых записей.

Вопрос 3. Команда мигрирует из Fivetran (~$3k/мес за Postgres -> ClickHouse) в ClickPipes. Какие три проверки сделать до старта миграции, чтобы не получить сюрпризов в production?

Ответ

(1) Pre-flight checks ClickPipes — пройти 50+ автоматических проверок Postgres (wal_level=logical, replication slots, права юзера, PK/REPLICA IDENTITY на всех таблицах). Большинство production-аварий ловится тут. (2) Регион CH Cloud vs регион Postgres — кросс-регион даёт egress costs, которые могут съесть экономию vs Fivetran. (3) Объём initial snapshot и compute tier — параллельный snapshot до 10x быстрее, но крупные БД (>10 ТБ) требуют поднятия soft limits через support; одновременно compute units биллится с момента запуска, поэтому стоит запускать когда есть capacity для постоянного потока, а не “потестим раз в неделю”. Дополнительно: подготовить мониторинг (system.clickpipes_log + Prometheus) и план schema evolution (какие type changes требуют resync — заранее договориться с Postgres-командой о window).

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 11