Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Шаблон Dataflow для передачи из Pub/Sub в ClickHouse

Шаблон 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.

Чтобы сопоставить входящие сообщения с целевой таблицей, конвейер при запуске выполняет следующие действия:

  1. Получает схему целевой таблицы ClickHouse.
  2. Создает схему Beam Row на основе схемы ClickHouse.
  3. Для каждого входящего сообщения 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.

  1. Нажмите кнопку CREATE JOB FROM TEMPLATE.

    Консоль Dataflow
  2. Когда откроется форма шаблона, введите имя задачи и выберите нужный регион.

  3. В поле Dataflow Template введите ClickHouse или Pub/Sub и выберите шаблон Pub/Sub to ClickHouse.

  4. После выбора форма развернётся. Заполните:

    • входную подписку Pub/Sub в формате projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>.
    • URL конечной точки ClickHouse — для ClickHouse Cloud используйте https://<HOST>:8443.
    • базу данных ClickHouse, целевую таблицу, имя пользователя и пароль.
    • как минимум один пункт назначения dead-letter: таблицу ClickHouse или топик Pub/Sub (или оба варианта).
  5. При необходимости настройте батчинг (windowSeconds, batchRowCount) и параметры тонкой настройки ClickHouseIO, как подробно описано в разделе Параметры шаблона.

Отслеживание задачи

Перейдите на вкладку задач Dataflow в Google Cloud Console, чтобы отслеживать состояние задачи. Там вы найдете сведения о задаче, включая ход выполнения и возможные ошибки:

Консоль Dataflow с выполняющейся задачей Pub/Sub to ClickHouse

Шаблон также отправляет следующие пользовательские метрики в пространстве имен 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).

Исходный код шаблона

Исходный код шаблона доступен в следующих репозиториях:

Navigation