コンテンツにスキップ
ブリッジ - NATS

ブリッジ - 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」オプションを指定します
NameCONNECT時に、クライアントを識別するためサーバーに送る、省略可能な名前ラベルです。
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

方式はappendwriteの2つです。NATSなどのストリーミング環境では、appendを推奨します。

  • append: appendモードで書き込み
  • write: INSERT SQLで書き込み

table_name

対象テーブル名を指定します。大文字と小文字は区別しません。

format

  • json(既定)
  • csv

compress

現在はgzipに対応しています。:{compress}を省略すると、データを圧縮しません。

options

?の後に、URLエンコードした追加オプションを指定できます。

名前既定値説明
timeformatns時刻の形式:s、ms、us、ns
tzUTCタイムゾーン:UTC、Local、地域指定
delimiter,CSV区切り文字。CSV以外の場合は無視
headingfalseCSVにヘッダーがある場合、trueで先頭行をスキップ

サブスクライバーの保留メッセージの上限は、nats.ioのドキュメントを参照してください。

例:

  • db/append/EXAMPLE:csv?timeformat=s&heading=true
  • db/write/EXAMPLE:csv:gzip?timeformat=s
  • db/append/EXAMPLE:json?timeformat=s&pendingMsgLimit=1048576

NATSクライアントアプリケーション

NATSサーバーのiot.sensorサブジェクトにCSVデータを送信する、簡単なGoアプリケーションを作成します。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
package main

import (
	"fmt"
	"strings"
	"time"

	"github.com/nats-io/nats.go"
)

func main() {
    // NATSサーバーに接続
	opts := nats.GetDefaultOptions()
	opts.Servers = []string{"nats://127.0.0.1:4222"}
	conn, err := opts.Connect()
	if err != nil {
		panic(err)
	}
	defer conn.Close()

	tick := time.Now()

    // CSVデータを作成
    lines := []string{}
	for i := 0; i < 10; i++ {
        // NAME,TIME,VALUE
		line := fmt.Sprintf("hello-nats,%d,3.1415", tick.Add(time.Duration(i)).UnixNano())
		lines = append(lines, line)
	}
	reqData := []byte(strings.Join(lines, "\n"))

	// A) 要求・応答方式
	if rsp, err := conn.Request("iot.sensor", reqData, 100*time.Millisecond); err != nil {
		panic(err)
	} else {
		fmt.Println("RESP:", string(rsp.Data))
	}
	// B) 応答を待たない送信方式
	// if err := conn.Publish("iot.sensor", reqData); err != nil {
	// 	panic(err)
	// }
}

実行すると、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として保存します。

1
2
3
4
CSV(payload())
MAPVALUE(1, parseTime(value(1), "ns"))
MAPVALUE(2, parseFloat(value(2)))
APPEND( table("example") )

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)
}
最終更新日