> ## Documentation Index
> Fetch the complete documentation index at: https://clickhouse.com/docs/llms.txt
> Use this file to discover all available pages before exploring further.

> Вы можете потоково передавать JSON-сообщения из Pub/Sub в ClickHouse с помощью шаблона Google Dataflow

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

export const Image = ({img, alt, size = "lg"}) => {
  const normalizedSize = ["sm", "md", "lg"].includes(size) ? size : "lg";
  return <div className={`ch-image-${normalizedSize}`}>
      <Frame>
        <img src={img} alt={alt} />
      </Frame>
    </div>;
};

Шаблон Pub/Sub to ClickHouse — это стриминговый конвейер, который читает сообщения в формате JSON из подписки Pub/Sub и записывает их в таблицу ClickHouse.
Сообщения, которые не удалось разобрать или сопоставить с целевой схемой, направляются в пункт назначения dead-letter: таблицу ClickHouse, тему Pub/Sub или в оба сразу.

<div id="pipeline-requirements">
  ## Требования к конвейеру
</div>

* Исходная подписка Pub/Sub должна существовать.
* Сообщения, публикуемые в подписку, должны быть корректным JSON.
* Целевая таблица ClickHouse должна существовать, а имена её столбцов должны совпадать с именами полей в полезной нагрузке JSON.
* Хост ClickHouse должен быть доступен с машин воркеров Dataflow.
* Должен быть указан как минимум один пункт назначения dead-letter (`clickHouseDeadLetterTable` или `deadLetterTopic`). Если указаны оба, сообщения с ошибками направляются в оба пункта назначения одновременно.
* Если задан `clickHouseDeadLetterTable`, таблица dead-letter уже должна существовать в ClickHouse со схемой, показанной в разделе [Обработка dead-letter](#dead-letter-handling).
* Если задан `deadLetterTopic`, топик Pub/Sub уже должен существовать.

<div id="template-parameters">
  ## Параметры шаблона
</div>

<br />

<br />

| Название параметра          | Описание параметра                                                                                                                                                                        | Обязательно | Примечания                                                                                                                                                                                                                      |
| --------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `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](#dead-letter-handling).    |
| `deadLetterTopic`           | Топик Pub/Sub, в который публикуются сообщения с ошибками. Пример: `projects/<PROJECT_ID>/topics/<TOPIC_NAME>`.                                                                           |             | Должен быть указан как минимум один из параметров `clickHouseDeadLetterTable` или `deadLetterTopic`. Полезные нагрузки сообщений с ошибкой публикуются в топик с `errorMessage` и `failedAt`, заданными как атрибуты сообщения. |
| `windowSeconds`             | Длительность временных окон батчинга в секундах.                                                                                                                                          |             | См. раздел [Batching and windowing](#batching-and-windowing) о взаимодействии с `batchRowCount`. Если не задан ни один из параметров, в комбинированном режиме используются значения по умолчанию `30s` и `1000` строк.         |
| `batchRowCount`             | Количество строк, накапливаемых перед сбросом в ClickHouse.                                                                                                                               |             | См. раздел [Batching and windowing](#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`.                                                                                                                                                                                                           |

<Note>
  Значения по умолчанию для всех параметров `ClickHouseIO` можно найти в разделе [`ClickHouseIO` Apache Beam Connector](/docs/ru/integrations/connectors/data-ingestion/etl-tools/apache-beam#clickhouseiowrite-parameters).
</Note>

<div id="message-format-and-schema-mapping">
  ## Формат сообщения и сопоставление со схемой
</div>

Сообщения Pub/Sub должны быть объектами JSON, в которых имена полей верхнего уровня в точности совпадают с именами столбцов целевой таблицы ClickHouse.

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

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

<br />

<Warning>
  Имена полей JSON должны в точности совпадать с именами столбцов ClickHouse (сопоставление чувствительно к регистру). Поля в сообщении, которые не соответствуют столбцам ClickHouse, игнорируются. Если для столбца ClickHouse в полезной нагрузке JSON нет соответствующего поля, конвейер пытается записать в этот столбец `NULL` — что возможно только в том случае, если столбец объявлен как [`Nullable`](/docs/ru/reference/data-types/nullable). Сообщения, которые не удается разобрать, значения которых нельзя привести к типу столбца или которые привели бы к записи `NULL` в столбец, не допускающий NULL, направляются в пункт назначения dead-letter.
</Warning>

<div id="type-conversion">
  ### Преобразование типов
</div>

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

| Тип ClickHouse                                                                     | Примечания                                                                                                                       |
| ---------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------- |
| [`Float32`](/docs/ru/reference/data-types/float)                                        | Разбирается с помощью `Float.valueOf`.                                                                                           |
| [`Float64`](/docs/ru/reference/data-types/float)                                        | Разбирается с помощью `Double.valueOf`.                                                                                          |
| [`Date`](/docs/ru/reference/data-types/date)                                            | Разбирается как строка даты в формате ISO-8601.                                                                                  |
| [`DateTime`](/docs/ru/reference/data-types/datetime)                                    | Разбирается как строка даты и времени в формате ISO-8601 (например, `2026-01-15T12:34:56Z`).                                     |
| [`Array(T)`](/docs/ru/reference/data-types/array)                                       | JSON-массив; каждый элемент преобразуется к типу элемента `T`. Для пустых или отсутствующих массивов возвращается пустой массив. |
| Integer types (`Int8`/`Int16`/`Int32`/`Int64`, `UInt8`/`UInt16`/`UInt32`/`UInt64`) | Разбираются из числового значения JSON или его строкового представления.                                                         |
| [`String`](/docs/ru/reference/data-types/string)                                        | Для текстовых полей используется как есть; нетекстовые узлы JSON сериализуются в строковое представление JSON.                   |

<div id="batching-and-windowing">
  ## Батчинг и оконная обработка
</div>

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

| `windowSeconds`  | `batchRowCount`  | Поведение                                                                                                                    |
| ---------------- | ---------------- | ---------------------------------------------------------------------------------------------------------------------------- |
| задано           | не задано        | Фиксированные окна по времени длительностью `windowSeconds`.                                                                 |
| не задано        | задано           | Глобальное окно с триггером по количеству строк; срабатывает каждые `batchRowCount` строк.                                   |
| заданы оба       | заданы оба       | Глобальное окно с комбинированным триггером; срабатывает при выполнении первого из условий (время **или** количество строк). |
| не задан ни один | не задан ни один | Комбинированный режим со значениями по умолчанию: `30` секунд или `1000` строк — в зависимости от того, что наступит раньше. |

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

<div id="dead-letter-handling">
  ## Обработка dead-letter
</div>

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

<div id="clickhouse-dead-letter-table">
  ### Таблица ClickHouse dead-letter
</div>

Если задан параметр `clickHouseDeadLetterTable`, таблица dead-letter уже должна существовать со следующей фиксированной схемой:

| Столбец         | Тип        | Описание                                                                   |
| --------------- | ---------- | -------------------------------------------------------------------------- |
| `raw_message`   | `String`   | Исходная полезная нагрузка сообщения Pub/Sub в виде текста UTF-8.          |
| `error_message` | `String`   | Сообщение об исключении, объясняющее, почему не удалось обработать строку. |
| `stack_trace`   | `String`   | Полная трассировка стека Java, зафиксированная в момент сбоя.              |
| `failed_at`     | `DateTime` | Временная метка момента обработки, в который произошел сбой строки.        |

Минимальное определение для одновузлового развертывания:

```sql theme={null}
CREATE TABLE my_table_dead_letter (
    raw_message   String,
    error_message String,
    stack_trace   String,
    failed_at     DateTime
) ENGINE = MergeTree()
ORDER BY failed_at;
```

<Note>
  Адаптируйте движок и предложение `ORDER BY` под своё развертывание — используйте `ReplicatedMergeTree` для реплицируемых таблиц, добавьте `ON CLUSTER` для распределённых развертываний и при необходимости настройте партиционирование или TTL.
</Note>

<div id="pubsub-dead-letter-topic">
  ### dead-letter-топик Pub/Sub
</div>

Если задан `deadLetterTopic`, каждое сообщение, обработка которого завершилась ошибкой, повторно публикуется в топик со следующим содержимым:

* **Полезная нагрузка**: исходные байты сообщения.
* **Атрибут** `errorMessage`: сообщение исключения, зафиксированное в момент сбоя.
* **Атрибут** `failedAt`: временная метка времени обработки, соответствующая моменту сбоя строки.

Это позволяет удобно повторно обработать сообщения с ошибками после устранения проблемы в схеме или у продьюсера.

<div id="running-the-template">
  ## Запуск шаблона
</div>

Шаблон `Pub/Sub to ClickHouse` доступен в Google Cloud Console.

<Note>
  Обязательно ознакомьтесь с этим документом, особенно с разделами выше, чтобы полностью понять требования к конфигурации шаблона и необходимые предварительные условия.
</Note>

Войдите в Google Cloud Console и найдите Dataflow.

1. Нажмите кнопку `CREATE JOB FROM TEMPLATE`.
   <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/google-dataflow/create_job_from_template_button.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=ca429a13d8a9e99c43ae477bf14ad1a9" border alt="Консоль Dataflow" width="1872" height="886" data-path="images/integrations/data-ingestion/google-dataflow/create_job_from_template_button.webp" />

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`, как подробно описано в разделе [Параметры шаблона](#template-parameters).

<div id="monitor-the-job">
  ### Отслеживание задачи
</div>

Перейдите на [вкладку задач Dataflow](https://console.cloud.google.com/dataflow/jobs) в Google Cloud Console, чтобы отслеживать состояние задачи. Там вы найдете сведения о задаче, включая ход выполнения и возможные ошибки:

<Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/google-dataflow/pubsub-inqueue-job.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=c5922b3ad406648be710f93d856f5fe8" size="lg" border alt="Консоль Dataflow с выполняющейся задачей Pub/Sub to ClickHouse" width="3016" height="470" data-path="images/integrations/data-ingestion/google-dataflow/pubsub-inqueue-job.webp" />

Шаблон также отправляет следующие пользовательские метрики в пространстве имен `PubSubToClickHouse`; их можно просмотреть на странице задачи Dataflow:

| Метрика                 | Type         | Описание                                                                                                                |
| ----------------------- | ------------ | ----------------------------------------------------------------------------------------------------------------------- |
| `messages-received`     | Counter      | Общее количество сообщений Pub/Sub, полученных на этапе разбора.                                                        |
| `rows-parsed-ok`        | Counter      | Сообщения, успешно преобразованные в строку и направленные в основной выход.                                            |
| `rows-parse-failed`     | Counter      | Сообщения, для которых не удалось выполнить разбор или сопоставление со схемой и которые были направлены в dead-letter. |
| `message-payload-bytes` | Distribution | Распределение размеров полезной нагрузки входящих сообщений Pub/Sub в байтах.                                           |

<div id="troubleshooting">
  ## Устранение неполадок
</div>

<div id="code-241-dbexception-memory-limit-total-exceeded">
  ### Ошибка превышения общего лимита памяти (код 241)
</div>

Эта ошибка возникает, когда ClickHouse не хватает памяти при обработке больших батчей данных. Чтобы устранить проблему:

* Увеличьте ресурсы инстанса: переведите ClickHouse server на более крупный инстанс с большим объёмом памяти, чтобы он справлялся с нагрузкой при обработке данных.
* Уменьшите размер батча: сократите `batchRowCount` (и/или `maxInsertBlockSize`) в конфигурации задачи Dataflow, чтобы отправлять в ClickHouse меньшие фрагменты данных и снизить потребление памяти на батч.

<div id="all-messages-going-to-dlq">
  ### Все сообщения отправляются в пункт назначения dead-letter
</div>

Наиболее распространённые причины:

* Имена JSON-полей не совпадают в точности с именами столбцов ClickHouse (сопоставление чувствительно к регистру).
* Значение JSON невозможно привести к типу столбца (например, строку не в формате ISO-8601 в столбце `DateTime`).
* Схема целевой таблицы изменилась после запуска конвейера — схема загружается один раз при запуске. Перезапустите задачу после внесения изменений в схему.

Проверьте столбцы `error_message` и `stack_trace` в таблице dead-letter ClickHouse (или атрибут `errorMessage` в сообщениях Pub/Sub dead-letter), чтобы определить первопричину.

<div id="no-rows-arriving">
  ### Конвейер запускается, но строки не поступают в ClickHouse
</div>

* Убедитесь, что подписка получает сообщения — проверьте метрику `messages-received` на странице задачи Dataflow.
* В режиме по времени (только `windowSeconds`) строки сбрасываются на диск только на границах окна. Уменьшите `windowSeconds`, чтобы проверить, происходят ли сбросы.
* Проверьте сетевую доступность между воркерами Dataflow и конечной точкой ClickHouse (брандмауэр, пиринг VPC или Private Service Connect).

<div id="template-source-code">
  ## Исходный код шаблона
</div>

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

* [`GoogleCloudPlatform/DataflowTemplates`](https://github.com/GoogleCloudPlatform/DataflowTemplates/tree/main/v2/googlecloud-to-clickhouse) — основной репозиторий Google Cloud Platform.
* [`ClickHouse/DataflowTemplates`](https://github.com/ClickHouse/DataflowTemplates) — форк ClickHouse.
