NATS permite:
- Publicar em ou assinar subjects de mensagens.
- Processar novas mensagens à medida que se tornem disponíveis.
Criando uma tabela
nats_url– host:port (por exemplo,localhost:4222).nats_subjects– Lista de subjects da tabela NATS para assinar/publicar. Suporta subjects curinga comofoo.*.baroubaz.>nats_format– Formato da mensagem. Usa a mesma notação da função SQLFORMAT, comoJSONEachRow. Para mais informações, consulte a seção Formatos.
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 raizschema.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 timeoutnats_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. Senats_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- Setrue, um ciclo de streaming em segundo plano permanece aberto durante todo o intervalo de descarregamento (nats_flush_interval_msoustream_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 sobrescrevernats_urlnemnats_server_listda 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 sobrescrevernats_urlnemnats_server_listda 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 sobrescrevernats_urlnemnats_server_listda 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 cujonats_urlenats_server_listnã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 emnats_credentials.nats_credentials- Conteúdo das credentials do NATS (a mesma carga útil de um arquivo.credscom JWT de usuário e seed). Como é a única forma que uma consulta pode usar, substitui umnats_credential_fileherdado 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. Requernats_secure. Assim comonats_credential_file, é aceito apenas de uma coleção nomeada definida no arquivo de configuração do servidor cujonats_urlenats_server_listnã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. Requernats_secureenats_client_key_file. Aceito das mesmas fontes quenats_ca_file.nats_client_key_file- Caminho para a chave privada denats_client_cert_file. Aceito das mesmas fontes quenats_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_errore_raw_message).
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:
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:
- Use o engine para criar um consumer do NATS e tratá-lo como um fluxo de dados.
- Crie uma tabela com a estrutura desejada.
- Crie uma visão materializada que converta os dados do engine e os insira em uma tabela criada anteriormente.
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:
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.
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).
_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árionats do pacote NATS CLI:
criando um stream
criando um stream
criando um consumer pull durável
criando um consumer pull durável
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 recebefsync, 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.