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

Параметры шаблона



Значения по умолчанию для всех параметров ClickHouseIO можно найти в разделе ClickHouseIO Apache Beam Connector.

Формат сообщения и сопоставление со схемой

Сообщения Pub/Sub должны быть объектами JSON, в которых имена полей верхнего уровня в точности совпадают с именами столбцов целевой таблицы ClickHouse. Чтобы сопоставить входящие сообщения с целевой таблицей, конвейер при запуске выполняет следующие действия:
  1. Получает схему целевой таблицы ClickHouse.
  2. Создает схему Beam Row на основе схемы ClickHouse.
  3. Для каждого входящего сообщения Pub/Sub разбирает полезную нагрузку JSON и формирует строку, считывая поля с именами из схемы ClickHouse.

Имена полей JSON должны в точности совпадать с именами столбцов ClickHouse (сопоставление чувствительно к регистру). Поля в сообщении, которые не соответствуют столбцам ClickHouse, игнорируются. Если для столбца ClickHouse в полезной нагрузке JSON нет соответствующего поля, конвейер пытается записать в этот столбец NULL — что возможно только в том случае, если столбец объявлен как Nullable. Сообщения, которые не удается разобрать, значения которых нельзя привести к типу столбца или которые привели бы к записи NULL в столбец, не допускающий NULL, направляются в пункт назначения dead-letter.

Преобразование типов

Значения JSON приводятся к соответствующему типу столбца ClickHouse:

Батчинг и оконная обработка

Поскольку конвейер работает в режиме стриминга, входящие строки накапливаются в окнах перед записью в ClickHouse. Стратегия оконной обработки выбирается на основе указанных вами параметров: Подбирая эти значения, вы находите баланс между задержкой и эффективностью вставки. Меньшие окна снижают сквозную задержку; большие окна дают меньшее число более крупных батчей INSERT.

Обработка dead-letter

Сообщения, для которых не удалось выполнить парсинг JSON, сопоставление со схемой или приведение типов, направляются в настроенные пункты назначения dead-letter. Необходимо указать как минимум один из параметров: clickHouseDeadLetterTable или deadLetterTopic; если заданы оба, сообщения с ошибками будут отправлены в оба.

Таблица ClickHouse dead-letter

Если задан параметр clickHouseDeadLetterTable, таблица dead-letter уже должна существовать со следующей фиксированной схемой: Минимальное определение для одновузлового развертывания:
Адаптируйте движок и предложение ORDER BY под своё развертывание — используйте ReplicatedMergeTree для реплицируемых таблиц, добавьте ON CLUSTER для распределённых развертываний и при необходимости настройте партиционирование или TTL.

dead-letter-топик Pub/Sub

Если задан deadLetterTopic, каждое сообщение, обработка которого завершилась ошибкой, повторно публикуется в топик со следующим содержимым:
  • Полезная нагрузка: исходные байты сообщения.
  • Атрибут errorMessage: сообщение исключения, зафиксированное в момент сбоя.
  • Атрибут failedAt: временная метка времени обработки, соответствующая моменту сбоя строки.
Это позволяет удобно повторно обработать сообщения с ошибками после устранения проблемы в схеме или у продьюсера.

Запуск шаблона

Шаблон Pub/Sub to ClickHouse доступен в Google Cloud Console.
Обязательно ознакомьтесь с этим документом, особенно с разделами выше, чтобы полностью понять требования к конфигурации шаблона и необходимые предварительные условия.
Войдите в Google Cloud Console и найдите Dataflow.
  1. Нажмите кнопку CREATE JOB FROM TEMPLATE.
  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, чтобы отслеживать состояние задачи. Там вы найдете сведения о задаче, включая ход выполнения и возможные ошибки: Шаблон также отправляет следующие пользовательские метрики в пространстве имен PubSubToClickHouse; их можно просмотреть на странице задания Dataflow:

Устранение неполадок

Ошибка превышения общего лимита памяти (код 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).

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

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