> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-parallel-read-in-order-multi-part.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

> 该引擎支持将 ClickHouse 与 NATS 集成，以发布或订阅消息 subject，并在有新消息可用时进行处理。

# NATS 表引擎

该引擎支持将 ClickHouse 与 [NATS](https://nats.io/) 集成。

`NATS` 可让你：

* 发布或订阅消息 subject。
* 在有新消息可用时进行处理。

## 创建表

```sql theme={null}
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
    name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
    name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
    ...
) ENGINE = NATS SETTINGS
    nats_url = 'host:port',
    nats_subjects = 'subject1,subject2,...',
    nats_format = 'data_format'[,]
    [nats_schema = '',]
    [nats_num_consumers = N,]
    [nats_queue_group = 'group_name',]
    [nats_secure = false,]
    [nats_max_reconnect = N,]
    [nats_reconnect_wait = N,]
    [nats_server_list = 'host1:port1,host2:port2,...',]
    [nats_skip_broken_messages = N,]
    [nats_max_block_size = N,]
    [nats_flush_interval_ms = N,]
    [nats_username = 'user',]
    [nats_password = 'password',]
    [nats_token = 'clickhouse',]
    [nats_credentials = '-----BEGIN NATS USER JWT----- ...',]
    [nats_startup_connect_tries = 5,]
    [nats_max_rows_per_message = 1,]
    [nats_commit_on_select = false,]
    [nats_handle_error_mode = 'default']
```

必需参数：

* `nats_url` – host:port (例如 `localhost:4222`) 。
* `nats_subjects` – NATS 表要订阅/发布的 subject 列表。支持通配符 subject，例如 `foo.*.bar` 或 `baz.>`
* `nats_format` – 消息格式。使用与 SQL `FORMAT` 函数相同的表示法，例如 `JSONEachRow`。更多信息，请参见[格式](/zh/reference/formats/index)部分。

可选参数：

* `nats_schema` – 如果格式需要 schema 定义，则必须使用此参数。例如，[Cap'n Proto](https://capnproto.org/) 需要提供 schema 文件的 path 以及根对象 `schema.capnp:Message` 的名称。
* `nats_stream` – NATS JetStream 中现有 stream 的名称。
* `nats_consumer_name` – NATS JetStream 中现有持久化拉取消费者的名称。
* `nats_num_consumers` – 每个表的消费者数量。默认值：`1`。仅适用于 NATS core：如果单个消费者的吞吐量不足，可指定更多消费者。
* `nats_queue_group` – NATS 订阅者的 queue group 名称。默认值为表名。
* `nats_max_reconnect` – 已弃用且不起作用；系统会按 `nats_reconnect_wait` timeout 永久执行重连。
* `nats_reconnect_wait` – 每次重连尝试之间的休眠时间 (毫秒) 。默认值：`2000`。
* `nats_server_list` - 用于 connection 的 server 列表。可用于连接到 NATS cluster。
* `nats_skip_broken_messages` - NATS 消息解析器对每个块中与 schema 不兼容消息的容忍数量。默认值：`0`。如果 `nats_skip_broken_messages = N`，则该引擎会跳过 *N* 条无法解析的 NATS 消息 (1 条消息等于 1 行数据) 。
* `nats_max_block_size` - 为从 NATS flush 数据而通过 poll 收集的行数。默认值：[max\_insert\_block\_size](/zh/reference/settings/session-settings/max-insert#max_insert_block_size)。
* `nats_flush_interval_ms` - flush 从 NATS 读取的数据的 timeout。默认值：[stream\_flush\_interval\_ms](/zh/reference/settings/session-settings/stream#stream_flush_interval_ms)。
* `nats_wait_for_flush_interval` - 如果为 `true`，后台流式周期将在整个 flush 时间间隔内保持开启 (`nats_flush_interval_ms`，否则为 `stream_flush_interval_ms`) ，而不是在消费者 queue 耗尽后立即结束；这样可使更多消息累积到单个块中，但会额外增加最多一个 flush 时间间隔的摄取延迟。默认值：`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 凭据文件的 path。仅接受来自服务器配置文件中定义的命名集合的值，且该集合的 `nats_url` 和 `nats_server_list` 未被查询覆盖，因为服务器会使用自身权限打开该 path。在查询中，请改为通过 `nats_credentials` 传入文件内容。
* `nats_credentials` - NATS 凭据内容 (与包含用户 JWT 和 seed 的 `.creds` 文件中的负载相同) 。由于这是查询唯一可用的写法，它会替换从命名集合继承的 `nats_credential_file`，而不会与之冲突；除非操作员通过 `<nats_credential_file overridable="false">` 锁定了该 path。不能将其赋值为空字符串来移除命名集合携带的凭据。
* `nats_ca_file` - 包含受信任 CA 证书的文件 path，用于验证 NATS 服务器证书。需要 `nats_secure`。与 `nats_credential_file` 一样，仅接受来自服务器配置文件中定义的命名集合的值，且该集合的 `nats_url` 和 `nats_server_list` 未被查询覆盖，因为服务器会使用自身权限打开该 path。
* `nats_client_cert_file` - 提供给 NATS 服务器的客户端证书的 path。需要 `nats_secure` 和 `nats_client_key_file`。接受与 `nats_ca_file` 相同来源的值。
* `nats_client_key_file` - `nats_client_cert_file` 私钥的 path。接受与 `nats_ca_file` 相同来源的值。
* `nats_startup_connect_tries` - 启动时的连接尝试次数。默认值：`5`。
* `nats_max_rows_per_message` — 对于按行组织的格式，一条 NATS 消息中写入的最大行数。默认值：`1`。
* `nats_commit_on_select` - 发出查询时提交消息。仅适用于 JetStream；NATS core 不提供确认机制。默认值：`0`。
* `nats_handle_error_mode` — NATS 引擎的错误处理方式。可选值：default (如果消息解析失败，则抛出异常) ；stream (异常消息和原始消息将保存在虚拟列 `_error` 和 `_raw_message` 中) 。

SSL 连接：

要建立安全连接，请使用 `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` 提供相应文件。这三个都是运维设置：
它们来自服务器配置文件中定义的命名集合。表建立连接时会读取每个文件，因此
文件无法读取或格式错误会导致查询失败，而不是握手失败。

写入 NATS 表：

如果表只从一个 subject 读取，则任何 insert 都会发布到该 subject。
但如果表从多个 subject 读取，就需要指定要发布到哪个 subject。
因此，每当向具有多个 subject 的表中 insert 数据时，都需要设置 `stream_like_engine_insert_queue`。
你可以从该表读取的 subject 中选择一个，并将数据发布到该 subject。例如：

```sql theme={null}
CREATE TABLE queue (
    key UInt64,
    value UInt64
  ) ENGINE = NATS
    SETTINGS nats_url = 'localhost:4444',
             nats_subjects = 'subject1,subject2',
             nats_format = 'JSONEachRow';

INSERT INTO queue
SETTINGS stream_like_engine_insert_queue = 'subject2'
VALUES (1, 1);
```

此外，还可以添加格式设置以及 NATS 相关设置。

示例：

```sql theme={null}
CREATE TABLE queue (
    key UInt64,
    value UInt64,
    date DateTime
  ) ENGINE = NATS
    SETTINGS nats_url = 'localhost:4444',
             nats_subjects = 'subject1',
             nats_format = 'JSONEachRow',
             date_time_input_format = 'best_effort';
```

可通过 ClickHouse 配置文件添加 NATS 服务器配置。
更具体地说，可以添加 NATS 引擎的密码：

```xml theme={null}
<nats>
    <user>click</user>
    <password>house</password>
    <token>clickhouse</token>
</nats>
```

## 说明

`SELECT` 并不特别适合用于读取消息 (调试除外) ，因为每条消息只能读取一次。更实用的做法是使用 [materialized views](/zh/reference/statements/create/view) 创建实时处理链路。为此：

1. 使用该引擎创建一个 NATS 消费者，并将其视为数据 stream。
2. 创建一个具有所需结构的表。
3. 创建一个 materialized view，将该引擎中的数据转换后写入前面创建的表中。

当 `MATERIALIZED VIEW` 连接到该引擎后，就会开始在后台收集数据。这样一来，你就可以持续接收来自 NATS 的消息，并使用 `SELECT` 将其转换为所需格式。
一个 NATS 表可以拥有任意数量的 materialized view；它们不会直接从该表读取数据，而是接收新的记录 (以块的形式) ，因此你可以写入多个明细粒度不同的表 (带分组聚合和不带分组聚合) 。

示例：

```sql theme={null}
CREATE TABLE queue (
    key UInt64,
    value UInt64
  ) ENGINE = NATS
    SETTINGS nats_url = 'localhost:4444',
             nats_subjects = 'subject1',
             nats_format = 'JSONEachRow',
             date_time_input_format = 'best_effort';

CREATE TABLE daily (key UInt64, value UInt64)
    ENGINE = MergeTree() ORDER BY key;

CREATE MATERIALIZED VIEW consumer TO daily
    AS SELECT key, value FROM queue;

SELECT key, value FROM daily ORDER BY key;
```

若要停止接收 stream 数据或更改转换逻辑，请分离 materialized view：

```sql theme={null}
DETACH TABLE consumer;
ATTACH TABLE consumer;
```

如果你想使用 `ALTER` 修改目标表，我们建议先禁用物化视图，以避免目标表与视图数据之间出现不一致。

## 虚拟列

* `_subject` - NATS 消息的 subject。数据类型：`String`。

当 `nats_handle_error_mode='stream'` 时，会提供以下额外的虚拟列：

* `_raw_message` - 无法成功解析的原始消息。数据类型：`Nullable(String)`。
* `_error` - 解析失败时产生的异常消息。数据类型：`Nullable(String)`。

注意：`_raw_message` 和 `_error` 这两个虚拟列仅会在解析过程中发生异常时填充；如果消息解析成功，它们始终为 `NULL`。

## 数据格式支持

NATS 引擎 支持 ClickHouse 支持的所有[格式](/zh/reference/formats/index)。
单条 NATS 消息中的行数取决于格式是按行还是按块：

* 对于按行的格式，可通过设置 `nats_max_rows_per_message` 来控制单条 NATS 消息中的行数。
* 对于按块的格式，无法将块拆分成更小的部分，但一个块中的行数可通过通用设置 [max\_block\_size](/zh/reference/settings/session-settings/max#max_block_size) 控制。

## 使用 JetStream

在结合 NATS JetStream 使用 NATS 引擎之前，您必须先创建一个 NATS stream 和一个持久化拉取消费者。为此，您可以使用 [NATS CLI](https://github.com/nats-io/natscli) 包中的 `nats` 工具，例如：

<Accordion title="创建 stream">
  ```bash theme={null}
  $ nats stream add
  ? Stream Name stream_name
  ? Subjects stream_subject
  ? Storage file
  ? Replication 1
  ? Retention Policy Limits
  ? Discard Policy Old
  ? Stream Messages Limit -1
  ? Per Subject Messages Limit -1
  ? Total Stream Size -1
  ? Message TTL -1
  ? Max Message Size -1
  ? Duplicate tracking time window 2m0s
  ? Allow message Roll-ups No
  ? Allow message deletion Yes
  ? Allow purging subjects or the entire stream Yes
  Stream stream_name was created

  Information for Stream stream_name created 2025-10-03 14:12:51

                  Subjects: stream_subject
                  Replicas: 1
                   Storage: File

  Options:

                 Retention: Limits
           Acknowledgments: true
            Discard Policy: Old
          Duplicate Window: 2m0s
                Direct Get: true
         Allows Msg Delete: true
              Allows Purge: true
  Allows Per-Message TTL: false
            Allows Rollups: false

  Limits:

          Maximum Messages: unlimited
       Maximum Per Subject: unlimited
             Maximum Bytes: unlimited
               Maximum Age: unlimited
      Maximum Message Size: unlimited
         Maximum Consumers: unlimited

  State:

                  Messages: 0
                     Bytes: 0 B
            First Sequence: 0
             Last Sequence: 0
          Active Consumers: 0
  ```
</Accordion>

<Accordion title="创建持久化拉取消费者">
  ```bash theme={null}
  $ nats consumer add
  ? Select a Stream stream_name
  ? Consumer name consumer_name
  ? Delivery target (empty for Pull Consumers)
  ? Start policy (all, new, last, subject, 1h, msg sequence) all
  ? Acknowledgment policy explicit
  ? Replay policy instant
  ? Filter Stream by subjects (blank for all)
  ? Maximum Allowed Deliveries -1
  ? Maximum Acknowledgments Pending 0
  ? Deliver headers only without bodies No
  ? Add a Retry Backoff Policy No
  Information for Consumer stream_name > consumer_name created 2025-10-03T14:13:51+03:00

  Configuration:

                      Name: consumer_name
                 Pull Mode: true
            Deliver Policy: All
                Ack Policy: Explicit
                  Ack Wait: 30.00s
             Replay Policy: Instant
           Max Ack Pending: 1,000
         Max Waiting Pulls: 512

  State:

  Last Delivered Message: Consumer sequence: 0 Stream sequence: 0
      Acknowledgment Floor: Consumer sequence: 0 Stream sequence: 0
          Outstanding Acks: 0 out of maximum 1,000
      Redelivered Messages: 0
      Unprocessed Messages: 0
             Waiting Pulls: 0 of maximum 512
  ```
</Accordion>

创建好 stream 和持久化拉取消费者后，我们就可以创建一个使用 NATS 引擎的表。为此，您需要设置：nats\_stream、nats\_consumer\_name 和 nats\_subjects：

```SQL theme={null}
CREATE TABLE nats_jet_stream (
    key UInt64,
    value UInt64
  ) ENGINE NATS
    SETTINGS  nats_url = 'localhost:4222',
              nats_stream = 'stream_name',
              nats_consumer_name = 'consumer_name',
              nats_subjects = 'stream_subject',
              nats_format = 'JSONEachRow';
```

JetStream 表提供至少一次投递保证：消息只有在插入其依赖的 materialized views 后才会被确认，因此插入失败或中断的消息会保持未确认状态，并被重新投递。Core NATS (不含 JetStream) 没有确认或重放机制，因此仅提供至多一次语义，中断的消息会丢失。

## 数据持久性

本节仅适用于 JetStream。Core NATS 没有确认机制，且如上所述采用至多一次语义，因此不存在已确认消息可能丢失的时间窗口。

如果在插入的数据写入磁盘前 OS page cache 被丢弃，JetStream 表可能会在无提示的情况下丢失已消费的行。批次被推送到依赖的 materialized view 后，消费者会确认这些消息，从而使 stream 越过这些消息继续推进。但是，插入的行只有在目标 parts 被 fsync 后才真正持久化；默认情况下不会同步执行此操作 (`fsync_after_insert = 0`) 。如果在确认之后、目标 parts 被 fsync 之前 page cache 丢失，消息不会再被重新投递，因此这些行会在没有任何错误的情况下丢失，`count()` 结果只会变小。普通的进程 kill 不会暴露这个问题，因为 kernel 会保留 page cache，并最终将其写回。page cache 丢失则会暴露该问题，例如设备级断电，以及主机或 kernel 的非正常重置。

对于推荐的 materialized-view 消费路径 (仅在整个插入管道完成后才发送确认) ，在目标 `MergeTree` 表上设置 `fsync_after_insert = 1` (以及 `fsync_part_directory = 1`) ，可确保插入的 parts 在发送确认前持久化，从而大幅缩小这一时间窗口。必须在批次插入的每个 `MergeTree` 表上启用该设置，包括级联 materialized-view 的目标表；任何仍使用默认设置的此类表仍可能丢失其 parts。异步中间层不会仅凭此设置获得持久性：例如，当 `distributed_foreground_insert = 0` 时，`Distributed` 目标端会在后台插入；这是 ClickHouse Cloud 之外的默认设置，因此它需要自身的持久性设置或同步插入。此缓解措施也不适用于带有 `nats_commit_on_select = 1` 的直接 `INSERT ... SELECT ... FROM <nats_table>`；在这种情况下，消息会在读取结束时被确认，而不是在目标端写入持久化 parts 后。
