Skip to main content
このエンジンを使用すると、ClickHouse を NATS と統合できます。 NATS では次のことができます。
  • メッセージ subject をパブリッシュまたはサブスクライブする。
  • 新しいメッセージが利用可能になり次第処理する。

テーブルの作成

パラメータ:
  • nats_url – host:port (例: localhost:4222) ..
  • nats_subjects – NATS table がサブスクライブ/パブリッシュする subject の一覧。foo.*.bar や baz.> のようなワイルドカード subject をサポートします
  • nats_format – メッセージのフォーマット。JSONEachRow など、SQL の FORMAT 関数と同じ記法を使用します。詳細は フォーマット セクションを参照してください。
パラメータ:
  • nats_schema – フォーマットでスキーマ定義が必要な場合に使用する必要があるパラメータです。たとえば、Cap’n Proto では、スキーマファイルへのパスとルート 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_reconnect_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。
  • nats_flush_interval_ms - NATS から読み取ったデータを flush するまでのタイムアウト。デフォルト: 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のいずれかを選択して、そこにデータをパブリッシュできます。例:
nats 関連の設定に加えて、フォーマット設定も追加できます。 例:
NATS server の設定は、ClickHouse の設定ファイルを使用して追加できます。 具体的には、NATS エンジン用の password を追加できます:

説明

SELECT は、メッセージの読み取りにはあまり適していません (debugging 目的を除く) 。各メッセージは一度しか読み取れないためです。より実用的なのは、materialized 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、なし) 。 例:
ストリームデータの受信を停止するか、変換ロジックを変更するには、materialized viewをデタッチします:
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 でサポートされているすべてのフォーマットに対応しています。 1 つの NATS メッセージに含まれる行数は、そのフォーマットが行ベースかブロックベースかによって異なります。
  • 行ベースのフォーマットでは、1 つの NATS メッセージに含める行数を nats_max_rows_per_message の設定で制御できます。
  • ブロックベースのフォーマットでは、ブロックをより小さなパーツに分割することはできませんが、1 つのブロックに含まれる行数は一般設定の max_block_size で制御できます。

JetStream の使用

NATS JetStream で NATS エンジンを使用する前に、NATS の stream と durable pull コンシューマーを作成する必要があります。これには、たとえば NATS CLI パッケージの nats ユーティリティを使用できます。
stream と durable pull コンシューマーを作成したら、NATS エンジンを使用してテーブルを作成できます。そのためには、nats_stream、nats_consumer_name、nats_subjects を設定する必要があります。
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 されるのは、宛先が耐久性のあるパーツを書き込んだ後ではなく、読み取りが終端に達した時点だからです。
最終更新日 2026年9月26日