DatalayersへのMQTTデータ取り込み
Datalayersは、産業用IoT、IoV、エネルギーなどの分野向けに設計されたマルチモーダルかつハイパーコンバージドなデータベースです。高いデータスループットと安定したパフォーマンスを備えており、IoTアプリケーションに最適です。EMQXは現在、Sinkを介してメッセージやデータをDatalayersに格納することをサポートしており、データ分析や可視化を容易にします。
本ページでは、EMQXとDatalayersのデータ統合について詳細に解説し、ルールおよびSinkの作成方法を実践的に案内します。
動作概要
Datalayersデータ統合はEMQXの標準機能であり、デバイスからのMQTTメッセージをDatalayersへシームレスに転送し、保存および分析を可能にします。ルールとSinkを設定することで、処理済みのMQTTデータを柔軟にDatalayersへルーティングできます。
以下の図は、エネルギー貯蔵シナリオにおけるEMQXとDatalayersの典型的な連携アーキテクチャを示しています:

このアーキテクチャでは、EMQXがデバイスの接続管理、メッセージ転送、ルールベースの処理を担当し、Datalayersがデータの保存、分析、可視化を担います。両者が連携することで、エネルギー消費のリアルタイムデータを効率的に収集・分析するスケーラブルなIoTプラットフォームを構築できます。
EMQX 6.0.0以降、DatalayersはApache Arrowをベースにした高性能バイナリ通信プロトコルであるArrow Flight SQLをサポートしています。従来のInfluxDB Line Protocolと比較して、Arrow Flight SQLはより効率的なデータ転送と構造化データ書き込みの強力なサポートを提供します。
注意
Arrow FlightドライバーはRustで実装され、Erlang VMにNative Implemented Function(NIF)を通じて統合されています。この機能は現在実験的であり、テスト環境での利用を推奨します。
具体的なワークフローは以下の通りです:
メッセージのパブリッシュと受信:デバイスはMQTT経由でEMQXに接続し、電力、電流、電圧などのエネルギー関連メトリクスを定期的にパブリッシュします。EMQXはこれらのメッセージを受信し、ルールエンジンに渡します。
ルールエンジンによるメッセージ処理:EMQXの組み込みルールエンジンはトピックパターンに基づいてメッセージをマッチングし、ペイロードの変換、フィールドのフィルタリング、コンテキスト情報の付加などの処理を行います。
Datalayersへの書き込み:ルールがトリガーされると、処理済みデータをDatalayersに書き込むSinkアクションが実行されます。SinkはフィールドをDatalayersのテーブルやカラムにマッピングするカスタマイズ可能なSQLテンプレートをサポートします。
EMQXは以下の2つの書き込み方式をサポートしています:
- InfluxDB Line Protocol
- Arrow Flight SQLドライバー
Sinkの設定は選択した方式により異なります。
エネルギー貯蔵データがDatalayersに書き込まれた後は、対応ツールを利用してデータ分析が柔軟に行えます。例えば:
- Grafanaなどの可視化ツールに接続し、エネルギー貯蔵データのチャートを生成・表示する。
- 業務システムと連携し、エネルギー貯蔵デバイスの状態監視やアラートを実施する。
特長と利点
Datalayersデータ統合は以下の特長と利点を提供します:
- 効率的なデータ処理:EMQXは多数のIoTデバイス接続とメッセージスループットを処理可能であり、Datalayersはデータ書き込み、保存、クエリに優れているため、システムに過負荷をかけずにIoTシナリオのデータ処理要件を満たします。
- メッセージ変換:メッセージはEMQXルール内で大規模な処理・変換が可能であり、Datalayersに書き込む前に柔軟に加工できます。
- スケーラビリティ:EMQXとDatalayersは共にクラスター対応しており、ビジネスの成長に応じて水平スケールが可能です。
- 豊富なクエリ機能:Datalayersはタイムスタンプデータの効率的なクエリと分析のために最適化された関数、演算子、インデックス技術を提供し、IoT時系列データから価値ある洞察を抽出します。
- 効率的なストレージ:Datalayersは高圧縮エンコード方式を採用し、ストレージコストを大幅に削減します。また、データ保持期間のカスタマイズが可能で、不要なデータがストレージを占有するのを防ぎます。
はじめる前に
このセクションでは、EMQXでDatalayers Sinkを作成する前の準備として、Datalayersのインストール、データベース作成、テーブル構造定義について説明します。
前提条件
- ルールの理解
- データ統合の理解
- 書き込みに使用するドライバータイプに応じて、InfluxDB Line ProtocolまたはArrow Flight SQLの理解
Datalayersのインストールとセットアップ
Dockerを使用してDatalayersをインストールし起動します。詳細手順はInstall Datalayersを参照してください。
bash# Datalayersコンテナを起動 docker run -d --name datalayers -p 8360:8360 -p 8361:8361 datalayers/datalayers:latest- ポート
8360はArrow Flight SQLのデフォルトgRPCポートです。 - ポート
8361はHTTPポートで、主にLine Protocol書き込みや管理APIに使用されます。
- ポート
Datalayersサービス起動後、デフォルトのユーザー名・パスワード
admin/publicでDatalayers CLIにログインし、データベースを作成します。Datalayersコンテナにアクセス:
bashdocker exec -it datalayers bashDatalayers CLIを起動:
bashdlsql -u admin -p publicデータベースを作成(例:
mqtt):sqlcreate database mqtt
Arrow Flight SQLドライバーを使用する場合は、対象テーブルを事前に作成する必要があります。
注意
InfluxDB Line Protocolを使用する場合は、テーブルの事前作成は不要です。Datalayersは受信したLine Protocolデータの
measurementおよびフィールド定義に基づき自動的にテーブルを作成します。例として、
t_mqtt_msgというテーブルを以下のSQLで作成します:sqlCREATE TABLE IF NOT EXISTS `t_mqtt_msg` ( time TIMESTAMP(3) NOT NULL, msgid STRING NOT NULL, sender STRING NOT NULL, topic STRING NOT NULL, qos INT8 NOT NULL, payload STRING, arrived TIMESTAMP(3) NOT NULL, timestamp key(time) ) PARTITION BY HASH (msgid, sender) PARTITIONS 1 ENGINE = TimeSeries WITH (ttl = '14d');
Datalayersコネクターの作成
このセクションでは、EMQXでDatalayersサーバーに接続するためのコネクター作成方法を説明します。
以下の手順はEMQXとDatalayersがローカルで稼働していることを前提としています。別環境やリモート環境にデプロイしている場合は、接続設定を適宜更新してください。
EMQXダッシュボードで、Integration -> Connectorsをクリックします。
ページ右上のCreateをクリックします。
Create ConnectorページでDatalayersを選択し、Nextをクリックします。
Configurationページでコネクターの詳細を入力します:
- Connector Name:英数字で始まり、英数字、ハイフン、アンダースコアのみ使用可能。例:
my_datalayers - Description(任意):後で識別しやすいよう説明を追加可能
Datalayersサーバー接続設定:
Driver Type:
InfluxDB Line Protocol:InfluxDB互換のLine Protocolでデータを取り込み。テーブルは自動作成されます。Arrow Flight:SQLテンプレートを用いた高性能な構造化データ書き込みを可能にします。スキーマ制御が厳格で高い書き込みスループットが必要な場合に最適です。注意
Arrow FlightドライバーはRustで実装され、Erlang VMにNIFで統合されています。現在実験的機能であり、テスト環境での評価を推奨します。
Server Host:
- デフォルト:
127.0.0.1:8361 Arrow Flightドライバー使用時はgRPC通信のためポート8360を使用します。
- デフォルト:
Database Name:Datalayers上の対象データベース名(例:
mqtt)Username / Password:Datalayersアクセス用認証情報(例:
admin/public)Enable TLS(任意):暗号化接続を有効化。証明書パスや検証オプションの設定が可能です。詳細は外部リソースアクセスのTLS有効化を参照してください。
注意
Arrow Flight SQLプロトコル使用時は、証明書検証のスキップ(
verify_none)はライブラリの制約によりサポートされません。gRPCサーバー証明書のCommon Name(CN)がサーバーホストと一致していることを確認してください。
- Connector Name:英数字で始まり、英数字、ハイフン、アンダースコアのみ使用可能。例:
ドライバーに
Arrow Flightを選択すると、Enable Prepared Statementsオプションが表示されます。これはSinkがSQLテンプレートを使用してデータ挿入を行うかを制御し、デフォルトで有効です。Createをクリックする前に、Test ConnectivityでDatalayersサーバーへの接続確認が可能です。
ページ下部のCreateをクリックしてコネクター作成を完了します。ポップアップでBack to Connector ListまたはCreate Ruleを選択できます。ルールとSinkの作成手順はCreate a Datalayers Ruleを参照してください。
Datalayersルールの作成
このセクションでは、EMQXでソーストピックt/#からのMQTTメッセージを処理し、設定済みのSinkを使ってDatalayersに送信するルールの作成方法を示します。
SQLを定義したルールの作成
EMQXダッシュボードの左メニューからData Integration -> Rulesに移動します。
Rulesページ右上のCreateボタンをクリックします。
ルール作成フォームでルールID(例:
my_rule)を入力します。SQL Editorでルールロジックを定義します。トピック
t/#にパブリッシュされたMQTTメッセージをDatalayersに保存するには、以下のSQLを使用できます:注意
カスタムSQLルールを書く場合、Sinkテンプレートで参照するすべての変数(例:
${clientid},${payload.temp})がルールのSELECT句に含まれていることを確認してください。SELECT * FROM "t/#"TIP
EMQXのSQLに不慣れな場合は、SQL ExamplesやEnable Debugをクリックしてサンプルクエリを試し、出力を確認できます。
ルールにDatalayers Sinkを追加し、処理結果をDatalayersに書き込みます。
- InfluxDB Line Protocolを使用する場合は、Add an InfluxDB Line Protocol Sinkを参照してください。
- Arrow Flight SQLドライバーを使用する場合は、Add an Arrow Flight SQL Sinkを参照してください。
Create Ruleページで設定内容を確認し、Saveをクリックしてルールを作成します。
作成したルールはRules一覧に表示されます。対象ルールの**Actions (Sink)**タブをクリックすると、関連するDatalayers Sinkを確認できます。
また、Integrations -> Flow Designerでトポロジーグラフを表示すると、トピックt/#のメッセージがmy_ruleルールで処理され、Datalayersに書き込まれている様子が視覚的に確認できます。
InfluxDB Line Protocol Sinkの追加
このセクションでは、InfluxDB Line Protocolを用いて処理済みデータをDatalayersに書き込むSinkをルールに追加する方法を説明します。
ルールエディター右側のAdd Actionボタンをクリックし、ルール条件にマッチした際にトリガーされるアクションを定義します。このアクションが処理済みメッセージをDatalayersに転送します。
Type of Actionドロップダウンで
Datalayersを選択し、ActionはデフォルトのCreate Actionのままにします。既存のDatalayers Sinkを選択することも可能ですが、本例では新規作成を想定しています。Sinkの名前(例:
dl_sink_influx)を入力します。名前は英数字の組み合わせで構いません。Connectorドロップダウンから、
InfluxDB Line Protocolドライバーで設定済みのコネクターを選択します。利用可能なコネクターがない場合は隣のボタンから新規作成してください。詳細はCreate a Datalayers Connectorを参照。Time Precisionはデフォルトでミリ秒に設定します。
Datalayersへのデータ解析・書き込みに用いるData Formatと内容を定義します。
JSONまたはLine Protocolを選択可能です:JSON:
Measurement、Fields、Timestamp、Tagsを指定します。キーと値は定数またはプレースホルダー(例:
${payload.temp})が利用可能です。書式ルールはInfluxDB Line Protocolを参照してください。FieldsはCSVファイルを使った一括設定もサポートしています。詳細はUse CSV to Batch Configure Fieldsを参照してください。
Line Protocol:
テーブル、フィールド、タイムスタンプ、タグを含む単一のLine Protocol文字列を定義できます。キーと値は定数またはプレースホルダーが利用可能です。構文はInfluxDB Line Protocolを参照してください。
TIP
Datalayersに書き込むデータはInfluxDB v1のLine Protocolと完全互換のため、InfluxDB Line Protocolも参考にできます。
例えば、符号付き整数値を入力する場合、プレースホルダーの後に
iを付けます(例:${payload.int}i)。詳細はInfluxDB 1.8で整数値を書き込む方法を参照してください。Line Protocolの例:
sqldevices,clientid=${clientid} temp=${payload.temp},hum=${payload.hum},precip=${payload.precip}i ${timestamp}
Fallback Actions(任意):メッセージ配信失敗時の信頼性向上のため、フォールバックアクションを1つ以上設定可能です。詳細はFallback Actionsを参照してください。
Advanced Settingsを展開し、必要に応じて詳細設定を行います。詳細はAdvanced Settingsを参照してください。
Createをクリックする前に、Test ConnectivityでSinkがDatalayersサーバーに接続できるかテスト可能です。
CreateをクリックしてSink作成を完了します。Create Ruleページに戻ると、Action Outputsタブに新規Sinkが表示されます。
CSVを使ったフィールド一括設定
TIP
この機能はInfluxDB Line ProtocolのSinkで、データフォーマットがJSONの場合にのみ利用可能です。フィールド設定を一括インポートできます。
Datalayersのデータエントリーは数百のフィールドを含むことが多く、データフォーマット設定が煩雑になることがあります。EMQXはこれを解決するため、フィールドの一括設定機能を提供しています。
JSON形式でデータフォーマットを設定する際、CSVファイルからフィールドのキー・値ペアを一括インポート可能です。
FieldsテーブルのBatch Settingsボタンをクリックし、Import Batch Settingsポップアップを開きます。
指示に従いテンプレートファイルをダウンロードし、フィールドのキー・値ペアを記入します。テンプレートのデフォルト内容は以下の通りです:
Field Value 備考(任意) temp ${payload.temp} hum ${payload.hum} precip ${payload.precip}i フィールド値の後ろに iを付けるとDatalayersは整数型として保存- Field:フィールドキー。定数または
${var}形式のプレースホルダーをサポート。 - Value:フィールド値。定数またはプレースホルダーをサポートし、Line Protocolに従い型識別子を付加可能。
- 備考:CSV内のコメント用で、EMQXへのインポート対象外。
バッチ設定CSVファイルは最大2048行までです。
- Field:フィールドキー。定数または
記入済みテンプレートファイルを保存し、Import Batch SettingsポップアップにアップロードしてImportをクリックし、一括設定を完了します。
インポート後、Fields設定テーブルでキー・値ペアをさらに調整可能です。
Arrow Flight SQL Sinkの追加
このセクションでは、Arrow Flight SQLドライバーを使用し、SQL挿入文でDatalayersにデータを書き込むSinkをルールに追加する方法を説明します。
注意
Arrow Flight SQLドライバーは現在実験的です。商用環境での利用は慎重に行ってください。
ルールエディター右側のAdd Actionボタンをクリックし、ルールマッチ時にトリガーされるアクションを定義します。このアクションが処理済みデータをDatalayersに転送します。
Type of Actionドロップダウンで
Datalayersを選択し、ActionはデフォルトのCreate Actionのままにします。既存のDatalayers Sinkを選択することも可能ですが、本例では新規作成を想定しています。Sinkの名前(例:
dl_sink_arrow)を入力します。英数字の組み合わせが推奨されます。Connectorドロップダウンから、
Arrow Flightドライバーで設定済みのコネクターを選択します。存在しない場合は隣のボタンから新規作成してください。詳細はCreate a Datalayers Connectorを参照。データを対象テーブルに挿入する方法を定義するSQLテンプレートを設定します。
TIP
これはプリプロセッシングSQLテンプレートです。フィールド名を引用符で囲まず、SQL文の末尾にセミコロン
;を含めないでください。すべての${}プレースホルダーはルールSQLで選択したフィールドと一致させる必要があります。TIP
コネクターで設定したデータベース以外にデータを挿入する場合は、SQLテンプレート内で対象データベース名を明示的に指定してください。なお、コネクターは対象データベースの存在を引き続きチェックします。
例:
sqlinsert into t_mqtt_msg(time, msgid, sender, topic, qos, payload, arrived) values (${timestamp}, ${id}, ${clientid}, ${topic}, ${qos}, ${payload}, ${timestamp})Fallback Actions(任意):信頼性向上のため、Sinkがメッセージ処理に失敗した場合にトリガーされるフォールバックアクションを1つ以上設定可能です。詳細はFallback Actionsを参照してください。
Advanced Settingsを展開し、必要に応じて詳細設定を行います。詳細はAdvanced Settingsを参照してください。
Createをクリックする前に、Test ConnectionボタンでSinkがDatalayersサーバーに接続できるか検証可能です。
CreateをクリックしてSink作成を完了します。Create Ruleページに戻ると、Action Outputsタブに新規Sinkが表示されます。
ルールとSinkのテスト
ルールとSinkの設定後、テスト用MQTTメッセージをパブリッシュしてDatalayersへのデータ書き込みが成功しているか確認できます。
MQTTXを使い、トピック
t/1にメッセージを送信します。これによりセッションイベント(クライアントのオンライン/オフラインなど)がトリガーされる場合もあります:bashmqttx pub -i emqx_c -t t/1 -m '{ "temp": "23.5", "hum": "62", "precip": 2 }'このメッセージはルールエンジンをトリガーし、設定済みのDatalayers Sinkに転送されます。ルールにクライアント接続・切断などのセッションイベントが含まれている場合も、この操作でトリガーされます。
Sinkの実行統計を確認します。EMQXダッシュボードのRulesページで対象ルールを探し、Actions (Sink)タブに切り替えます。対象SinkのMatchedおよびSuccessカウントが1増加していることを確認してください。
CLIを使ってDatalayers内のデータを検証します。
Datalayersコンテナにアクセスし、CLIツールを起動します:
bashdocker exec -it datalayers bash dlsql -u admin -p public書き込み方式に応じてSQLクエリを実行します:
InfluxDB Line Protocol使用時は、Sink設定の
measurementで指定したテーブル名(例:devices)がデフォルトです:sqluse mqtt select * from devicesArrow Flight SQL使用時は、事前作成した対象テーブル(例:
t_mqtt_msg)をクエリします:sqluse mqtt select * from t_mqtt_msg
詳細設定
このセクションでは、DatalayersコネクターおよびSinkの詳細設定オプションについて説明します。ダッシュボードでコネクターやSinkを設定する際、Advanced Settingsを展開して以下のパラメータをニーズに応じて調整できます。
| フィールド名 | 説明 | デフォルト |
|---|---|---|
| Buffer Pool Size | バッファワーカープロセスの数を指定します。これらのプロセスはEMQXとDatalayersのエグレス型Sink間のデータフローを管理し、データを一時的に保存・処理してからターゲットサービスに送信します。エグレスシナリオでのパフォーマンス最適化やスムーズなデータ送信に重要です。イングレスのみを扱うブリッジでは適用されないため0に設定可能です。 | 4 |
| Request TTL | リクエストTTL(Time to Live)は、リクエストがバッファに入ってから有効とみなされる最大時間(秒)を指定します。TTLを超えたリクエストや、送信済みだがDatalayersからの応答・アックがタイムリーに得られない場合、そのリクエストは期限切れと判断されます。 | 45 |
| Health Check Interval | SinkがDatalayersとの接続状態を自動的にヘルスチェックする間隔(秒)を指定します。 | 15 |
| Max Buffer Queue Size | Datalayers Sinkの各バッファワーカープロセスがバッファリング可能な最大バイト数を指定します。バッファワーカーはデータを一時保存し、効率的にデータストリームを処理します。システム性能やデータ送信要件に応じて調整してください。 | 1 |
| Batch Size | EMQXからDatalayersへ単一転送操作で送信するデータバッチの最大サイズを指定します。これによりデータ転送の効率とパフォーマンスを調整可能です。Batch Sizeが1の場合、データレコードはバッチ化されず個別に送信されます。 | 100 |
| Query Mode | synchronousまたはasynchronousのリクエストモードを選択し、メッセージ送信を最適化します。非同期モードではDatalayersへの書き込みがMQTTメッセージのパブリッシュ処理をブロックしませんが、クライアントがメッセージをDatalayers到達前に受信する可能性があります。 | Asynch |
| Inflight Window | 「インフライトキューリクエスト」は開始済みで応答・アック待ちのリクエストを指します。この設定はSinkとDatalayers間の同時インフライトリクエスト最大数を制御します。Request Modeがasynchronousの場合、同一MQTTクライアントからのメッセージを厳密に順序処理したい場合はこの値を1に設定してください。 | 100 |