> ## 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 table がサブスクライブ/パブリッシュする subject の一覧。`foo.*.bar` や `baz.>` のようなワイルドカード subject をサポートします
* `nats_format` – メッセージのフォーマット。`JSONEachRow` など、SQL の `FORMAT` 関数と同じ記法を使用します。詳細は [フォーマット](/ja/reference/formats/index) セクションを参照してください。

パラメータ:

* `nats_schema` – フォーマットでスキーマ定義が必要な場合に使用する必要があるパラメータです。たとえば、[Cap'n Proto](https://capnproto.org/) では、スキーマファイルへのパスとルート `schema.capnp:Message` オブジェクト名が必要です。
* `nats_stream` – NATS JetStream 内の既存の stream 名。
* `nats_consumer_name` – NATS JetStream 内の既存の durable pull コンシューマー名。
* `nats_num_consumers` – テーブルごとのコンシューマー数。デフォルト: `1`。NATS core のみを使用していて、1 つのコンシューマーのスループットが不十分な場合は、より多くのコンシューマーを指定します。
* `nats_queue_group` – NATS subscriber の queue group 名。デフォルトはテーブル名です。
* `nats_max_reconnect` – 非推奨であり、効果はありません。再接続は `nats&#95;reconnect&#95;wait` タイムアウトで恒久的に実行されます。
* `nats_reconnect_wait` – 再接続試行のたびに待機する時間 (ミリ秒単位) 。デフォルト: `2000`。
* `nats_server_list` - 接続先の server 一覧。NATS クラスターに接続するために指定できます。
* `nats_skip_broken_messages` - ブロックごとに許容する、スキーマ非互換の NATS メッセージ の数。デフォルト: `0`。`nats_skip_broken_messages = N` の場合、このエンジンは解析できない *N* 件の NATS メッセージ をスキップします (1 メッセージ は 1 行のデータに相当します) 。
* `nats_max_block_size` - NATS からデータを flush するために poll で収集する行数。デフォルト: [max\_insert\_block\_size](/ja/reference/settings/session-settings/max-insert#max_insert_block_size)。
* `nats_flush_interval_ms` - NATS から読み取ったデータを flush するまでのタイムアウト。デフォルト: [stream\_flush\_interval\_ms](/ja/reference/settings/session-settings/stream#stream_flush_interval_ms)。
* `nats_wait_for_flush_interval` - `true` の場合、background streaming cycle は、コンシューマー queue が空になるとすぐに終了するのではなく、flush interval 全体 (`nats_flush_interval_ms`、指定されていない場合は `stream_flush_interval_ms`) にわたって開いたままとなります。これにより、最大 1 flush interval 分の追加インジェストレイテンシーと引き換えに、より多くの メッセージ を単一のブロックに蓄積できます。デフォルト: `false` (低レイテンシーの drain-and-go 動作) 。
* `nats_username` - NATS username。サーバー設定ファイルで定義された named collection に保存されている場合、クエリでその collection の `nats_url` または `nats_server_list` を override することはできません。
* `nats_password` - NATS password。サーバー設定ファイルで定義された named collection に保存されている場合、クエリでその collection の `nats_url` または `nats_server_list` を override することはできません。
* `nats_token` - NATS auth token。サーバー設定ファイルで定義された named collection に保存されている場合、クエリでその collection の `nats_url` または `nats_server_list` を override することはできません。
* `nats_credential_file` - NATS credentials file への path。server が自身の権限で path を開くため、クエリによって `nats_url` と `nats_server_list` が override されない、サーバー設定ファイルで定義された named collection からのみ受け付けられます。クエリでは、代わりにファイルの内容を `nats_credentials` に渡します。
* `nats_credentials` - NATS credentials の内容 (user JWT と seed を含む `.creds` file と同じペイロード)。クエリで使用できる唯一の指定方法であるため、`<nats_credential_file overridable="false">` により operator がその path をロックしている場合を除き、競合するのではなく named collection から継承した `nats_credential_file` を置き換えます。named collection が保持する credentials を削除するために空文字列を割り当てることはできません。
* `nats_ca_file` - NATS server certificate の検証に使用する、信頼された CA certificates を含む file への path。`nats_secure` が必要です。`nats_credential_file` と同様に、server が自身の権限で path を開くため、クエリによって `nats_url` と `nats_server_list` が override されない、サーバー設定ファイルで定義された named collection からのみ受け付けられます。
* `nats_client_cert_file` - NATS server に提示する client certificate への 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` — 行ベースのフォーマットで、1 つの NATS メッセージ に書き込まれる最大行数。 (デフォルト: `1`) 。
* `nats_commit_on_select` - クエリ実行時に メッセージ を commit します。JetStream にのみ適用されます。core NATS には acknowledgement はありません。デフォルト: `0`。
* `nats_handle_error_mode` — NATS エンジンでの error の処理方法。設定可能な値: default (メッセージ の解析に失敗すると exception が throw されます) 、stream (exception メッセージ と raw メッセージ が仮想カラム `_error` および `_raw_message` に保存されます) 。

SSL 接続:

安全な接続には、`nats_secure = 1`を使用します。
証明書の検証は、`CLICKHOUSE_NATS_TLS_SECURE`環境変数によって制御されます。
証明書が期限切れ、自己署名、欠落、またはその他の理由で無効な場合は、`CLICKHOUSE_NATS_TLS_SECURE=0`を設定して検証を無効にします。

private CA によって署名された server certificate は、`nats_ca_file`で CA certificate を指定することで検証します。
これは検証を無効にするよりも望ましい方法です。server が client certificates を必要とする場合は、
`nats_client_cert_file`と`nats_client_key_file`で指定します。これら3つはすべて運用者設定です:
サーバー設定ファイルで定義された named collection から取得されます。各ファイルはテーブルの接続時に読み取られるため、
読み取り不能または形式不正のファイルがあると、ハンドシェイクではなくクエリが失敗します。

NATS table への書き込み:

テーブルが1つのsubjectのみを読み取る場合、いかなるINSERTも同じsubjectにパブリッシュされます。
ただし、テーブルが複数のsubjectを読み取る場合は、どのsubjectにパブリッシュするかを指定する必要があります。
そのため、複数のsubjectを持つテーブルにINSERTする際は、`stream_like_engine_insert_queue`を設定する必要があります。
テーブルが読み取る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';
```

NATS server の設定は、ClickHouse の設定ファイルを使用して追加できます。
具体的には、NATS エンジン用の password を追加できます:

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

## 説明

`SELECT` は、メッセージの読み取りにはあまり適していません (`debugging` 目的を除く) 。各メッセージは一度しか読み取れないためです。より実用的なのは、[materialized view](/ja/reference/statements/create/view) を使用してリアルタイムのスレッドを作成することです。これを行うには、次の手順に従います。

1. engine を使用して NATS コンシューマーを作成し、それをデータストリームとして扱います。
2. 必要な structure を持つ table を作成します。
3. engine からのデータを変換し、あらかじめ作成した table に格納する materialized view を作成します。

`MATERIALIZED VIEW` を engine に接続すると、バックグラウンドでデータの収集を開始します。これにより、NATS からメッセージを継続的に受信し、`SELECT` を使って必要なフォーマットに変換できます。
1 つの NATS table には、必要な数だけ materialized view を作成できます。これらは table から直接データを読み取るのではなく、新しいレコードをブロック単位で受け取ります。そのため、詳細度の異なる複数の table に書き込むことができます (グループ化あり - aggregation、なし) 。

例:

```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;
```

ストリームデータの受信を停止するか、変換ロジックを変更するには、materialized viewをデタッチします:

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

`ALTER` を使用してターゲットテーブルを変更する場合は、ターゲットテーブルとビューのデータの不整合を避けるため、materialized viewを無効化することを推奨します。

## 仮想カラム

* `_subject` - NATS メッセージの subject。データ型: `String`。

`nats_handle_error_mode='stream'` の場合は、次の仮想カラムも利用できます。

* `_raw_message` - 正常にパースできなかった生のメッセージ。データ型: `Nullable(String)`。
* `_error` - パース失敗時に発生した例外メッセージ。データ型: `Nullable(String)`。

注: 仮想カラム `_raw_message` と `_error` に値が入るのは、パース中に例外が発生した場合のみです。メッセージが正常にパースされた場合、これらは常に `NULL` です。

## データフォーマットのサポート

NATS エンジンは、ClickHouse でサポートされているすべての[フォーマット](/ja/reference/formats/index)に対応しています。
1 つの NATS メッセージに含まれる行数は、そのフォーマットが行ベースかブロックベースかによって異なります。

* 行ベースのフォーマットでは、1 つの NATS メッセージに含める行数を `nats_max_rows_per_message` の設定で制御できます。
* ブロックベースのフォーマットでは、ブロックをより小さなパーツに分割することはできませんが、1 つのブロックに含まれる行数は一般設定の [max\_block\_size](/ja/reference/settings/session-settings/max#max_block_size) で制御できます。

## JetStream の使用

NATS JetStream で NATS エンジンを使用する前に、NATS の stream と durable pull コンシューマーを作成する必要があります。これには、たとえば [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="durable pull コンシューマーの作成">
  ```bash theme={null}
  $ nats consumer add
  ? Select a Stream stream_name
  ? Consumer name consumer_name
  ? Delivery target (Pull Consumers の場合は空)
  ? Start policy (all, new, last, subject, 1h, msg sequence) all
  ? Acknowledgment policy explicit
  ? Replay policy instant
  ? Filter Stream by subjects (すべての場合は空)
  ? 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 と durable pull コンシューマーを作成したら、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 テーブルでは少なくとも 1 回の配信が保証されます。メッセージは依存する materialized view への挿入後にのみ確認応答されるため、挿入に失敗したり中断されたりしたメッセージは未確認応答のままとなり、再配信されます。JetStream を使用しない Core NATS には確認応答や再生機能がないため、最大 1 回の配信となり、中断されたメッセージは失われます。

## データ耐久性

このセクションは JetStream にのみ該当します。Core NATS には acknowledgement がなく、上述のとおり at-most-once であるため、acknowledgement 済みのメッセージが失われうる window は存在しません。

JetStream テーブルでは、挿入されたデータがディスクに書き込まれる前に OS page cache が破棄されると、消費済みの行がエラーもなく失われることがあります。バッチが依存先の materialized view にプッシュされると、consumer はそれらのメッセージを acknowledge し、stream はそこから先へ進みます。しかし、挿入された行が耐久化されるのはターゲットのパーツが fsync された時点であり、これはデフォルトでは同期的に行われません (`fsync_after_insert = 0`)。acknowledgement の後、ターゲットのパーツが fsync される前に page cache が失われると、メッセージは再配信されないため、エラーが出ないまま行が失われ、`count()` の値が単に小さくなるだけになります。単純なプロセスの kill ではこの問題は表面化しません。kernel が page cache を保持し、最終的に書き戻すためです。顕在化するのは page cache 自体が失われた場合で、デバイスレベルの電源喪失や、ホストまたは kernel の不正なリセットがその例です。

推奨される materialized view 経由の consumption 経路 (挿入 pipeline 全体が完了した後にのみ acknowledgement が送信される) では、ターゲットの `MergeTree` テーブルに `fsync_after_insert = 1` (および `fsync_part_directory = 1`) を設定することで、acknowledgement の送信前に挿入されたパーツが耐久化され、この window を大幅に狭められます。この設定は、cascade された materialized view のターゲットを含め、バッチの挿入先となるすべての `MergeTree` テーブルで有効にする必要があります。1 つでもデフォルトのままのテーブルがあれば、そのパーツは依然として失われる可能性があります。非同期の中間層は、この設定だけでは耐久性を得られません。たとえば `Distributed` ターゲットは `distributed_foreground_insert = 0` (ClickHouse Cloud 以外でのデフォルト) の場合バックグラウンドで挿入するため、独自の耐久性設定または同期挿入が必要です。また、この緩和策は `nats_commit_on_select = 1` を伴う直接的な `INSERT ... SELECT ... FROM <nats_table>` には適用されません。この場合、メッセージが acknowledge されるのは、宛先が耐久性のあるパーツを書き込んだ後ではなく、読み取りが終端に達した時点だからです。
