Шаблон Pub/Sub to ClickHouse — это стриминговый конвейер, который читает сообщения в формате JSON из подписки Pub/Sub и записывает их в таблицу ClickHouse. Сообщения, которые не удалось разобрать или сопоставить с целевой схемой, направляются в пункт назначения dead-letter: таблицу ClickHouse, тему Pub/Sub или в оба сразу.
Требования к конвейеру
- Исходная подписка Pub/Sub должна существовать.
- Сообщения, публикуемые в подписку, должны быть корректным JSON.
- Целевая таблица ClickHouse должна существовать, а имена её столбцов должны совпадать с именами полей в полезной нагрузке JSON.
- Хост ClickHouse должен быть доступен с машин воркеров Dataflow.
- Должен быть указан как минимум один пункт назначения dead-letter (
clickHouseDeadLetterTableилиdeadLetterTopic). Если указаны оба, сообщения с ошибками направляются в оба пункта назначения одновременно. - Если задан
clickHouseDeadLetterTable, таблица dead-letter уже должна существовать в ClickHouse со схемой, показанной в разделе Обработка dead-letter. - Если задан
deadLetterTopic, топик Pub/Sub уже должен существовать.
Параметры шаблона
| Название параметра | Описание параметра | Обязательно | Примечания |
|---|---|---|---|
inputSubscription |
Подписка Pub/Sub, из которой считываются сообщения. Пример: projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>. |
✅ | Сообщения должны быть закодированы в JSON. |
clickHouseUrl |
URL конечной точки ClickHouse. Используйте https:// для SSL-соединений (ClickHouse Cloud) или http:// для соединений без SSL. Пример: https://<HOST>:8443 или http://<HOST>:8123. |
✅ | Для ClickHouse Cloud используйте конечную точку HTTPS на порту 8443. |
clickHouseDatabase |
Имя базы данных ClickHouse, в которой находится целевая таблица. Пример: default. |
✅ | |
clickHouseTable |
Имя таблицы ClickHouse, в которую записываются данные. | ✅ | Таблица должна существовать до запуска конвейера. |
clickHouseUsername |
Имя пользователя для аутентификации в ClickHouse. | ✅ | |
clickHousePassword |
Пароль для аутентификации в ClickHouse. | ✅ | |
clickHouseDeadLetterTable |
Таблица ClickHouse, в которую записываются сообщения с ошибками. Пример: my_table_dead_letter. |
Должен быть указан как минимум один из параметров clickHouseDeadLetterTable или deadLetterTopic. Таблица должна существовать со схемой dead-letter, показанной в разделе Обработка dead-letter. |
|
deadLetterTopic |
Топик Pub/Sub, в который публикуются сообщения с ошибками. Пример: projects/<PROJECT_ID>/topics/<TOPIC_NAME>. |
Должен быть указан как минимум один из параметров clickHouseDeadLetterTable или deadLetterTopic. Полезные нагрузки сообщений с ошибкой публикуются в топик с errorMessage и failedAt, заданными как атрибуты сообщения. |
|
windowSeconds |
Длительность временных окон батчинга в секундах. | См. раздел Batching and windowing о взаимодействии с batchRowCount. Если не задан ни один из параметров, в комбинированном режиме используются значения по умолчанию 30s и 1000 строк. |
|
batchRowCount |
Количество строк, накапливаемых перед сбросом в ClickHouse. | См. раздел Batching and windowing о взаимодействии с windowSeconds. |
|
maxInsertBlockSize |
Максимальное количество строк в одном операторе INSERT, отправляемом в ClickHouse. По умолчанию — 1,000,000. |
Опция ClickHouseIO. |
|
maxRetries |
Максимальное количество повторных попыток для неудачных вставок в ClickHouse. По умолчанию — 5. |
Опция ClickHouseIO. |
|
insertDeduplicate |
Включать ли дедупликацию для запросов INSERT в реплицируемых таблицах ClickHouse. По умолчанию — true. |
Опция ClickHouseIO. |
|
insertQuorum |
Для запросов INSERT в реплицируемых таблицах ожидать, пока указанное количество реплик не подтвердит запись и не линеаризует добавление данных. 0 отключает запись по quorum. |
Опция ClickHouseIO. Отключено в настройках сервера по умолчанию. |
|
insertDistributedSync |
Если параметр включен, запросы INSERT в distributed таблицы ожидают, пока данные не будут отправлены на все узлы кластера. По умолчанию — true. |
Опция ClickHouseIO. |
Формат сообщения и сопоставление со схемой
Сообщения Pub/Sub должны быть объектами JSON, в которых имена полей верхнего уровня в точности совпадают с именами столбцов целевой таблицы ClickHouse.
Чтобы сопоставить входящие сообщения с целевой таблицей, конвейер при запуске выполняет следующие действия:
- Получает схему целевой таблицы ClickHouse.
- Создает схему Beam
Rowна основе схемы ClickHouse. - Для каждого входящего сообщения Pub/Sub разбирает полезную нагрузку JSON и формирует строку, считывая поля с именами из схемы ClickHouse.
Преобразование типов
Значения JSON приводятся к соответствующему типу столбца ClickHouse:
| Тип ClickHouse | Примечания |
|---|---|
Float32 |
Разбирается с помощью Float.valueOf. |
Float64 |
Разбирается с помощью Double.valueOf. |
Date |
Разбирается как строка даты в формате ISO-8601. |
DateTime |
Разбирается как строка даты и времени в формате ISO-8601 (например, 2026-01-15T12:34:56Z). |
Array(T) |
JSON-массив; каждый элемент преобразуется к типу элемента T. Для пустых или отсутствующих массивов возвращается пустой массив. |
Integer types (Int8/Int16/Int32/Int64, UInt8/UInt16/UInt32/UInt64) |
Разбираются из числового значения JSON или его строкового представления. |
String |
Для текстовых полей используется как есть; нетекстовые узлы JSON сериализуются в строковое представление JSON. |
Батчинг и оконная обработка
Поскольку конвейер работает в режиме стриминга, входящие строки накапливаются в окнах перед записью в ClickHouse. Стратегия оконной обработки выбирается на основе указанных вами параметров:
windowSeconds |
batchRowCount |
Поведение |
|---|---|---|
| задано | не задано | Фиксированные окна по времени длительностью windowSeconds. |
| не задано | задано | Глобальное окно с триггером по количеству строк; срабатывает каждые batchRowCount строк. |
| заданы оба | заданы оба | Глобальное окно с комбинированным триггером; срабатывает при выполнении первого из условий (время или количество строк). |
| не задан ни один | не задан ни один | Комбинированный режим со значениями по умолчанию: 30 секунд или 1000 строк — в зависимости от того, что наступит раньше. |
Подбирая эти значения, вы находите баланс между задержкой и эффективностью вставки. Меньшие окна снижают сквозную задержку; большие окна дают меньшее число более крупных батчей INSERT.
Обработка dead-letter
Сообщения, для которых не удалось выполнить парсинг JSON, сопоставление со схемой или приведение типов, направляются в настроенные пункты назначения dead-letter. Необходимо указать как минимум один из параметров: clickHouseDeadLetterTable или deadLetterTopic; если заданы оба, сообщения с ошибками будут отправлены в оба.
Таблица ClickHouse dead-letter
Если задан параметр clickHouseDeadLetterTable, таблица dead-letter уже должна существовать со следующей фиксированной схемой:
| Столбец | Тип | Описание |
|---|---|---|
raw_message |
String |
Исходная полезная нагрузка сообщения Pub/Sub в виде текста UTF-8. |
error_message |
String |
Сообщение об исключении, объясняющее, почему не удалось обработать строку. |
stack_trace |
String |
Полная трассировка стека Java, зафиксированная в момент сбоя. |
failed_at |
DateTime |
Временная метка момента обработки, в который произошел сбой строки. |
Минимальное определение для одновузлового развертывания:
CREATE TABLE my_table_dead_letter (
raw_message String,
error_message String,
stack_trace String,
failed_at DateTime
) ENGINE = MergeTree()
ORDER BY failed_at;dead-letter-топик Pub/Sub
Если задан deadLetterTopic, каждое сообщение, обработка которого завершилась ошибкой, повторно публикуется в топик со следующим содержимым:
- Полезная нагрузка: исходные байты сообщения.
- Атрибут
errorMessage: сообщение исключения, зафиксированное в момент сбоя. - Атрибут
failedAt: временная метка времени обработки, соответствующая моменту сбоя строки.
Это позволяет удобно повторно обработать сообщения с ошибками после устранения проблемы в схеме или у продьюсера.
Запуск шаблона
Шаблон Pub/Sub to ClickHouse доступен в Google Cloud Console.
Войдите в Google Cloud Console и найдите Dataflow.
-
Нажмите кнопку
CREATE JOB FROM TEMPLATE.
-
Когда откроется форма шаблона, введите имя задачи и выберите нужный регион.
-
В поле
Dataflow TemplateвведитеClickHouseилиPub/Subи выберите шаблонPub/Sub to ClickHouse. -
После выбора форма развернётся. Заполните:
- входную подписку Pub/Sub в формате
projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>. - URL конечной точки ClickHouse — для ClickHouse Cloud используйте
https://<HOST>:8443. - базу данных ClickHouse, целевую таблицу, имя пользователя и пароль.
- как минимум один пункт назначения dead-letter: таблицу ClickHouse или топик Pub/Sub (или оба варианта).
- входную подписку Pub/Sub в формате
-
При необходимости настройте батчинг (
windowSeconds,batchRowCount) и параметры тонкой настройкиClickHouseIO, как подробно описано в разделе Параметры шаблона.
Отслеживание задачи
Перейдите на вкладку задач Dataflow в Google Cloud Console, чтобы отслеживать состояние задачи. Там вы найдете сведения о задаче, включая ход выполнения и возможные ошибки:

