> ## 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.

> Você pode transmitir mensagens JSON do Pub/Sub para o ClickHouse usando um template do Google Dataflow

# Template do Dataflow de Pub/Sub para 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>;
};

O template do Pub/Sub para ClickHouse é um pipeline de streaming que lê mensagens codificadas em JSON de uma assinatura do Pub/Sub e as grava em uma tabela do ClickHouse.
Mensagens cuja análise falha ou que não podem ser mapeadas para o esquema de destino são encaminhadas para um destino dead-letter: uma tabela do ClickHouse, um tópico do Pub/Sub ou ambos.

<div id="pipeline-requirements">
  ## Requisitos do pipeline
</div>

* A assinatura Pub/Sub de origem deve existir.
* As mensagens publicadas na assinatura devem ser JSON válido.
* A tabela ClickHouse de destino deve existir, e os nomes de suas colunas devem corresponder aos nomes dos campos no payload JSON.
* O host do ClickHouse deve estar acessível a partir das máquinas dos workers do Dataflow.
* Pelo menos um destino dead-letter (`clickHouseDeadLetterTable` ou `deadLetterTopic`) deve ser fornecido. Se ambos forem fornecidos, as mensagens com falha serão roteadas para os dois destinos simultaneamente.
* Quando `clickHouseDeadLetterTable` estiver definido, a tabela dead-letter já deverá existir no ClickHouse com o esquema mostrado em [Tratamento de dead-letter](#dead-letter-handling).
* Quando `deadLetterTopic` estiver definido, o tópico Pub/Sub já deverá existir.

<div id="template-parameters">
  ## Parâmetros do template
</div>

<br />

<br />

| Nome do parâmetro           | Descrição do parâmetro                                                                                                                                                            | Obrigatório | Observações                                                                                                                                                                                                        |
| --------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| `inputSubscription`         | A assinatura do Pub/Sub da qual as mensagens serão lidas. Exemplo: `projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>`.                                                     | ✅           | As mensagens devem estar codificadas em JSON.                                                                                                                                                                      |
| `clickHouseUrl`             | A URL do endpoint do ClickHouse. Use `https://` para conexões SSL (ClickHouse Cloud) ou `http://` para conexões sem SSL. Exemplo: `https://<HOST>:8443` ou `http://<HOST>:8123`.  | ✅           | Para ClickHouse Cloud, use o endpoint HTTPS na porta `8443`.                                                                                                                                                       |
| `clickHouseDatabase`        | O nome do banco de dados do ClickHouse onde a tabela de destino está localizada. Exemplo: `default`.                                                                              | ✅           |                                                                                                                                                                                                                    |
| `clickHouseTable`           | O nome da tabela do ClickHouse na qual os dados serão gravados.                                                                                                                   | ✅           | A tabela deve existir antes de executar o pipeline.                                                                                                                                                                |
| `clickHouseUsername`        | O nome de usuário para autenticação no ClickHouse.                                                                                                                                | ✅           |                                                                                                                                                                                                                    |
| `clickHousePassword`        | A senha para autenticação no ClickHouse.                                                                                                                                          | ✅           |                                                                                                                                                                                                                    |
| `clickHouseDeadLetterTable` | A tabela do ClickHouse na qual gravar mensagens com falha. Exemplo: `my_table_dead_letter`.                                                                                       |             | Pelo menos um entre `clickHouseDeadLetterTable` ou `deadLetterTopic` deve ser informado. A tabela deve existir com o esquema de dead-letter mostrado em [Tratamento de dead-letter](#dead-letter-handling).        |
| `deadLetterTopic`           | O tópico do Pub/Sub no qual publicar mensagens com falha. Exemplo: `projects/<PROJECT_ID>/topics/<TOPIC_NAME>`.                                                                   |             | Pelo menos um entre `clickHouseDeadLetterTable` ou `deadLetterTopic` deve ser informado. Os payloads com falha são publicados no tópico com `errorMessage` e `failedAt` definidos como attributes da mensagem.     |
| `windowSeconds`             | Duração, em segundos, das janelas de agrupamento em lotes baseadas em tempo.                                                                                                      |             | Veja [Agrupamento em lotes e janelamento](#batching-and-windowing) para entender a interação com `batchRowCount`. Se nenhum dos dois for definido, o modo combinado usará os valores padrão `30s` e `1000` linhas. |
| `batchRowCount`             | Número de linhas a acumular antes de fazer o flush para o ClickHouse.                                                                                                             |             | Veja [Agrupamento em lotes e janelamento](#batching-and-windowing) para entender a interação com `windowSeconds`.                                                                                                  |
| `maxInsertBlockSize`        | Número máximo de linhas por instrução `INSERT` enviada ao ClickHouse. O padrão é `1,000,000`.                                                                                     |             | Uma opção de `ClickHouseIO`.                                                                                                                                                                                       |
| `maxRetries`                | Número máximo de tentativas de retry para inserts do ClickHouse com falha. O padrão é `5`.                                                                                        |             | Uma opção de `ClickHouseIO`.                                                                                                                                                                                       |
| `insertDeduplicate`         | Indica se a desduplicação deve ser habilitada para queries `INSERT` em tabelas replicadas do ClickHouse. O padrão é `true`.                                                       |             | Uma opção de `ClickHouseIO`.                                                                                                                                                                                       |
| `insertQuorum`              | Para queries `INSERT` em tabelas replicadas, aguarda o número especificado de réplicas confirmar a gravação e linearizar a adição dos dados. `0` desabilita gravações com quorum. |             | Uma opção de `ClickHouseIO`. Desabilitada nas configurações padrão do servidor.                                                                                                                                    |
| `insertDistributedSync`     | Se habilitado, queries `INSERT` em tabelas distribuídas aguardam até que os dados sejam enviados a todos os nós do cluster. O padrão é `true`.                                    |             | Uma opção de `ClickHouseIO`.                                                                                                                                                                                       |

<Note>
  Os valores padrão de todos os parâmetros de `ClickHouseIO` podem ser encontrados em [`ClickHouseIO` Apache Beam Connector](/docs/pt-BR/integrations/connectors/data-ingestion/etl-tools/apache-beam#clickhouseiowrite-parameters).
</Note>

<div id="message-format-and-schema-mapping">
  ## Formato da mensagem e mapeamento de esquema
</div>

As mensagens do Pub/Sub devem ser objetos JSON cujos nomes de campos de nível superior correspondam exatamente aos nomes das colunas da tabela ClickHouse de destino.

Para mapear as mensagens recebidas para a tabela de destino, o pipeline faz o seguinte na inicialização:

1. Obtém o esquema da tabela ClickHouse de destino.
2. Cria um esquema `Row` do Beam com base nesse esquema do ClickHouse.
3. Para cada mensagem recebida do Pub/Sub, analisa o payload JSON e monta uma linha lendo os campos nomeados no esquema do ClickHouse.

<br />

<Warning>
  Os nomes dos campos JSON devem corresponder exatamente aos nomes das colunas do ClickHouse (a correspondência diferencia maiúsculas de minúsculas). Os campos da mensagem que não correspondem a uma coluna do ClickHouse são ignorados. Se uma coluna do ClickHouse não tiver um campo correspondente no payload JSON, o pipeline tentará gravar `NULL` nessa coluna — o que só funciona quando a coluna é declarada como [`Nullable`](/docs/pt-BR/reference/data-types/nullable). Mensagens que não puderem ser analisadas, cujos valores não puderem ser convertidos para o tipo da coluna, ou que tentariam gravar `NULL` em uma coluna que não aceita `NULL`, são encaminhadas para o destino dead-letter.
</Warning>

<div id="type-conversion">
  ### Conversão de tipos
</div>

Os valores JSON são convertidos para o tipo de coluna correspondente no ClickHouse:

| Tipo do ClickHouse                                                                 | Observações                                                                                                            |
| ---------------------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------- |
| [`Float32`](/docs/pt-BR/reference/data-types/float)                                     | Interpretado com `Float.valueOf`.                                                                                      |
| [`Float64`](/docs/pt-BR/reference/data-types/float)                                     | Interpretado com `Double.valueOf`.                                                                                     |
| [`Date`](/docs/pt-BR/reference/data-types/date)                                         | Interpretado como uma string de data ISO-8601.                                                                         |
| [`DateTime`](/docs/pt-BR/reference/data-types/datetime)                                 | Interpretado como uma string de data e hora ISO-8601 (por exemplo, `2026-01-15T12:34:56Z`).                            |
| [`Array(T)`](/docs/pt-BR/reference/data-types/array)                                    | JSON array; cada elemento é convertido para o tipo de elemento `T`. Arrays vazios ou ausentes geram um array vazio.    |
| Integer types (`Int8`/`Int16`/`Int32`/`Int64`, `UInt8`/`UInt16`/`UInt32`/`UInt64`) | Interpretados a partir do número JSON ou de sua representação em string.                                               |
| [`String`](/docs/pt-BR/reference/data-types/string)                                     | Usado como está para campos textuais; nós JSON não textuais são serializados para sua representação de string em JSON. |

<div id="batching-and-windowing">
  ## Agrupamento em lotes e janelamento
</div>

Como o pipeline opera em streaming, as linhas recebidas são acumuladas em janelas antes de serem gravadas no ClickHouse. A estratégia de janelamento é selecionada com base nos parâmetros que você fornece:

| `windowSeconds` | `batchRowCount` | Comportamento                                                                                                              |
| --------------- | --------------- | -------------------------------------------------------------------------------------------------------------------------- |
| definido        | não definido    | Janelas fixas baseadas em tempo de `windowSeconds`.                                                                        |
| não definido    | definido        | Janela global com disparo por contagem; é acionada a cada `batchRowCount` linhas.                                          |
| ambos definidos | ambos definidos | Janela global com disparo combinado; é acionada pela condição que for atendida primeiro (tempo **ou** contagem de linhas). |
| nenhum definido | nenhum definido | Modo combinado com os valores padrão: `30` segundos ou `1000` linhas, o que ocorrer primeiro.                              |

Ao ajustar esses valores, você pode equilibrar latência e eficiência de insert. Janelas menores reduzem a latência de ponta a ponta; janelas maiores produzem menos lotes `INSERT`, porém maiores.

<div id="dead-letter-handling">
  ## Tratamento de dead-letter
</div>

As mensagens que falharem no parsing de JSON, no mapeamento de esquema ou na coerção de tipo serão encaminhadas para os destinos dead-letter configurados. Pelo menos um entre `clickHouseDeadLetterTable` e `deadLetterTopic` deve ser informado; se ambos forem definidos, as mensagens com falha serão enviadas para ambos.

<div id="clickhouse-dead-letter-table">
  ### Tabela dead-letter do ClickHouse
</div>

Quando `clickHouseDeadLetterTable` é definido, a tabela dead-letter já deve existir com este esquema fixo:

| Coluna          | Tipo       | Descrição                                                     |
| --------------- | ---------- | ------------------------------------------------------------- |
| `raw_message`   | `String`   | O payload original da mensagem do Pub/Sub em texto UTF-8.     |
| `error_message` | `String`   | A mensagem da exceção que descreve por que a linha falhou.    |
| `stack_trace`   | `String`   | A stack trace completa do Java capturada no momento da falha. |
| `failed_at`     | `DateTime` | O timestamp de processamento em que a linha falhou.           |

Uma definição mínima para uma implantação de nó único:

```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>
  Adapte o engine e a cláusula `ORDER BY` à sua implantação — use `ReplicatedMergeTree` para tabelas replicadas, adicione `ON CLUSTER` em implantações distribuídas e ajuste o particionamento ou o TTL conforme necessário.
</Note>

<div id="pubsub-dead-letter-topic">
  ### Tópico dead-letter do Pub/Sub
</div>

Quando `deadLetterTopic` é definido, cada mensagem que falha é republicada no tópico com:

* **Payload**: os bytes originais da mensagem.
* **Atributo** `errorMessage`: a mensagem da exceção capturada no momento da falha.
* **Atributo** `failedAt`: o timestamp de processamento no momento em que a linha falhou.

Isso facilita reprocessar mensagens com falha assim que o problema subjacente de esquema ou do produtor tiver sido resolvido.

<div id="running-the-template">
  ## Executando o template
</div>

O template Pub/Sub para ClickHouse está disponível no Google Cloud Console.

<Note>
  Revise este documento, especialmente as seções acima, para entender plenamente os requisitos de configuração e os pré-requisitos do template.
</Note>

Faça login no Google Cloud Console e pesquise por Dataflow.

1. Clique no botão `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="Console do Dataflow" width="1872" height="886" data-path="images/integrations/data-ingestion/google-dataflow/create_job_from_template_button.webp" />

2. Quando o formulário do template abrir, insira um nome para o job e selecione a região desejada.

3. No campo `Dataflow Template`, digite `ClickHouse` ou `Pub/Sub` e selecione o template `Pub/Sub para ClickHouse`.

4. Depois de selecionado, o formulário se expande. Preencha:

   * A assinatura de entrada do Pub/Sub, no formato `projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>`.
   * A URL do endpoint do ClickHouse — para o ClickHouse Cloud, use `https://<HOST>:8443`.
   * O banco de dados do ClickHouse, a tabela de destino, o nome de usuário e a senha.
   * Pelo menos um destino dead-letter: uma tabela do ClickHouse ou um tópico do Pub/Sub (ou ambos).

5. Opcionalmente, personalize os parâmetros de agrupamento em lotes (`windowSeconds`, `batchRowCount`) e os parâmetros de ajuste do `ClickHouseIO`, conforme detalhado na seção [Template parameters](#template-parameters).

<div id="monitor-the-job">
  ### Monitore o job
</div>

Acesse a [aba Dataflow Jobs](https://console.cloud.google.com/dataflow/jobs) no Google Cloud Console para monitorar o status do job. Nela, você encontrará os detalhes do job, incluindo o progresso e eventuais erros:

<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="Console do Dataflow mostrando um job do Pub/Sub para ClickHouse em execução" width="3016" height="470" data-path="images/integrations/data-ingestion/google-dataflow/pubsub-inqueue-job.webp" />

O template também emite as seguintes métricas personalizadas no espaço de nomes `PubSubToClickHouse`, visíveis na página do job do Dataflow:

| Métrica                 | Tipo         | Descrição                                                                                                      |
| ----------------------- | ------------ | -------------------------------------------------------------------------------------------------------------- |
| `messages-received`     | Contador     | Total de mensagens do Pub/Sub recebidas pela etapa de parsing.                                                 |
| `rows-parsed-ok`        | Contador     | Mensagens convertidas com sucesso em uma linha e encaminhadas para a saída principal.                          |
| `rows-parse-failed`     | Contador     | Mensagens em que houve falha no parsing ou no mapeamento de esquema e que foram encaminhadas para dead-letter. |
| `message-payload-bytes` | Distribuição | Distribuição dos tamanhos dos payloads das mensagens recebidas do Pub/Sub, em bytes.                           |

<div id="troubleshooting">
  ## Solução de problemas
</div>

<div id="code-241-dbexception-memory-limit-total-exceeded">
  ### Erro de limite de memória (total) excedido (código 241)
</div>

Esse erro ocorre quando o ClickHouse fica sem memória ao processar grandes lotes de dados. Para resolver esse problema:

* Aumente os recursos da instância: faça upgrade do seu servidor ClickHouse para uma instância maior, com mais memória, para dar conta da carga de processamento de dados.
* Diminua o tamanho do lote: reduza `batchRowCount` (e/ou `maxInsertBlockSize`) na configuração do seu job do Dataflow para enviar fragmentos menores de dados ao ClickHouse, reduzindo o consumo de memória por lote.

<div id="all-messages-going-to-dlq">
  ### Todas as mensagens estão indo para o destino dead-letter
</div>

As causas mais comuns são:

* Os nomes dos campos JSON não correspondem exatamente aos nomes das colunas do ClickHouse (a correspondência diferencia maiúsculas de minúsculas).
* Não é possível converter o tipo de uma coluna a partir do valor JSON (por exemplo, uma string fora do padrão ISO-8601 em uma coluna `DateTime`).
* O esquema da tabela de destino mudou desde que o pipeline foi iniciado — o esquema é obtido uma vez na inicialização. Reinicie o job após aplicar as alterações no esquema.

Inspecione as colunas `error_message` e `stack_trace` da tabela dead-letter do ClickHouse (ou o atributo `errorMessage` nas mensagens dead-letter do Pub/Sub) para identificar a causa raiz.

<div id="no-rows-arriving">
  ### O pipeline inicia, mas nenhuma linha chega ao ClickHouse
</div>

* Confirme se a assinatura está recebendo mensagens — verifique a métrica `messages-received` na página do job do Dataflow.
* No modo baseado em tempo (apenas `windowSeconds`), as linhas só são gravadas ao fim de cada janela. Reduza `windowSeconds` para verificar se os flushes estão ocorrendo.
* Verifique se há conectividade de rede entre os workers do Dataflow e o endpoint do ClickHouse (firewall, VPC peering ou private service connect).

<div id="template-source-code">
  ## Código-fonte do Template
</div>

O código-fonte do Template está disponível em:

* [`GoogleCloudPlatform/DataflowTemplates`](https://github.com/GoogleCloudPlatform/DataflowTemplates/tree/main/v2/googlecloud-to-clickhouse) — o repositório original do Google Cloud Platform.
* [`ClickHouse/DataflowTemplates`](https://github.com/ClickHouse/DataflowTemplates) — o fork da ClickHouse.
