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

ブリッジ - MQTT

MQTTブリッジを使うと、machbase-neoと外部のMQTTブローカーの間でメッセージを送受信できます。

📢
MQTTベースのプラットフォームにmachbase-neoを導入する場合、ブリッジを接続すれば、既存システムを変更する必要はありません。
  • 外部MQTTブローカーへのメッセージ送信
    flowchart LR
  machbase-neo --PUBLISH-->external-system
  subgraph machbase-neo
      direction LR
      machbase[("machbase
                  エンジン")] --読み取り--> tql
      tql["TQLスクリプト"] --> bridge("bridge(mqtt)")
  end
  subgraph external-system
    direction LR
    broker[[MQTTブローカー]] --> subscriber["アプリケーション
                                        (購読側)"]
  end
  
  • 外部MQTTブローカーからのメッセージ受信
    flowchart RL
    external-system --PUBLISH--> machbase-neo
    machbase-neo --SUBSCRIBE--> external-system
    subgraph machbase-neo
        direction RL
        bridge("bridge(mqtt)") --> subscriber
        subscriber["TQLスクリプト"] --書き込み--> machbase
        machbase[("machbase
                    エンジン")]
    end
    subgraph external-system
        direction RL
        client["アプリケーション
              (発行側)"] --PUBLISH--> mqtt[["MQTTブローカー"]]
    end
  

外部MQTTブリッジの登録

ブリッジを登録します。

bridge add -t mqtt my_mqtt broker=127.0.0.1:1883 id=client-id;

MQTTブリッジは、machbase-neoから外部ブローカーへの接続方法を定義します。 メッセージの受信については、後述のサブスクライバーを参照してください。

使用できる接続オプション

オプション説明
brokerブローカーアドレス。接続先が冗長化されている場合は、複数の「broker」オプションを使用しますbroker=192.0.1.100:1883
idクライアントID
usernameユーザー名
passwordパスワード
keepalive期間形式のkeepalivekeepalive=30s
cleansessionクリーンセッションcleansession=1 cleansession=false
cafileCA証明書(*.pem)のファイルパスTLS
keyクライアント秘密鍵(*.pem)のファイルパスTLS
certクライアント証明書(*.pem)のファイルパスTLS

cafilekeycertをすべて指定すると、TLSによる安全なMQTT接続が有効になります。

メッセージの送信

まず、mosquitto_subをデバッグモード(-d)で実行します。 machbase-neoがneo/messagesトピックにメッセージを発行すると、ブローカー経由で受信します。

mosquitto_sub -d -h 127.0.0.1 -p 1883 -i client-app -t neo/messages
Client client-app sending CONNECT
Client client-app received CONNACK (0)
Client client-app sending SUBSCRIBE (Mid: 1, Topic: neo/messages, QoS: 0, Options: 0x00)
Client client-app received SUBACK
Subscribed (mid: 1): 0

ブリッジのpublish()関数を呼び出すTQLスクリプトを作成します。

TIMER この例は、簡単にするためFAKE()で手動実行します。 ブリッジのpublish機能は、タイマーと組み合わせて自動的にデータを送信する場合に便利です。

現在のJavaScriptランタイムでは、次のTQLで同じ外部ブローカーに送信できます。 この例はmqttモジュールでクライアント接続を作成します。 登録済みブリッジのpublish()を使う旧Tengoの例は、後述の参照用コードに残しています。

SCRIPT({
  const mqtt = require('mqtt');
  const values = [0, 2.5, 5, 7.5, 10];
  const client = new mqtt.Client({servers: ['tcp://127.0.0.1:1883']});
  let published = 0;
  const timeout = setTimeout(() => client.close(), 5000);
  client.on('open', () => {
    for (const value of values) {
      client.publish('neo/messages', 'The message number is ' + value, {qos: 1});
      $.yield(value);
    }
  });
  client.on('published', () => {
    if (++published === values.length) {
      clearTimeout(timeout);
      client.close();
    }
  });
  client.on('error', err => {
    console.error(err);
    clearTimeout(timeout);
    client.close();
  });
})
CSV()
旧Tengoの例(現在のランタイムでは実行できません)
1
2
3
4
5
6
7
8
FAKE(linspace(0,10, 5))
SCRIPT("tengo", {
  ctx := import("context")
  br := ctx.bridge("my_mqtt")
  br.publish("neo/messages", "The message number is "+ctx.value(0))
  ctx.yieldKey(ctx.key(), ctx.value()...)
})
CSV()

スクリプトを実行すると、mosquitto_subが受信したメッセージを直ちに出力します。

mosquitto_sub -d -h 127.0.0.1 -p 1883 -i client-app -t neo/messages
... omit ...
Client client-app received PUBLISH (d0, q0, r0, m0, 'neo/messages', ... (23 bytes))
The message number is 0
Client client-app received PUBLISH (d0, q0, r0, m0, 'neo/messages', ... (25 bytes))
The message number is 2.5
Client client-app received PUBLISH (d0, q0, r0, m0, 'neo/messages', ... (23 bytes))
The message number is 5
Client client-app received PUBLISH (d0, q0, r0, m0, 'neo/messages', ... (25 bytes))
The message number is 7.5
Client client-app received PUBLISH (d0, q0, r0, m0, 'neo/messages', ... (24 bytes))
The message number is 10

メッセージ受信とサブスクライバー

次に、MQTTブローカーから受信したメッセージを、ブリッジ経由でデータベースに保存する例を示します。 デモでは、mosquittoをブローカー、mosquitto_pubをMQTTクライアントとして使用し、外部システムを模擬します。

    flowchart RL
    external-system --PUBLISH--> machbase-neo
    machbase-neo --SUBSCRIBE--> external-system
    subgraph machbase-neo
        direction RL
        bridge("bridge(mq)") --> subscriber
        subscriber["mqttsubr.tql"] --書き込み--> machbase
        machbase[("machbase
                    エンジン")]
    end
    subgraph external-system
        direction RL
        client["mosquitto_pub"] --PUBLISH--> mqtt[["mosquitto"]]
    end
  

1. MQTTブローカーの起動

machbase-neoのMQTTブリッジは、MQTT v3.1.1に準拠するブローカーと互換性があります。 ブローカーがない場合は、デモ用にmosquittoをインストールして実行してください。 https://mosquitto.org

$ mosquitto -p 1883

1691466522: mosquitto version 2.0.15 starting
1691466522: Using default config.
1691466522: Starting in local only mode. Connections will only be possible from clients running on this machine.
1691466522: Create a configuration file which defines a listener to allow remote access.
1691466522: For more details see https://mosquitto.org/documentation/authentication-methods/
1691466522: Opening ipv4 listen socket on port 1883.
1691466522: Opening ipv6 listen socket on port 1883.
1691466522: mosquitto version 2.0.15 running

2. ブリッジの登録

machbase-neoシェルで、以下のコマンドを実行してブリッジを追加します。

bridge add -t mqtt my_mqtt broker=127.0.0.1:1883 id=demo;

このコマンドは、指定したブローカーへの接続方法を定義します。

machbase-neo» bridge list;
╭─────────┬──────────┬─────────────────────────────────╮
│ NAME    │ TYPE     │ CONNECTION                      │
├─────────┼──────────┼─────────────────────────────────┤
│ my_mqtt │ mqtt     │ broker=127.0.0.1:1883 id=demo   │
╰─────────┴──────────┴─────────────────────────────────╯

my_mqttブリッジの登録に成功すると、machbase-neoがブローカーに接続し、 mosquittoログに、以下の接続記録が表示されます。 ネットワーク障害やブローカー障害があっても、machbase-neoは定期的に再接続を試みます。

1691466529: New connection from 127.0.0.1:65440 on port 1883.
1691466529: New client connected from 127.0.0.1:65440 as demo (p2, c1, k30).

3-A. 書き込み記述子を使用するサブスクライバー

ブリッジとテーブルを関連付けるサブスクライバーを登録します。

subscriber add --autostart mqtt_subr my_mqtt iot/sensor db/append/EXAMPLE:csv;

subscriber listで登録を確認します。

┌───────────┬─────────┬────────────┬───────────────────────┬───────────┬─────────┐
│ NAME      │ BRIDGE  │ TOPIC      │ DESTINATION           │ AUTOSTART │ STATE   │
├───────────┼─────────┼────────────┼───────────────────────┼───────────┼─────────┤
│ MQTT_SUBR │ my_mqtt │ iot/sensor │ db/append/EXAMPLE:csv │ true      │ RUNNING │
└───────────┴─────────┴────────────┴───────────────────────┴───────────┴─────────┘

各引数の意味は以下のとおりです。

  • --autostart: machbase-neoの起動時にサブスクライバーを自動起動します。省略すると、手動で開始・停止できます。
  • mqtt_subr: サブスクライバー名です。
  • my_mqtt: 使用するブリッジ名です。
  • iot/sensor: 購読するトピック(MQTTのトピック構文を使用)。
  • db/append/EXAMPLE:csv: 書き込み記述子です。入力データがCSVで、EXAMPLEテーブルにappendモードで書き込むことを示します。

書き込み記述子の代わりに、TQLスクリプトのパスも指定できます。後半で例を示します。

書き込み記述子の形式は以下のとおりです。

db/{method}/{table_name}:{format}:{compress}?{options}

method

方式はappendwriteの2つです。ストリーミング環境では、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で先頭行をスキップ
  • db/append/EXAMPLE:csv?timeformat=s&heading=true
  • db/write/EXAMPLE:csv:gzip?timeformat=s

mosquitto_pubでのメッセージ発行

以下のdata.csvファイルを用意します。

mqtt-demo.temp,1691470297923000000,34.1
mqtt-demo.humidity,1691470297923000000,67.8

mosquitto_pubで、data.csvをMQTTブローカーに発行します。

mosquitto_pub -d -h 127.0.0.1 -p 1883 -t iot/sensor -f data.csv

保存したデータを検索します。

machbase-neo» select * from example where name in ('mqtt-demo.temp', 'mqtt-demo.humidity');
╭────────┬────────────────────┬─────────────────────────┬───────────╮
│ ROWNUM │ NAME               │ TIME(LOCAL)             │ VALUE     │
├────────┼────────────────────┼─────────────────────────┼───────────┤
1 │ mqtt-demo.temp     │ 2023-08-08 13:51:37.923 │ 34.100000 │
2 │ mqtt-demo.humidity │ 2023-08-08 13:51:37.923 │ 67.800000 │
╰────────┴────────────────────┴─────────────────────────┴───────────╯

3-B. TQLを使用するサブスクライバー

データ書き込み用TQLスクリプト

machbase-neoのTQLエディターで、以下のコードをmqttsubr.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 --qos 1 mqttsubr my_mqtt iot/sensor /mqttsubr.tql;

各オプションの意味は以下のとおりです。

  • --autostart: machbase-neoとともに自動起動
  • --qos 1: QoS 1で購読(MQTTブリッジはQoS 0と1に対応)
  • mqttsubr: サブスクライバー名
  • my_mqtt: 使用するブリッジ名
  • iot/sensor: 購読トピック。#+などの標準MQTTトピック構文に対応

--autostartを指定したため、登録したサブスクライバーがRUNNING状態であることを確認します。

machbase-neo» subscriber list;
╭──────────┬─────────┬────────────┬───────────────┬───────────┬─────────╮
│ NAME     │ BRIDGE  │ TOPIC      │ TQL           │ AUTOSTART │ STATE   │
├──────────┼─────────┼────────────┼───────────────┼───────────┼─────────┤
│ MQTTSUBR │ my_mqtt │ iot/sensor │ /mqttsubr.tql │ true      │ RUNNING │
╰──────────┴─────────┴────────────┴───────────────┴───────────┴─────────╯

mosquitto_pubでのメッセージ発行

前述と同じdata.csvを使用します。

mqtt-demo.temp,1691470297923000000,34.1
mqtt-demo.humidity,1691470297923000000,67.8

mosquitto_pubでデータを発行します。

mosquitto_pub -d -h 127.0.0.1 -p 1883 -t iot/sensor -f data.csv

保存したデータを検索します。

machbase-neo» select * from example where name in ('mqtt-demo.temp', 'mqtt-demo.humidity');
╭────────┬────────────────────┬─────────────────────────┬───────────╮
│ ROWNUM │ NAME               │ TIME(LOCAL)             │ VALUE     │
├────────┼────────────────────┼─────────────────────────┼───────────┤
1 │ mqtt-demo.temp     │ 2023-08-08 13:51:37.923 │ 34.100000 │
2 │ mqtt-demo.humidity │ 2023-08-08 13:51:37.923 │ 67.800000 │
╰────────┴────────────────────┴─────────────────────────┴───────────╯
最終更新日