Skip to main content
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 (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.
  • Quando deadLetterTopic estiver definido, o tópico Pub/Sub já deverá existir.

Parâmetros do template



Os valores padrão de todos os parâmetros de ClickHouseIO podem ser encontrados em ClickHouseIO Apache Beam Connector.

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

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

Conversão de tipos

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

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: 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: Uma definição mínima para uma implantação de nó único:
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.

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.
Revise este documento, especialmente as seções acima, para entender plenamente os requisitos de configuração e os pré-requisitos do template.
Faça login no Google Cloud Console e pesquise por Dataflow.
  1. Clique no botão CREATE JOB FROM TEMPLATE.
  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.

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:

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

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

Código-fonte do Template

O código-fonte do Template está disponível em:
Última modificação em 23 de julho de 2026