Skip to main content
Este engine permite integrar o ClickHouse ao NATS. NATS permite:
  • Publicar em ou assinar subjects de mensagens.
  • Processar novas mensagens à medida que se tornem disponíveis.

Criando uma tabela

Parâmetros obrigatórios:
  • nats_url – host:port (por exemplo, localhost:4222).
  • nats_subjects – Lista de subjects da tabela NATS para assinar/publicar. Suporta subjects curinga como foo.*.bar ou baz.>
  • nats_format – Formato da mensagem. Usa a mesma notação da função SQL FORMAT, como JSONEachRow. Para mais informações, consulte a seção Formatos.
Parâmetros opcionais:
  • nats_schema – Parâmetro que deve ser usado se o formato exigir uma definição de schema. Por exemplo, Cap’n Proto exige o caminho para o arquivo de schema e o nome do objeto raiz schema.capnp:Message.
  • nats_stream – O nome de um stream existente no NATS JetStream.
  • nats_consumer_name – O nome de um consumer pull durável existente no NATS JetStream.
  • nats_num_consumers – O número de consumers por tabela. Padrão: 1. Especifique mais consumers se a taxa de transferência de um consumer for insuficiente, somente para NATS Core.
  • nats_queue_group – Nome do queue group dos assinantes NATS. O padrão é o nome da tabela.
  • nats_max_reconnect – Obsoleto e sem efeito; a reconexão é realizada permanentemente com o timeout nats_reconnect_wait.
  • nats_reconnect_wait – Tempo de espera, em milissegundos, entre cada tentativa de reconexão. Padrão: 2000.
  • nats_server_list - Lista de servidores para conexão. Pode ser especificada para conectar a um cluster NATS.
  • nats_skip_broken_messages - Tolerância do parser de mensagens NATS a mensagens incompatíveis com o schema por bloco. Padrão: 0. Se nats_skip_broken_messages = N, o engine ignora N mensagens NATS que não podem ser convertidas (uma mensagem equivale a uma linha de dados).
  • nats_max_block_size - Número de linhas coletadas por poll(s) para descarregar dados do NATS. Padrão: max_insert_block_size.
  • nats_flush_interval_ms - Timeout para descarregar os dados lidos do NATS. Padrão: stream_flush_interval_ms.
  • nats_wait_for_flush_interval - Se true, um ciclo de streaming em segundo plano permanece aberto durante todo o intervalo de descarregamento (nats_flush_interval_ms ou stream_flush_interval_ms, caso contrário), em vez de terminar assim que a fila do consumer é esvaziada, permitindo que mais mensagens se acumulem em um único bloco, ao custo de até um intervalo de descarregamento adicional na latência de ingestão. Padrão: false (comportamento de esvaziar e seguir com baixa latência).
  • nats_username - Nome de usuário do NATS. Quando é armazenado em uma coleção nomeada definida no arquivo de configuração do servidor, a consulta não pode sobrescrever nats_url nem nats_server_list da coleção.
  • nats_password - Senha do NATS. Quando é armazenada em uma coleção nomeada definida no arquivo de configuração do servidor, a consulta não pode sobrescrever nats_url nem nats_server_list da coleção.
  • nats_token - Token de autenticação do NATS. Quando é armazenado em uma coleção nomeada definida no arquivo de configuração do servidor, a consulta não pode sobrescrever nats_url nem nats_server_list da coleção.
  • nats_credential_file - Caminho para um arquivo de credentials do NATS. Ele é aceito apenas de uma coleção nomeada definida no arquivo de configuração do servidor cujo nats_url e nats_server_list não são sobrescritos pela consulta, pois o servidor abre o caminho com seus próprios privilégios. Em uma consulta, informe o conteúdo do arquivo em nats_credentials.
  • nats_credentials - Conteúdo das credentials do NATS (a mesma carga útil de um arquivo .creds com JWT de usuário e seed). Como é a única forma que uma consulta pode usar, substitui um nats_credential_file herdado de uma coleção nomeada em vez de entrar em conflito com ele, a menos que o operador tenha bloqueado esse caminho com <nats_credential_file overridable="false">. Não pode receber uma string vazia para remover as credentials fornecidas por uma coleção nomeada.
  • nats_ca_file - Caminho para um arquivo com os certificados de CA confiáveis usados para verificar o certificado do servidor NATS. Requer nats_secure. Assim como nats_credential_file, é aceito apenas de uma coleção nomeada definida no arquivo de configuração do servidor cujo nats_url e nats_server_list não são sobrescritos pela consulta, pois o servidor abre o caminho com seus próprios privilégios.
  • nats_client_cert_file - Caminho para o certificado do cliente apresentado ao servidor NATS. Requer nats_secure e nats_client_key_file. Aceito das mesmas fontes que nats_ca_file.
  • nats_client_key_file - Caminho para a chave privada de nats_client_cert_file. Aceito das mesmas fontes que nats_ca_file.
  • nats_startup_connect_tries - Número de tentativas de conexão na inicialização. Padrão: 5.
  • nats_max_rows_per_message — O número máximo de linhas gravadas em uma mensagem NATS para formatos baseados em linha. (padrão: 1).
  • nats_commit_on_select - Faz commit das mensagens quando uma consulta é realizada. Aplica-se somente ao JetStream; o NATS Core não tem confirmações. Padrão: 0.
  • nats_handle_error_mode — Como lidar com erros no engine NATS. Valores possíveis: default (a exceção será lançada se não conseguirmos converter uma mensagem), stream (a mensagem de exceção e a mensagem bruta serão salvas nas colunas virtuais _error e _raw_message).
