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.
Requisitos do pipeline
- 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 (
clickHouseDeadLetterTableoudeadLetterTopic) deve ser fornecido. Se ambos forem fornecidos, as mensagens com falha serão roteadas para os dois destinos simultaneamente. - Quando
clickHouseDeadLetterTableestiver definido, a tabela dead-letter já deverá existir no ClickHouse com o esquema mostrado em Tratamento de dead-letter. - Quando
deadLetterTopicestiver definido, o tópico Pub/Sub já deverá existir.
Parâmetros do template
| 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. |
|
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 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 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. |
Formato da mensagem e mapeamento de esquema
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:
- Obtém o esquema da tabela ClickHouse de destino.
- Cria um esquema
Rowdo Beam com base nesse esquema do ClickHouse. - Para cada mensagem recebida do Pub/Sub, analisa o payload JSON e monta uma linha lendo os campos nomeados no esquema do ClickHouse.
Conversão de tipos
Os valores JSON são convertidos para o tipo de coluna correspondente no ClickHouse:
| Tipo do ClickHouse | Observações |
|---|---|
Float32 |
Interpretado com Float.valueOf. |
Float64 |
Interpretado com Double.valueOf. |
Date |
Interpretado como uma string de data ISO-8601. |
DateTime |
Interpretado como uma string de data e hora ISO-8601 (por exemplo, 2026-01-15T12:34:56Z). |
Array(T) |
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 |
Usado como está para campos textuais; nós JSON não textuais são serializados para sua representação de string em JSON. |
Agrupamento em lotes e janelamento
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.
Tratamento de dead-letter
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.
Tabela dead-letter do ClickHouse
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:
CREATE TABLE my_table_dead_letter (
raw_message String,
error_message String,
stack_trace String,
failed_at DateTime
) ENGINE = MergeTree()
ORDER BY failed_at;Tópico dead-letter do Pub/Sub
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.
Executando o template
O template Pub/Sub para ClickHouse está disponível no Google Cloud Console.
Faça login no Google Cloud Console e pesquise por Dataflow.
-
Clique no botão
CREATE JOB FROM TEMPLATE.
-
Quando o formulário do template abrir, insira um nome para o job e selecione a região desejada.
-
No campo
Dataflow Template, digiteClickHouseouPub/Sube selecione o templatePub/Sub para ClickHouse. -
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).
- A assinatura de entrada do Pub/Sub, no formato
-
Opcionalmente, personalize os parâmetros de agrupamento em lotes (
windowSeconds,batchRowCount) e os parâmetros de ajuste doClickHouseIO, conforme detalhado na seção Template parameters.
Monitore o job
Acesse a aba 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:

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. |
Solução de problemas
Erro de limite de memória (total) excedido (código 241)
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/oumaxInsertBlockSize) na configuração do seu job do Dataflow para enviar fragmentos menores de dados ao ClickHouse, reduzindo o consumo de memória por lote.
Todas as mensagens estão indo para o destino dead-letter
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.
O pipeline inicia, mas nenhuma linha chega ao ClickHouse
- Confirme se a assinatura está recebendo mensagens — verifique a métrica
messages-receivedna página do job do Dataflow. - No modo baseado em tempo (apenas
windowSeconds), as linhas só são gravadas ao fim de cada janela. ReduzawindowSecondspara 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).
Código-fonte do Template
O código-fonte do Template está disponível em:
GoogleCloudPlatform/DataflowTemplates— o repositório original do Google Cloud Platform.ClickHouse/DataflowTemplates— o fork da ClickHouse.