NATS permite:
- Publicar o suscribirse a subjects de NATS.
- Procesar mensajes nuevos a medida que estén disponibles.
Crear una tabla
nats_url– host:port (por ejemplo,localhost:4222)..nats_subjects– Lista de subject de la tabla NATS a los que suscribirse o en los que publicar. Admite subjects comodín comofoo.*.barobaz.>nats_format– Formato del mensaje. Usa la misma notación que la función SQLFORMAT, comoJSONEachRow. Para obtener más información, consulta la sección Formats.
nats_schema– Parámetro que debe usarse si el formato requiere una definición de schema. Por ejemplo, Cap’n Proto requiere la ruta al archivo de schema y el nombre del objeto raízschema.capnp:Message.nats_stream– El nombre de un stream existente en NATS JetStream.nats_consumer_name– El nombre de un consumidor pull duradero existente en NATS JetStream.nats_num_consumers– El número de consumidores por tabla. Valor predeterminado:1. Especifica más consumidores si el throughput de un consumidor no es suficiente, solo para NATS core.nats_queue_group– Nombre del queue group de los suscriptores de NATS. El valor predeterminado es el nombre de la tabla.nats_max_reconnect– Deprecated y no tiene efecto; la reconexión se realiza permanentemente con el timeoutnats_reconnect_wait.nats_reconnect_wait– Tiempo de espera en milisegundos entre cada intento de reconexión. Valor predeterminado:2000.nats_server_list- Lista de servidores para la conexión. Puede especificarse para conectarse a un cluster de NATS.nats_skip_broken_messages- Tolerancia del parser de mensajes NATS a mensajes incompatibles con el schema por bloque. Valor predeterminado:0. Sinats_skip_broken_messages = N, el motor omite N mensajes NATS que no pueden parsearse (un mensaje equivale a una fila de datos).nats_max_block_size- Número de filas recopiladas por poll(s) para el flushing de datos desde NATS. Valor predeterminado: max_insert_block_size.nats_flush_interval_ms- Timeout para el flushing de los datos leídos desde NATS. Valor predeterminado: stream_flush_interval_ms.nats_wait_for_flush_interval- Si estrue, un ciclo de streaming en segundo plano permanece abierto durante todo el intervalo de flushing (nats_flush_interval_mso, de lo contrario,stream_flush_interval_ms) en lugar de finalizar tan pronto como se vacía la cola del consumidor, lo que permite que se acumulen más mensajes en un único bloque a costa de hasta un intervalo de flushing de latencia de ingestión adicional. Valor predeterminado:false(comportamiento de vaciado y continuación de baja latencia).nats_username- Nombre de usuario de NATS. Cuando se almacena en una colección con nombre definida en el archivo de configuración del servidor, la consulta no puede sobrescribir elnats_urlni elnats_server_listde la colección.nats_password- Contraseña de NATS. Cuando se almacena en una colección con nombre definida en el archivo de configuración del servidor, la consulta no puede sobrescribir elnats_urlni elnats_server_listde la colección.nats_token- Token de autenticación de NATS. Cuando se almacena en una colección con nombre definida en el archivo de configuración del servidor, la consulta no puede sobrescribir elnats_urlni elnats_server_listde la colección.nats_credential_file- Ruta a un archivo de credentials de NATS. Solo se acepta desde una colección con nombre definida en el archivo de configuración del servidor cuyonats_urlynats_server_listno sean sobrescritos por la consulta, porque el servidor abre la ruta con sus propios privilegios. En una consulta, pasa en su lugar el contenido del archivo ennats_credentials.nats_credentials- Contenido de las credentials de NATS (la misma carga útil que en un archivo.credscon JWT de usuario y seed). Dado que es la única forma que puede usar una consulta, reemplaza unnats_credential_fileheredado de una colección con nombre en lugar de entrar en conflicto con él, a menos que el operador haya bloqueado esa ruta con<nats_credential_file overridable="false">. No se le puede asignar la cadena vacía para eliminar las credentials que contiene una colección con nombre.nats_ca_file- Ruta a un archivo con los certificados de CA de confianza usados para verificar el certificado del servidor NATS. Requierenats_secure. Al igual quenats_credential_file, solo se acepta desde una colección con nombre definida en el archivo de configuración del servidor cuyonats_urlynats_server_listno sean sobrescritos por la consulta, porque el servidor abre la ruta con sus propios privilegios.nats_client_cert_file- Ruta al certificado de client presentado al servidor NATS. Requierenats_secureynats_client_key_file. Se acepta desde las mismas fuentes quenats_ca_file.nats_client_key_file- Ruta a la clave privada denats_client_cert_file. Se acepta desde las mismas fuentes quenats_ca_file.nats_startup_connect_tries- Número de intentos de conexión al inicio. Valor predeterminado:5.nats_max_rows_per_message— El número máximo de filas escritas en un mensaje NATS para formatos basados en filas. (valor predeterminado:1).nats_commit_on_select- Hace commit de los mensajes cuando se realiza una consulta. Solo se aplica a JetStream; NATS core no tiene confirmaciones. Valor predeterminado:0.nats_handle_error_mode— Cómo gestionar los errores del motor NATS. Valores posibles: default (se lanzará una excepción si no se puede parsear un mensaje), stream (el mensaje de excepción y el mensaje sin procesar se guardarán en las columnas virtuales_errory_raw_message).
nats_secure = 1.
La verificación de certificados se controla mediante la variable de entorno CLICKHOUSE_NATS_TLS_SECURE;
Si el certificado está expirado, es autofirmado, falta o no es válido por algún otro motivo, desactiva la verificación estableciendo CLICKHOUSE_NATS_TLS_SECURE=0.
Un certificado de servidor firmado por una CA privada se verifica indicando el certificado de CA en nats_ca_file,
lo cual es preferible a desactivar la verificación. Cuando el servidor requiere certificados de cliente,
proporciónalos con nats_client_cert_file y nats_client_key_file. Los tres son settings del operador:
provienen de una colección con nombre definida en el archivo de configuración del servidor. Cada archivo se lee cuando
la tabla se conecta, por lo que uno ilegible o malformado hace que falle la consulta en lugar del handshake.
Escritura en tabla NATS:
Si la tabla lee de un solo subject, cualquier insert publicará en ese mismo subject.
Sin embargo, si la tabla lee de varios subjects, es necesario especificar en qué subject queremos publicar.
Por eso, al insertar en una tabla con varios subjects, se necesita el setting stream_like_engine_insert_queue.
Puedes seleccionar uno de los subjects de los que lee la tabla y publicar allí tus datos. Por ejemplo:
Descripción
SELECT no es especialmente útil para leer mensajes (salvo para depuración), porque cada mensaje solo puede leerse una vez. Es más práctico crear flujos en tiempo real mediante vistas materializadas. Para ello:
- Use el motor para crear un consumidor de NATS y considérelo un flujo de datos.
- Cree una tabla con la estructura deseada.
- Cree una vista materializada que convierta los datos del motor y los inserte en una tabla creada previamente.
MATERIALIZED VIEW se conecta al motor, comienza a recopilar datos en segundo plano. Esto le permite recibir continuamente mensajes de NATS y convertirlos al formato requerido mediante SELECT.
Una tabla de NATS puede tener tantas vistas materializadas como desee; no leen datos de la tabla directamente, sino que reciben nuevos registros (en bloques). De este modo, puede escribir en varias tablas con distintos niveles de detalle (con agrupación - agregación y sin ella).
Ejemplo:
ALTER, le recomendamos deshabilitar la vista materializada para evitar discrepancias entre la tabla de destino y los datos de la vista.
Columnas virtuales
_subject- subject del mensaje de NATS. Tipo de dato:String.
nats_handle_error_mode='stream':
_raw_message- Mensaje sin procesar que no pudo analizarse correctamente. Tipo de dato:Nullable(String)._error- Mensaje de excepción producido durante un error de análisis. Tipo de dato:Nullable(String).
_raw_message y _error se rellenan solo en caso de excepción durante el análisis; siempre son NULL cuando el mensaje se ha analizado correctamente.
Compatibilidad con formatos de datos
El motor NATS es compatible con todos los formatos compatibles con ClickHouse. El número de filas en un mensaje de NATS depende de si el formato está basado en filas o en bloques:- En los formatos basados en filas, el número de filas en un mensaje de NATS puede controlarse mediante la configuración
nats_max_rows_per_message. - En los formatos basados en bloques, no es posible dividir un bloque en partes más pequeñas, pero el número de filas de un bloque puede controlarse mediante la configuración general max_block_size.
Uso de JetStream
Antes de usar el motor NATS con NATS JetStream, debe crear unstream de NATS y un consumidor pull duradero. Para ello, puede usar, por ejemplo, la utilidad nats del paquete NATS CLI:
creación de stream
creación de stream
creación de consumidor pull duradero
creación de consumidor pull duradero
stream y el consumidor pull duradero, podemos crear una tabla con el motor NATS. Para ello, debe inicializar: nats_stream, nats_consumer_name y nats_subjects:
Durabilidad de los datos
Esta sección se aplica únicamente a JetStream. Core NATS no tiene confirmaciones y ofrece entrega como máximo una vez, como se describió anteriormente, por lo que no existe ninguna ventana en la que pueda perderse un mensaje confirmado. Una tabla de JetStream puede perder silenciosamente filas ya consumidas si la caché de páginas del SO se descarta antes de que los datos insertados se escriban en disco. Después de enviar un lote a las vistas materializadas dependientes, el consumidor confirma esos mensajes, lo que permite que el stream avance más allá de ellos. Sin embargo, las filas insertadas solo son duraderas cuando la parte de destino se sincroniza mediantefsync, algo que, de forma predeterminada, no ocurre de manera síncrona (fsync_after_insert = 0). Si la caché de páginas se pierde después de la confirmación pero antes de que la parte de destino se sincronice mediante fsync, los mensajes ya no se vuelven a entregar, por lo que las filas se pierden sin que se produzca ningún error y count() simplemente es menor. La terminación forzada de un proceso no pone esto de manifiesto, porque el kernel conserva la caché de páginas y acaba escribiéndola en disco. En cambio, la pérdida de la caché de páginas sí lo revela; por ejemplo, una pérdida de alimentación a nivel de dispositivo o un reinicio no limpio del host o del kernel.
Para la ruta recomendada de consumo mediante vistas materializadas (la confirmación se envía solo después de que finalice toda la canalización de inserción), establecer fsync_after_insert = 1 (y fsync_part_directory = 1) en las tablas MergeTree de destino garantiza que las partes insertadas sean duraderas antes de enviar la confirmación, lo que reduce considerablemente esta ventana. La configuración debe habilitarse en cada tabla MergeTree en la que se inserte el lote, incluidos los destinos de vistas materializadas en cascada; cualquier tabla de este tipo que mantenga el valor predeterminado aún puede perder su parte. Los intermediarios asíncronos no obtienen durabilidad únicamente con esta configuración: por ejemplo, un destino Distributed inserta en segundo plano cuando distributed_foreground_insert = 0, que es el valor predeterminado fuera de ClickHouse Cloud, por lo que requiere su propia configuración de durabilidad o una inserción síncrona. Esta mitigación tampoco se aplica a un INSERT ... SELECT ... FROM <nats_table> directo con nats_commit_on_select = 1, donde los mensajes se confirman cuando la lectura llega a su fin, en lugar de después de que el destino haya escrito una parte duradera.