Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Template do Dataflow de Pub/Sub para ClickHouse

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



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:

  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.

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.

  1. Clique no botão CREATE JOB FROM TEMPLATE.

    Console do Dataflow
  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:

Console do Dataflow mostrando um job do Pub/Sub para ClickHouse em execução

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

Navigation