Шаблон также отправляет следующие пользовательские метрики в пространстве имен PubSubToClickHouse; их можно просмотреть на странице задачи Dataflow:
| Метрика | Type | Описание |
|---|---|---|
messages-received |
Counter | Общее количество сообщений Pub/Sub, полученных на этапе разбора. |
rows-parsed-ok |
Counter | Сообщения, успешно преобразованные в строку и направленные в основной выход. |
rows-parse-failed |
Counter | Сообщения, для которых не удалось выполнить разбор или сопоставление со схемой и которые были направлены в dead-letter. |
message-payload-bytes |
Distribution | Распределение размеров полезной нагрузки входящих сообщений Pub/Sub в байтах. |
Устранение неполадок
Ошибка превышения общего лимита памяти (код 241)
Эта ошибка возникает, когда ClickHouse не хватает памяти при обработке больших батчей данных. Чтобы устранить проблему:
- Увеличьте ресурсы инстанса: переведите ClickHouse server на более крупный инстанс с большим объёмом памяти, чтобы он справлялся с нагрузкой при обработке данных.
- Уменьшите размер батча: сократите
batchRowCount(и/илиmaxInsertBlockSize) в конфигурации задачи Dataflow, чтобы отправлять в ClickHouse меньшие фрагменты данных и снизить потребление памяти на батч.
Все сообщения отправляются в пункт назначения dead-letter
Наиболее распространённые причины:
- Имена JSON-полей не совпадают в точности с именами столбцов ClickHouse (сопоставление чувствительно к регистру).
- Значение JSON невозможно привести к типу столбца (например, строку не в формате ISO-8601 в столбце
DateTime). - Схема целевой таблицы изменилась после запуска конвейера — схема загружается один раз при запуске. Перезапустите задачу после внесения изменений в схему.
Проверьте столбцы error_message и stack_trace в таблице dead-letter ClickHouse (или атрибут errorMessage в сообщениях Pub/Sub dead-letter), чтобы определить первопричину.
Конвейер запускается, но строки не поступают в ClickHouse
- Убедитесь, что подписка получает сообщения — проверьте метрику
messages-receivedна странице задачи Dataflow. - В режиме по времени (только
windowSeconds) строки сбрасываются на диск только на границах окна. УменьшитеwindowSeconds, чтобы проверить, происходят ли сбросы. - Проверьте сетевую доступность между воркерами Dataflow и конечной точкой ClickHouse (брандмауэр, пиринг VPC или Private Service Connect).
Исходный код шаблона
Исходный код шаблона доступен в следующих репозиториях:
GoogleCloudPlatform/DataflowTemplates— основной репозиторий Google Cloud Platform.ClickHouse/DataflowTemplates— форк ClickHouse.