Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Подключение Streamkap к ClickHouse

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

Streamkap — это платформа интеграции данных в реальном времени, специализирующаяся на стриминге CDC (фиксация изменений данных) и потоковой обработке данных. Она построена на высокопроизводительном масштабируемом стеке с использованием Apache Kafka, Apache Flink и Debezium и предлагается как полностью управляемый сервис в вариантах развертывания SaaS или BYOC (Собственное облако).

Streamkap позволяет в потоковом режиме передавать в ClickHouse каждую вставку, обновление и удаление из исходных баз данных, таких как PostgreSQL, MySQL, SQL Server, MongoDB и других, с задержкой в миллисекунды.

Это делает платформу отличным выбором для аналитических панелей мониторинга в реальном времени, операционной аналитики и передачи актуальных данных в модели машинного обучения.

Ключевые возможности

  • CDC (фиксация изменений данных) в реальном времени: Streamkap считывает изменения напрямую из журналов вашей базы данных, благодаря чему данные в ClickHouse остаются актуальной репликой источника. Упрощённая потоковая обработка: преобразовывайте, обогащайте, маршрутизируйте, форматируйте данные и создавайте эмбеддинги в реальном времени до их загрузки в ClickHouse. Работает на базе Flink, но без связанной с ним сложности

  • Полностью управляемое и масштабируемое решение: Оно предоставляет готовый к промышленной эксплуатации и не требующий обслуживания конвейер, избавляя от необходимости управлять собственной инфраструктурой Kafka, Flink, Debezium или schema registry. Платформа рассчитана на высокую пропускную способность и может линейно масштабироваться для обработки миллиардов событий.

  • Автоматическая эволюция схемы: Streamkap автоматически обнаруживает изменения схемы в исходной базе данных и переносит их в ClickHouse. Он может добавлять новые столбцы и изменять типы столбцов без ручного вмешательства.

  • Оптимизировано для ClickHouse: Эта интеграция создана для эффективной работы с возможностями ClickHouse. По умолчанию она использует движок ReplacingMergeTree, чтобы без проблем обрабатывать обновления и удаления из исходной системы.

  • Надёжная доставка: Платформа обеспечивает гарантию доставки как минимум один раз, сохраняя согласованность данных между источником и ClickHouse. Для операций upsert она выполняет дедупликацию на основе первичного ключа.

Начало работы

В этом руководстве в общих чертах описано, как настроить конвейер Streamkap для загрузки данных в ClickHouse.

Предварительные требования

  • Аккаунт Streamkap.
  • Сведения о подключении к вашему кластеру ClickHouse: Hostname, Port, Username и Password.
  • Исходная база данных (например, PostgreSQL, SQL Server), настроенная для поддержки CDC (фиксация изменений данных). Подробные руководства по настройке можно найти в документации Streamkap.

Настройте источник в Streamkap

  1. Войдите в свой аккаунт Streamkap.
  2. На боковой панели перейдите в раздел Connectors и выберите вкладку Sources.
  3. Нажмите + Add и выберите тип исходной базы данных (например, SQL Server RDS).
  4. Заполните сведения о подключении, включая конечную точку, порт, имя базы данных и учетные данные пользователя.
  5. Сохраните коннектор.

Настройте пункт назначения ClickHouse

  1. В разделе Connectors выберите вкладку Destinations.
  2. Нажмите + Add и выберите ClickHouse из списка.
  3. Введите сведения о подключении для вашего сервиса ClickHouse:
    • Hostname: хост вашего экземпляра ClickHouse (например, abc123.us-west-2.aws.clickhouse.cloud)
    • Port: защищённый порт HTTPS, обычно 8443
    • Username and Password: учетные данные пользователя ClickHouse
    • Database: имя целевой базы данных в ClickHouse
  4. Сохраните пункт назначения.

Создайте и запустите конвейер

  1. Перейдите в Pipelines на боковой панели и нажмите + Create.
  2. Выберите источник и пункт назначения, которые вы только что настроили.
  3. Выберите схемы и таблицы, которые хотите передавать.
  4. Задайте имя конвейеру и нажмите Save.

После создания конвейер станет активным. Streamkap сначала создаст снимок существующих данных, а затем начнет передавать все новые изменения по мере их появления.

Проверьте данные в ClickHouse

Подключитесь к вашему кластеру ClickHouse и выполните запрос, чтобы увидеть данные, поступающие в целевую таблицу.

SELECT * FROM your_table_name LIMIT 10;

Как это работает с ClickHouse

Интеграция Streamkap предназначена для эффективной работы с данными CDC (фиксация изменений данных) в ClickHouse.

Движок таблицы и обработка данных

По умолчанию Streamkap использует режим upsert для ингестии. При создании таблицы в ClickHouse он использует движок ReplacingMergeTree. Этот движок идеально подходит для обработки событий CDC:

  • Первичный ключ исходной таблицы используется как ключ ORDER BY в определении таблицы ReplacingMergeTree.

  • Обновления в источнике записываются в ClickHouse как новые строки. В ходе фонового процесса merge в ReplacingMergeTree эти строки схлопываются, и сохраняется только самая новая версия на основе ключа ORDER BY.

  • Удаления обрабатываются с помощью флага метаданных, который передаёт значение в параметр ReplacingMergeTree is_deleted. Строки, удалённые в источнике, не удаляются сразу, а помечаются как удалённые.

    • При необходимости удалённые записи можно сохранять в ClickHouse для аналитики

Столбцы метаданных

Streamkap добавляет в каждую таблицу несколько столбцов метаданных, чтобы отслеживать состояние данных:

Имя столбца Описание
_STREAMKAP_SOURCE_TS_MS Временная метка события в исходной базе данных (в миллисекундах).
_STREAMKAP_TS_MS Временная метка момента, когда Streamkap обработал событие (в миллисекундах).
__DELETED Булев флаг (true/false), указывающий, была ли строка удалена в источнике.
_STREAMKAP_OFFSET Значение смещения из внутренних журналов Streamkap, полезное для упорядочивания и отладки.

Запрос самых свежих данных

Поскольку ReplacingMergeTree обрабатывает обновления и удаления в фоновом режиме, простой запрос SELECT * может показывать исторические или удалённые строки, пока слияние не завершено. Чтобы получить наиболее актуальное состояние данных, нужно отфильтровать удалённые записи и выбрать только последнюю версию каждой строки.

Это можно сделать с помощью модификатора FINAL — это удобно, но может сказаться на производительности запроса:

-- Использование FINAL для получения актуального состояния данных
SELECT * FROM your_table_name FINAL WHERE __DELETED = 'false';
SELECT * FROM your_table_name FINAL LIMIT 10;
SELECT * FROM your_table_name FINAL WHERE <filter by keys in ORDER BY clause>;
SELECT count(*) FROM your_table_name FINAL;

Для повышения производительности при работе с большими таблицами, особенно если вам не нужно читать все столбцы и если речь идёт о разовых аналитических запросах, можно использовать функцию argMax, чтобы вручную выбрать последнюю запись для каждого первичного ключа:

SELECT key,
       argMax(col1, version) AS col1,
       argMax(col2, version) AS col2
FROM t
WHERE <в

Для production-сценариев и параллельно выполняемых повторяющихся запросов конечных пользователей можно использовать Materialized Views, чтобы организовать данные так, чтобы они лучше соответствовали последующим сценариям доступа.

Дополнительные материалы

Navigation