NATS позволяет:
- Публиковать сообщения в subject и подписываться на них.
- Обрабатывать новые сообщения по мере их поступления.
Создание таблицы
nats_url– host:port (например,localhost:4222)..nats_subjects– Списокsubjectдля таблицы NATS, на которые нужно подписываться или в которые нужно публиковать сообщения. Поддерживаютсяsubjectс подстановочными символами, напримерfoo.*.barилиbaz.>nats_format– Формат сообщений. Используется та же нотация, что и в SQL-функцииFORMAT, напримерJSONEachRow. Подробнее см. в разделе Форматы.
nats_schema– Параметр, который необходимо использовать, если формат требует определения схемы. Например, Cap’n Proto требует путь к файлу схемы и имя корневого объектаschema.capnp:Message.nats_stream– Имя существующего stream в NATS JetStream.nats_consumer_name– Имя существующего durable pull consumer в NATS JetStream.nats_num_consumers– Количество consumers на таблицу. По умолчанию:1. Укажите больше consumers, если пропускной способности одного consumer недостаточно; применяется только к NATS core.nats_queue_group– Имя queue group для подписчиков NATS. По умолчанию используется имя таблицы.nats_max_reconnect– Устарел и не имеет эффекта; переподключение выполняется постоянно с тайм-аутомnats_reconnect_wait.nats_reconnect_wait– Время ожидания в миллисекундах между попытками переподключения. По умолчанию:2000.nats_server_list- Список серверов для подключения. Можно указать для подключения к кластеру NATS.nats_skip_broken_messages- Допустимое число несовместимых со схемой сообщений для parser сообщений NATS на блок. По умолчанию:0. Еслиnats_skip_broken_messages = N, то движок пропускает N сообщений NATS, которые не удаётся разобрать (одно сообщение соответствует одной строке данных).nats_max_block_size- Число строк, собираемых с помощью poll(s) для сброса данных из NATS. По умолчанию: max_insert_block_size.nats_flush_interval_ms- Тайм-аут для сброса данных, прочитанных из NATS. По умолчанию: stream_flush_interval_ms.nats_wait_for_flush_interval- Еслиtrue, фоновый цикл стриминга остаётся открытым на весь интервал сброса (nats_flush_interval_msили, в противном случае,stream_flush_interval_ms) вместо завершения сразу после опустошения очереди consumer, что позволяет накопить больше сообщений в одном блоке ценой дополнительной задержки ингестии до одного интервала сброса. По умолчанию:false(поведение с низкой задержкой: опустошить очередь и завершить).nats_username- Имя пользователя NATS. Если оно хранится в именованной коллекции, определённой в файле конфигурации сервера, запрос не может переопределитьnats_urlилиnats_server_listэтой коллекции.nats_password- Пароль NATS. Если он хранится в именованной коллекции, определённой в файле конфигурации сервера, запрос не может переопределитьnats_urlилиnats_server_listэтой коллекции.nats_token- Токен аутентификации NATS. Если он хранится в именованной коллекции, определённой в файле конфигурации сервера, запрос не может переопределитьnats_urlилиnats_server_listэтой коллекции.nats_credential_file- Путь к файлу учётных данных NATS. Принимается только из именованной коллекции, определённой в файле конфигурации сервера, если еёnats_urlиnats_server_listне переопределяются запросом, поскольку сервер открывает этот путь со своими правами. Вместо этого в запросе передайте содержимое файла вnats_credentials.nats_credentials- Содержимое учётных данных NATS (та же полезная нагрузка, что и в файле.credsс JWT пользователя и seed). Поскольку это единственный вариант, доступный в запросе, он заменяетnats_credential_file, унаследованный из именованной коллекции, а не конфликтует с ним, если только оператор не заблокировал этот путь с помощью<nats_credential_file overridable="false">. Ему нельзя присвоить пустую строку, чтобы удалить учётные данные, содержащиеся в именованной коллекции.nats_ca_file- Путь к файлу с доверенными CA‑сертификатами, используемыми для проверки сертификата сервера NATS. Требуетnats_secure. Как иnats_credential_file, принимается только из именованной коллекции, определённой в файле конфигурации сервера, если еёnats_urlиnats_server_listне переопределяются запросом, поскольку сервер открывает этот путь со своими правами.nats_client_cert_file- Путь к клиентскому сертификату, предоставляемому серверу NATS. Требуетnats_secureиnats_client_key_file. Принимается из тех же источников, что иnats_ca_file.nats_client_key_file- Путь к закрытому ключуnats_client_cert_file. Принимается из тех же источников, что иnats_ca_file.nats_startup_connect_tries- Количество попыток подключения при запуске. По умолчанию:5.nats_max_rows_per_message— Максимальное количество строк, записываемых в одно сообщение NATS для построчных форматов. (по умолчанию:1).nats_commit_on_select- Выполнять коммит сообщений при выполнении запроса. Применяется только к JetStream; в core NATS нет подтверждений. По умолчанию:0.nats_handle_error_mode— Как обрабатывать ошибки движка NATS. Возможные значения: default (если не удаётся разобрать сообщение, будет сгенерировано исключение), stream (сообщение об исключении и исходное сообщение будут сохранены в виртуальных столбцах_errorи_raw_message).
nats_secure = 1.
Проверка сертификата управляется переменной окружения CLICKHOUSE_NATS_TLS_SECURE;
если сертификат просрочен, самоподписан, отсутствует или иным образом недействителен, отключите проверку, задав CLICKHOUSE_NATS_TLS_SECURE=0.
Серверный сертификат, подписанный приватным CA, проверяется путём указания в nats_ca_file пути к CA‑сертификату,
что предпочтительнее отключения проверки. Если сервер требует клиентские сертификаты,
передайте их через nats_client_cert_file и nats_client_key_file. Все три являются настройками уровня оператора:
они берутся из named collection, определённой в файле конфигурации сервера. Каждый файл читается при
подключении таблицы, поэтому нечитаемый или некорректный файл приводит к ошибке запроса, а не рукопожатия.
Запись в таблицу NATS:
Если таблица читает только из одного subject, любая вставка будет публиковаться в тот же subject.
Однако, если таблица читает из нескольких subject, необходимо указать, в какой subject публиковать данные.
Поэтому при вставке в таблицу с несколькими subject требуется установка параметра stream_like_engine_insert_queue.
Вы можете выбрать один из subject, из которых читает таблица, и публиковать данные туда. Например:
Описание
SELECT не особенно полезен для чтения сообщений (кроме отладки), поскольку каждое сообщение можно прочитать только один раз. Гораздо практичнее создавать потоки в реальном времени с помощью materialized views. Для этого:
- Используйте движок, чтобы создать consumer NATS, и рассматривайте его как поток данных.
- Создайте таблицу с нужной структурой.
- Создайте materialized view, которое преобразует данные из движка и помещает их в ранее созданную таблицу.
MATERIALIZED VIEW подключается к движку, оно начинает собирать данные в фоновом режиме. Это позволяет непрерывно получать сообщения из NATS и преобразовывать их в требуемый формат с помощью SELECT.
Одна таблица NATS может иметь сколько угодно materialized views; они не читают данные из таблицы напрямую, а получают новые записи (блоками), поэтому вы можете записывать данные в несколько таблиц с разным уровнем детализации (с группировкой — агрегацией — и без неё).
Пример:
ALTER, мы рекомендуем отключить materialized view, чтобы избежать расхождений между целевой таблицей и данными из представления.
Виртуальные столбцы
_subject- subject сообщения NATS. Тип данных:String.
nats_handle_error_mode='stream':
_raw_message- Исходное сообщение, которое не удалось успешно разобрать. Тип данных:Nullable(String)._error- Сообщение об исключении, возникшем при сбое разбора. Тип данных:Nullable(String).
_raw_message и _error заполняются только в случае исключения при разборе; если сообщение успешно разобрано, они всегда имеют значение NULL.
Поддержка форматов данных
Движок NATS поддерживает все форматы, доступные в ClickHouse. Количество строк в одном сообщении NATS зависит от того, является ли формат построчным или блочным:- Для построчных форматов количеством строк в одном сообщении NATS можно управлять с помощью настройки
nats_max_rows_per_message. - Для блочных форматов блок нельзя разделить на более мелкие части, однако количество строк в одном блоке можно регулировать с помощью общей настройки max_block_size.
Использование JetStream
Перед использованием движка NATS с NATS JetStream необходимо создать stream в NATS и durable pull consumer. Для этого можно использовать, например, утилитуnats из пакета NATS CLI:
создание stream
создание stream
создание durable pull consumer
создание durable pull consumer
Долговечность данных
Этот раздел относится только к JetStream. В Core NATS подтверждений нет, а доставка выполняется по принципу «не более одного раза», как описано выше, поэтому там нет окна, в котором подтвержденное сообщение может быть потеряно. Таблица JetStream может незаметно потерять уже обработанные строки, если кэш страниц ОС будет сброшен до того, как вставленные данные окажутся записаны на диск. После передачи батча в зависимые materialized view consumer подтверждает эти сообщения, что позволяет stream продвинуться дальше. Однако вставленные строки становятся долговечными только после fsync целевой части, а по умолчанию он не выполняется синхронно (fsync_after_insert = 0). Если кэш страниц теряется после подтверждения, но до fsync целевой части, сообщения уже не будут доставлены повторно, поэтому строки пропадают без каких-либо ошибок, а count() просто оказывается меньше. Обычное завершение процесса эту проблему не выявляет, так как ядро сохраняет кэш страниц и в конечном итоге записывает его на диск. А вот потеря кэша страниц ее проявляет: например, при отключении питания на уровне устройства или некорректном сбросе хоста либо ядра.
Для рекомендуемого способа потребления через materialized view (подтверждение отправляется только после завершения всего конвейера вставки) установка fsync_after_insert = 1 (и fsync_part_directory = 1) для целевых таблиц MergeTree обеспечивает долговечность вставленных частей до отправки подтверждения, что существенно сужает это окно. Настройку необходимо включить для каждой таблицы MergeTree, в которую вставляется батч, включая цели каскадных materialized view; любая такая таблица, оставленная со значением по умолчанию, по-прежнему может потерять свою часть. Асинхронные промежуточные звенья одной лишь этой настройкой долговечности не получают: например, цель Distributed выполняет вставку в фоновом режиме при distributed_foreground_insert = 0, что является значением по умолчанию за пределами ClickHouse Cloud, поэтому ей нужны собственные настройки долговечности или синхронная вставка. Эта мера также не применима к прямому INSERT ... SELECT ... FROM <nats_table> с nats_commit_on_select = 1, где сообщения подтверждаются по достижении конца чтения, а не после того, как пункт назначения записал долговечную часть.