Conexão SSL: Para uma conexão segura, use nats_secure = 1. A verificação de certificados é controlada pela variável de ambiente CLICKHOUSE_NATS_TLS_SECURE; se o certificado estiver expirado, for autossinado, estiver ausente ou for inválido de outra forma, desative a verificação definindo CLICKHOUSE_NATS_TLS_SECURE=0. Um certificado de servidor assinado por uma CA privada é verificado apontando nats_ca_file para o certificado da CA, o que é preferível a desativar a verificação. Quando o servidor exige certificados de cliente, forneça-os com nats_client_cert_file e nats_client_key_file. Todos os três são configurações do operador: eles vêm de uma coleção nomeada definida no arquivo de configuração do servidor. Cada arquivo é lido quando a tabela se conecta; portanto, um arquivo ilegível ou malformado faz a consulta falhar em vez de falhar durante o handshake. Gravação na tabela NATS: Se a tabela lê de apenas um subject, qualquer insert publicará nesse mesmo subject. No entanto, se a tabela lê de múltiplos subjects, é preciso especificar em qual subject queremos publicar. Por isso, ao inserir em uma tabela com múltiplos subjects, é necessário definir stream_like_engine_insert_queue. Você pode selecionar um dos subjects dos quais a tabela lê e publicar seus dados nele. Por exemplo:
Também é possível adicionar configurações de formato junto com as configurações relacionadas ao NATS. Exemplo:
A configuração do servidor NATS pode ser adicionada usando o arquivo de configuração do ClickHouse. Mais especificamente, você pode adicionar sua senha para o engine NATS:

Descrição

SELECT não é particularmente útil para ler mensagens (exceto para depuração), porque cada mensagem pode ser lida apenas uma vez. É mais prático criar fluxos em tempo real usando visões materializadas. Para fazer isso:
  1. Use o engine para criar um consumer do NATS e tratá-lo como um fluxo de dados.
  2. Crie uma tabela com a estrutura desejada.
  3. Crie uma visão materializada que converta os dados do engine e os insira em uma tabela criada anteriormente.
Quando a MATERIALIZED VIEW é vinculada ao engine, ela começa a coletar dados em segundo plano. Isso permite que você receba continuamente mensagens do NATS e as converta para o formato necessário usando SELECT. Uma tabela do NATS pode ter quantas visões materializadas você quiser; elas não leem dados diretamente da tabela, mas recebem novos registros (em blocos). Assim, você pode gravar em várias tabelas com diferentes níveis de detalhamento (com agrupamento - agregação e sem). Exemplo:
Para parar de receber dados dos streams ou alterar a lógica de conversão, desanexe a visão materializada:
Se quiser alterar a tabela de destino com ALTER, recomendamos desabilitar a visão materializada para evitar discrepâncias entre a tabela de destino e os dados da view.

