ブリッジ - NATS
NATSブリッジ Since v8.0.20 を使うと、machbase-neoとNATSサーバー(https://nats.io)の間でメッセージを送受信できます。
NATSサーバーブリッジの登録
ブリッジを登録します。
bridge add -t nats my_nats server=nats://127.0.0.1:3000 name=client-name;NATSブリッジは、machbase-neoから外部NATSサーバーへの接続方法を定義します。 メッセージの受信については、後述のサブスクライバーを参照してください。
使用できる接続オプションは以下のとおりです。詳細は、NATSの公式ドキュメントを参照してください。
| オプション | 説明 |
|---|---|
Server | サーバーアドレス。接続先が冗長化されている場合は、複数の「server」オプションを指定します |
Name | CONNECT時に、クライアントを識別するためサーバーに送る、省略可能な名前ラベルです。 |
NoRandomize | サーバープールのランダム化を無効にするかどうかを指定します。 |
NoEcho | 同じ接続から送信したメッセージに一致する購読がある場合、サーバーからのエコーバックを無効にするかどうかを設定します。サーバーバージョン1.2以降、Proto 1以降で対応します。 |
Verbose | サーバーが正常に処理したコマンドに、OK ACKを返すように指定します。 |
Pedantic | サブジェクトを追加検証するかどうかをサーバーに指定します。 |
AllowReconnect | 現在のサーバーから切断された場合の再接続処理を有効にします。 |
MaxReconnect | 再接続を中止するまでの最大試行回数を設定します。 |
ReconnectWait | 以前接続していたサーバーへの再接続を試みた後の、待機時間を設定します。 |
Timeout | 接続のDial操作のタイムアウトを設定します。 |
PingInterval | クライアントからサーバーへのping送信間隔です。0以下で無効になります(例:PingInterval=2m)。 |
User | サーバーへの接続時に使用するユーザー名を設定します。 |
Password | サーバーへの接続時に使用するパスワードを設定します。 |
Token | サーバーへの接続時に使用するトークンを設定します。 |
RetryOnFailedConnect | 初期候補のサーバーに接続できない場合、直ちに再接続状態にします。 |
SkipHostLookup | サーバーホスト名のDNS検索をスキップします(例:SkipHostLookup=true)。 |
メッセージ受信とサブスクライバー
NATSサーバーからメッセージを受信し、ブリッジとサブスクライバーを介してデータベースに保存する例を示します。
1. NATSサーバーの起動
NATSサーバーのインストールは、https://nats.io を参照してください。単独実行モードは簡単にインストールできます。
$ nats-server
[61052] 2021/10/28 16:53:38.003205 [INF] Starting nats-server
[61052] 2021/10/28 16:53:38.003329 [INF] Version: 2.6.1
[61052] 2021/10/28 16:53:38.003333 [INF] Git: [not set]
[61052] 2021/10/28 16:53:38.003339 [INF] Name: NDUP6JO4T5LRUEXZUHWXMJYMG4IZAJDNWETTA4GPJ7DKXLJUXBN3UP3M
[61052] 2021/10/28 16:53:38.003342 [INF] ID: NDUP6JO4T5LRUEXZUHWXMJYMG4IZAJDNWETTA4GPJ7DKXLJUXBN3UP3M
[61052] 2021/10/28 16:53:38.004046 [INF] Listening for client connections on 0.0.0.0:4222
[61052] 2021/10/28 16:53:38.004683 [INF] Server is ready
...2. NATSブリッジの登録
machbase-neoシェルで、以下のコマンドを実行します。
bridge add -t nats my_nats server=nats://127.0.0.1:4222 name=demo;このコマンドは、machbase-neoからNATSサーバーへの接続方法を定義します。
┌──────────┬──────────┬──────────────────────────────────────────┐
│ NAME │ TYPE │ CONNECTION │
├──────────┼──────────┼──────────────────────────────────────────┤
│ my_nats │ nats │ server=nats://127.0.0.1:4222 name=demo │
└──────────┴──────────┴──────────────────────────────────────────┘3-A. 書き込み記述子を使用するサブスクライバー
ブリッジとデータベーステーブルを関連付けるサブスクライバーを追加します。
subscriber add --autostart nats_subr my_nats iot.sensor db/append/EXAMPLE:csv;subscriber listで確認します。
┌───────────┬─────────┬────────────┬───────────────────────┬───────────┬─────────┐
│ NAME │ BRIDGE │ TOPIC │ DESTINATION │ AUTOSTART │ STATE │
├───────────┼─────────┼────────────┼───────────────────────┼───────────┼─────────┤
│ NATS_SUBR │ my_nats │ iot.sensor │ db/append/EXAMPLE:csv │ true │ RUNNING │
└───────────┴─────────┴────────────┴───────────────────────┴───────────┴─────────┘各項目の意味は以下のとおりです。
--autostart: machbase-neoの起動時に自動起動します。省略すると、手動で開始・停止できます。nats_subr: サブスクライバー名です。my_nats: 使用するブリッジ名です。iot.sensor: 購読するNATSサブジェクト名です。db/append/EXAMPLE:csv: 書き込み記述子で、CSVデータをEXAMPLEテーブルにappendモードで書き込むことを示します。
書き込み記述子の代わりに、TQLスクリプトのパスも指定できます。後半で例を示します。
書き込み記述子の形式は以下のとおりです。
db/{method}/{table_name}:{format}:{compress}?{options}method
方式はappendとwriteの2つです。NATSなどのストリーミング環境では、appendを推奨します。
append: appendモードで書き込みwrite: INSERT SQLで書き込み
table_name
対象テーブル名を指定します。大文字と小文字は区別しません。
format
json(既定)csv
compress
現在はgzipに対応しています。:{compress}を省略すると、データを圧縮しません。
options
?の後に、URLエンコードした追加オプションを指定できます。
| 名前 | 既定値 | 説明 |
|---|---|---|
timeformat | ns | 時刻の形式:s、ms、us、ns |
tz | UTC | タイムゾーン:UTC、Local、地域指定 |
delimiter | , | CSV区切り文字。CSV以外の場合は無視 |
heading | false | CSVにヘッダーがある場合、trueで先頭行をスキップ |
サブスクライバーの保留メッセージの上限は、nats.ioのドキュメントを参照してください。
例:
db/append/EXAMPLE:csv?timeformat=s&heading=truedb/write/EXAMPLE:csv:gzip?timeformat=sdb/append/EXAMPLE:json?timeformat=s&pendingMsgLimit=1048576
NATSクライアントアプリケーション
NATSサーバーのiot.sensorサブジェクトにCSVデータを送信する、簡単なGoアプリケーションを作成します。
| |
実行すると、iot.sensorサブジェクトに10件のCSVレコードが送信され、
nats_subrサブスクライバーが受信して、EXAMPLEテーブルに保存します。
$ go run nats_pub.go ↵
RESP: {"success":true,"reason":"10 records appended","elapse":"2.186209ms"}3-B. TQLを使用するサブスクライバー
データ書き込み用TQLスクリプト
CSVデータを受信してexampleテーブルに書き込むTQLを作成し、test.tqlとして保存します。
| |
machbase-neoシェルで、ブリッジとTQLスクリプトを関連付けるサブスクライバーを追加します。
subscriber add --autostart nats_subr my_nats iot.sensor /test.tql;各項目の意味は以下のとおりです。
--autostart: machbase-neoの起動時に自動起動nats_subr: サブスクライバー名my_nats: 使用するブリッジ名iot.sensor: 購読するサブジェクト(NATS構文に対応)/test.tql: 受信データを処理するTQLファイルのパス
subscriber listで確認します。
┌───────────┬─────────┬────────────┬─────────────┬───────────┬─────────┐
│ NAME │ BRIDGE │ TOPIC │ DESTINATION │ AUTOSTART │ STATE │
├───────────┼─────────┼────────────┼─────────────┼───────────┼─────────┤
│ NATS_SUBR │ my_nats │ test.topic │ /test.tql │ true │ RUNNING │
└───────────┴─────────┴────────────┴─────────────┴───────────┴─────────┘NATSクライアントアプリケーション
前述のNATSクライアントアプリケーションを、そのまま実行します。
package main
import (
"flag"
"fmt"
"strings"
"time"
"github.com/nats-io/nats.go"
)
// このプログラムを実行する前に、
// 1. machbase-neoサーバーにブリッジを追加
// bridge add -t nats my_nats server=127.0.0.1:4222 name=hello
// 2. サブスクライバーを追加
// subscriber add hello-nats my_nats test.topic db/write/EXAMPLE:csv;
// 3. サブスクライバーを開始
// subscriber start hello-nats
// 4. Run
// go run nats_pub.go -server nats://<ip>:<port> -subject hello
func main() {
optServer := flag.String("server", "nats://127.0.0.1:4222", "nats server address")
optSubject := flag.String("subject", "hello", "subject to subscribe")
optRequest := flag.Bool("request", false, "request-response model")
flag.Parse()
opts := nats.GetDefaultOptions()
opts.Servers = []string{*optServer}
conn, err := opts.Connect()
if err != nil {
panic(err)
}
defer conn.Close()
tick := time.Now()
lines := []string{}
linesPerMsg := 1
msgCount := 1000000
serial := 0
for n := 0; n < msgCount; n++ {
for i := 0; i < linesPerMsg; i++ {
line := fmt.Sprintf("hello-nats,%d,1.2345", tick.Add(time.Duration(serial)*time.Microsecond).UnixNano())
lines = append(lines, line)
serial++
}
reqData := []byte(strings.Join(lines, "\n"))
lines = lines[0:0]
if *optRequest {
// A) 要求・応答方式
if rsp, err := conn.Request(*optSubject, reqData, 100*time.Millisecond); err != nil {
panic(err)
} else {
fmt.Println("RESP:", string(rsp.Data))
}
} else {
// B) 応答を待たない送信方式
if err := conn.Publish(*optSubject, reqData); err != nil {
panic(err)
}
}
}
fmt.Println("msg sent: ", conn.OutMsgs)
}