CDC для ClickHouse
Мы рады сообщить о новом коннекторе базы данных ClickHouse для потоковой передачи данных CDC (Change Data Capture) в ClickHouse.
Поскольку ClickHouse идеально подходит для приложений, работающих в режиме реального времени, мы создали высокопроизводительный коннектор ClickHouse, использующий CDC, которые можно масштабировать линейно без ущерба производительности.
Если Вы пока еще не очень хорошо знакомы с концепцией захвата изменения данных (CDC), прочитайте о CDC в контексте потоковой передаче данных, - это поможет Вам понять, будет ли для Вас полезен ClickHouse CDC или нет. Это особенно актуально в свете сообщений о приобретении компанией OpenAI компании Rockset, поскольку компаниям очень часто приходится искать разные варианты CDC, в том числе от MongoDB и DynamoDB.
В этой статье мы рассмотрим некоторые основные функции коннектора, а также его влияние на производительность.
Технологии
ClickHouse - это колоночно-ориентированная база данных с открытым исходным кодом, работающая в режиме реального времени. Структура, ориентированная на столбцы, означает, что данные хранятся и извлекаются по столбцам, а не по строкам. Несмотря на сходство с хранилищами данных, ключевым отличием ClickHouse является создание материализованных представлений при записи, что позволяет ускорить выполнение аналитических запросов в разы. Как раз это и делает ее подходящей для использования в режиме реального времени.
Streamkap - это потоковая платформа обработки данных с высокой пропускной способностью, не требующая особого обслуживания. Идеально подходит для тех случаев, когда Вам нужно действовать быстро, при этом экономить средства, а также заручиться круглосуточной поддержкой, использующей такие хорошо знакомые всем нам технологии, как Apache Kafka, Kafka Connect, Debezium и Apache Flink.
Представляем Вашему вниманию краткий обзор того, как Streamkap осуществляет потоковую передачу данных из базы данных в ClickHouse.
Поддерживаемые типы данных
Наш коннектор поддерживает следующие типы данных:
В настоящее время JSON-поля обрабатываются как строки, использование параметра allow_experimental_object_type=1 находится в стадии тестирования.
Режимы Insert/Upsert
Мы поддерживаем ввод данных в таблицы ClickHouse в режимах Insert и Upsert, при этом режим Upsert является режимом по умолчанию для нашего коннектора.
Режим Insert обеспечивает более высокую пропускную способность и сохраняет исторический набор изменений в таблицах ClickHouse.
Режим Upsert не сохраняет исторические изменения и может использоваться только в тех случаях, когда для обновления целевой таблицы можно использовать первичный ключ из источника.
Режим Insert (Append)
При вставке/добавлении каждое изменение отслеживается и вставляется в ClickHouse как новая строка. Операции удаления в источнике будут отмечены с помощью мета-значения как удаленные __deleted.
Для использования режима Insert (Append) используется движок MergeTree.
Режим Upsert
Upsert - это то, к чему Вы, возможно, привыкли. В этом режиме прекрасно сочетаются и вставки, и обновления. Если есть совпадение по первичному ключу строки, значение будет перезаписано. И наоборот, если совпадения нет, событие будет вставлено.
Режим Upsert реализован с помощью движка ReplacingMergeTree .
Данный механизм удаляет дубликаты данных во время фонового слияния на основе ключа упорядочивания, что позволяет очищать старые записи.
Пример Upsert с базовыми типами
В данном случае Upsert выполнен в формате JSON. Ключ имеет только одно поле `id`, которое является первичным ключом, по которому будут дедуплицироваться строки.
Результирующая таблица:
Данные:
Дедуплицированные данные с использованием FINAL:
Создание моментальных снимков
Создание моментальных снимков - это процесс загрузки существующих данных из базы данных в ClickHouse. Заполнение данных осуществляется с помощью Select, который в отличие от потокового режима считывает данные из журнала базы данных.
По умолчанию Streamkap использует инкрементные снимки. Такой метод подходит для больших таблиц, он практически никак не влияет на исходную базу данных. Благодаря водяному знаку процесс создания моментальных снимков можно возобновить с того места, где он был прерван.
Метаданные
Streamkap добавляет дополнительные столбцы метаданных к каждой вставке в таблицу ClickHouse, что делает анализ данных еще более эффективным.
В каждую таблицу ClickHouse добавляются следующие столбцы метаданных:
- _streamkap_ts_ms: временная метка CDC
- _streamkap_deleted: если текущее событие CDC является событием удаления, то для режима «upsert» для ReplacingMergeTree используется вычисляемый столбец secod streamkap deleted типа UInt8
- _streamkap_partition: smallint, представляющий внутренний номер раздела Streamkap, полученный путем последовательного хэширования ключевых полей исходных записей
- _streamkap_source_ts_ms: метка времени, когда произошло событие изменения в исходной базе данных
- -streamkap_op: тип операции события CDC (c insert, u update, d delete, r snapshot, t truncate)
Работа с полуструктурированными данными
Вложенные массивы и структуры
Ниже мы приводим несколько примеров того, как сложные структуры автоматически переводятся на типы ClickHouse.
Для поддержки массивов, содержащих различные структуры, Streamkap в ClickHouse нужно изменить на следующее значение, а flatten_nested - на 0:
ALTER ROLE STREAMKAP_ROLE SETTINGS flatten_nested = 0;
Поле Struct, содержащее подмассив
Здесь показана входная запись в формате JSON, где ключ имеет только одно поле id:
Результирующая таблица. Не видно, как для обработки сложной структуры в `Tuple(nb Int32, str String, sub_arr Array(Tuple(n Int32, s String)), sub_arr_str Array(String)) был отображен столбец `obj` `:
Данные:
Поле вложенного массива, содержащее подструктуру
Здесь показана входная запись в формате JSON, в которой ключ имеет только одно поле id:
Снова результирующая таблица, в которой столбец `arr` отображен на `Array(Tuple(nb Int32, str String))`.
Данные:
Гарантия согласованности и доставки данных
Streamkap гарантирует доставку данных ClickHouse не менее одного раза и по умолчанию использует режим upsert.
Для режима Insert/Append это может привести к тому, что в ClickHouse будут вставлены дополнительные дубликаты строк, но материализованные представления ClickHouse смогут вовремя их отфильтровать.
В режиме Upsert, используемом по умолчанию, мы выполняем дедупликацию по ключу исходной записи.
Преобразования
Streamkap поддерживает преобразования в конвейере, так что данные могут быть отправлены в ClickHouse уже предварительно обработанными. Это осуществляется с помощью Apache Flink, который считывает данные из темы Kafka (вставляются неизменяемые данные), преобразует их внутри Flink и записывает обратно в новую тему для вставки в ClickHouse.
Это особенно полезно для полуструктурированных данных, предварительной обработки и задач очистки. Такой метод может быть значительно эффективнее, чем работа с данными после обработки.
Ниже мы приводим некоторые наиболее распространенные преобразования, выполняемые Streamkap.
Устранение несоответствий в полуструктурированных данных
Рассмотрим исправление несогласованного полуструктурированного поля даты:
С помощью преобразований Streamkap все записи могут быть преобразованы в один общий формат, подходящий для столбца DateTime64:
Разделение больших полуструктурированных документов JSON
В базах данных документов дочерние сущности могут быть смоделированы как подмассивы, вложенные в документ родительской сущности.
В ClickHouse эти дочерние сущности целесообразно представлять в виде отдельных строк. Используя преобразования Streamkap, записи дочерних сущностей можно разделить на отдельные записи следующим образом:
Эволюция схемы
Эволюция схемы или обработка дрейфа - это процесс внесения изменений в целевые таблицы для отражения исходных изменений, например, добавления/удаления дополнительных столбцов.
Коннектор Streamkap автоматически обрабатывает дрейф схемы в ClickHouse в следующих случаях:
- Дополнительные столбцы: Будет обнаружено дополнительное поле и затем создан новый столбец , предназначенный для новых данных.
- Удаление столбцов: Этот столбец теперь будет игнорироваться, никаких дальнейших действий по отношению к нему предприниматься не будет.
- Изменение типа столбца: В таблице создается дополнительный столбец с суффиксом, обозначающим новый тип, например, ColumnName_type
Дополнительные таблицы могут быть добавлены в конвейер на любом этапе.
Ниже мы приводим несколько примеров такой эволюции схемы.
Добавление столбца
Рассмотрим следующую входную запись до эволюции схемы:
Новый столбец `new_double_col` добавлен в вышестоящую схему, что приводит к изменению схемы ClickHouse:
Данные ClickHouse:
Преобразование Int в String
Входная запись перед эволюцией схемы:
Новая запись попадает в систему :
Данные ClickHouse после добавления нового столбца IntColumn_str:
Производительность
ClickHouse предназначен для использования в режиме реального времени, поэтому мы сделали так, чтобы при увеличении ресурсов наш коннектор мог линейно масштабироваться. Ниже мы продемонстрировали линейную производительность до 85 тыс. записей CDC в секунду в режиме upsert, но мы уверены, что она может масштабироваться настолько сильно, насколько Вам нужно.
Для нашего теста мы использовали экземпляр кластера Clickhouse, состоящего из 3 узлов по 32 Гб каждый с 8 vCPU.
Формат входных записей содержит основные типы, средний размер строки - ~100 символов, большая строка содержит примерно 1000 символов. Мы использовали режим upsert, который по сравнению с режимом insert будет менее производительным.
Baseline единичная партиция
Baseline с одной задачей Streamkap и разделом Clickhouse с несколькими объемами.
Производительность:
Задержка в зависимости от объема:
В случае моментальных снимков/обратных заполнений имеет смысл использовать объемы более 100 000 записей, и это автоматически оптимизируется в Streamkap.
Для потокового режима обычно желательно использовать меньшие размеры массива, но это опять же автоматически оптимизируется.
Это лишь некоторые тесты с фиксированным размером массива, проведенные для того, чтобы продемонстрировать компромисс между пропускной способностью и задержкой. На практике размер массива меняется в зависимости от размера внутренней очереди, и Streamkap всегда его автоматически оптимизирует.
Масштабирование
В данном случае мы протестировали 100 000 записей и постепенно увеличивали количество задач: 1, 2, 4 и 8. В результате мы видим, что пропускная способность линейно зависит от количества задач.
Заключение
Streamkap создал самый высокопроизводительный CDC-коннектор для ClickHouse и добавил в него множество полезных функций.
