Colunas virtuais

  • _subject - subject da mensagem no NATS. Tipo de dado: String.
Colunas virtuais adicionais quando nats_handle_error_mode='stream':
  • _raw_message - Mensagem bruta que não pôde ser processada com sucesso. Tipo de dado: Nullable(String).
  • _error - Mensagem de exceção ocorrida durante uma falha no processamento. Tipo de dado: Nullable(String).
Observação: as colunas virtuais _raw_message e _error são preenchidas apenas em caso de exceção durante o processamento; elas são sempre NULL quando a mensagem é processada com sucesso.

Suporte a formatos de dados

O engine NATS oferece suporte a todos os formatos compatíveis com o ClickHouse. O número de linhas em uma mensagem NATS depende de o formato ser baseado em linhas ou em blocos:
  • Para formatos baseados em linhas, o número de linhas em uma mensagem NATS pode ser controlado pela configuração nats_max_rows_per_message.
  • Para formatos baseados em blocos, não é possível dividir um bloco em partes menores, mas o número de linhas em um bloco pode ser controlado pela configuração geral max_block_size.

Usando o JetStream

Antes de usar o engine NATS com o NATS JetStream, você deve criar um stream NATS e um consumer pull durável. Para isso, você pode usar, por exemplo, o utilitário nats do pacote NATS CLI:
Após criar o stream e o consumer pull durável, podemos criar uma tabela com o engine NATS. Para isso, você precisa definir: nats_stream, nats_consumer_name e nats_subjects:
As tabelas JetStream oferecem entrega de pelo menos uma vez: uma mensagem só é confirmada após ser inserida nas visões materializadas dependentes; portanto, uma mensagem cuja inserção falha ou é interrompida permanece sem confirmação e é entregue novamente. O NATS Core (sem JetStream) não oferece confirmação nem reprocessamento; portanto, tem entrega de no máximo uma vez, e uma mensagem interrompida é perdida.

Durabilidade dos dados

Esta seção se aplica somente ao JetStream. O NATS Core não tem confirmações e opera no modo de entrega no máximo uma vez, conforme descrito acima. Portanto, não há uma janela em que uma mensagem confirmada possa ser perdida. Uma tabela JetStream pode perder silenciosamente linhas já consumidas se o cache de páginas do SO for descartado antes que os dados inseridos sejam gravados em disco. Depois que um lote é enviado às visões materializadas dependentes, o consumidor confirma essas mensagens, o que permite que o stream avance além delas. No entanto, as linhas inseridas só se tornam duráveis quando a parte de destino recebe fsync, o que, por padrão, não ocorre de forma síncrona (fsync_after_insert = 0). Se o cache de páginas for perdido após a confirmação, mas antes de a parte de destino receber fsync, as mensagens não serão mais reenviadas; portanto, as linhas serão perdidas sem erro e count() será simplesmente menor. Encerrar um processo normalmente não expõe esse problema, pois o kernel mantém o cache de páginas e acaba gravando-o em disco. Já a perda do cache de páginas o expõe; exemplos incluem perda de energia no nível do dispositivo e uma reinicialização não limpa do host ou do kernel. No caminho recomendado de consumo por visão materializada (a confirmação só é enviada após a conclusão de todo o pipeline de inserção), definir fsync_after_insert = 1 (e fsync_part_directory = 1) nas tabelas MergeTree de destino torna as partes inseridas duráveis antes do envio da confirmação, reduzindo substancialmente essa janela. A configuração deve ser ativada em todas as tabelas MergeTree nas quais o lote é inserido, incluindo os destinos de visões materializadas em cascata; qualquer tabela desse tipo mantida com a configuração padrão ainda pode perder sua parte. Intermediários assíncronos não ganham durabilidade apenas com essa configuração: por exemplo, um destino Distributed insere em segundo plano quando distributed_foreground_insert = 0, que é o padrão fora do ClickHouse Cloud, portanto, ele precisa de suas próprias configurações de durabilidade ou de inserção síncrona. Essa mitigação também não se aplica a um INSERT ... SELECT ... FROM <nats_table> direto com nats_commit_on_select = 1, em que as mensagens são confirmadas quando a leitura chega ao fim, e não após o destino gravar uma parte durável.
Última modificação em 26 de setembro de 2026