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
- Войдите в свой аккаунт Streamkap.
- На боковой панели перейдите в раздел Connectors и выберите вкладку Sources.
- Нажмите + Add и выберите тип исходной базы данных (например, SQL Server RDS).
- Заполните сведения о подключении, включая конечную точку, порт, имя базы данных и учетные данные пользователя.
- Сохраните коннектор.
Настройте пункт назначения ClickHouse
- В разделе Connectors выберите вкладку Destinations.
- Нажмите + Add и выберите ClickHouse из списка.
- Введите сведения о подключении для вашего сервиса ClickHouse:
- Hostname: хост вашего экземпляра ClickHouse (например,
abc123.us-west-2.aws.clickhouse.cloud) - Port: защищённый порт HTTPS, обычно
8443 - Username and Password: учетные данные пользователя ClickHouse
- Database: имя целевой базы данных в ClickHouse
- Hostname: хост вашего экземпляра ClickHouse (например,
- Сохраните пункт назначения.
Создайте и запустите конвейер
- Перейдите в Pipelines на боковой панели и нажмите + Create.
- Выберите источник и пункт назначения, которые вы только что настроили.
- Выберите схемы и таблицы, которые хотите передавать.
- Задайте имя конвейеру и нажмите 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, чтобы организовать данные так, чтобы они лучше соответствовали последующим сценариям доступа.