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 を含む.credsfile と同じペイロード)。クエリで使用できる唯一の指定方法であるため、<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に保存されます) 。
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のいずれかを選択して、そこにデータをパブリッシュできます。例:
説明
SELECT は、メッセージの読み取りにはあまり適していません (debugging 目的を除く) 。各メッセージは一度しか読み取れないためです。より実用的なのは、materialized view を使用してリアルタイムのスレッドを作成することです。これを行うには、次の手順に従います。
- engine を使用して NATS コンシューマーを作成し、それをデータストリームとして扱います。
- 必要な structure を持つ table を作成します。
- engine からのデータを変換し、あらかじめ作成した table に格納する materialized view を作成します。
MATERIALIZED VIEW を engine に接続すると、バックグラウンドでデータの収集を開始します。これにより、NATS からメッセージを継続的に受信し、SELECT を使って必要なフォーマットに変換できます。
1 つの NATS table には、必要な数だけ materialized view を作成できます。これらは table から直接データを読み取るのではなく、新しいレコードをブロック単位で受け取ります。そのため、詳細度の異なる複数の table に書き込むことができます (グループ化あり - aggregation、なし) 。
例:
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 の作成
stream の作成
durable pull コンシューマーの作成
durable pull コンシューマーの作成
データ耐久性
このセクションは 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 されるのは、宛先が耐久性のあるパーツを書き込んだ後ではなく、読み取りが終端に達した時点だからです。