MySQLへのMQTTデータ取り込み
MySQLは、高い信頼性と安定性を持つ広く利用されているリレーショナルデータベースであり、迅速にインストール、設定、利用が可能です。MySQLデータ統合により、MQTTメッセージを効率的にMySQLデータベースに保存できるほか、イベントトリガーを通じてMySQL内のデータをリアルタイムに更新または削除することもサポートしています。MySQLデータ統合を活用することで、メッセージの保存、デバイスのオンライン/オフライン状態の更新、デバイスの動作記録などの機能を簡単に実装し、柔軟なIoTデータストレージおよびデバイス管理機能を実現できます。
本ページでは、EMQXとMySQL間のデータ統合について、実践的な作成および検証手順を紹介します。
動作概要
MySQLデータ統合はEMQXに標準搭載された機能であり、簡単な設定で複雑なビジネス開発を可能にします。典型的なIoTアプリケーションにおいて、EMQXはIoTプラットフォームとしてデバイス接続およびメッセージの中継を担当し、MySQLはデータストレージプラットフォームとしてデバイスの状態やメタデータ、メッセージデータの保存およびデータ分析を担います。

EMQXはルールエンジンとSinkを通じてデバイスイベントやデータをMySQLに転送します。アプリケーションはMySQL内のデータを読み取り、デバイスの状態を把握したり、デバイスのオンライン・オフライン記録を取得したり、デバイスデータを分析したりできます。具体的なワークフローは以下の通りです:
- IoTデバイスがEMQXに接続:IoTデバイスがMQTTプロトコルを通じて正常に接続されると、オンラインイベントがトリガーされます。イベントにはデバイスID、送信元IPアドレスなどの属性情報が含まれます。
- メッセージのパブリッシュと受信:デバイスは特定のトピックにテレメトリや状態データをパブリッシュします。EMQXはこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
- ルールエンジンによるメッセージ処理:組み込みのルールエンジンにより、特定のソースからのメッセージやイベントをトピックマッチングに基づいて処理します。ルールエンジンは対応するルールにマッチし、データ形式の変換、特定情報のフィルタリング、メッセージへのコンテキスト情報の付加などを行います。
- MySQLへの書き込み:ルールによりメッセージのMySQL書き込みがトリガーされます。SQLテンプレートを利用して、ルール処理結果からデータを抽出しSQLを構築、MySQLで実行することで、メッセージの特定フィールドをデータベースの対応するテーブルやカラムに書き込んだり更新したりします。
イベントおよびメッセージデータがMySQLに書き込まれた後は、MySQLに接続してデータを読み取り、以下のような柔軟なアプリケーション開発が可能です:
- Grafanaなどの可視化ツールに接続し、データに基づくグラフを生成してデータ変化を表示する。
- デバイス管理システムに接続し、デバイス一覧や状態を閲覧、異常なデバイス動作を検知し、潜在的な問題をタイムリーに解消する。
特長と利点
MySQLとのデータ統合は、以下のような特長とメリットをビジネスにもたらします:
- 柔軟なイベント処理:EMQXルールエンジンを通じて、MySQLはデバイスのライフサイクルイベントを処理でき、IoTアプリケーション実装に必要な各種管理・監視タスクの開発を大幅に容易にします。イベントデータを分析することで、デバイスの故障や異常動作、トレンド変化を迅速に検知し、適切な対策を講じることが可能です。
- メッセージ変換:メッセージはEMQXルールを介して多様な処理や変換を受けてからMySQLに書き込まれるため、保存や利用がより便利になります。
- 柔軟なデータ操作:MySQL Sinkが提供するSQLテンプレートを使うことで、特定フィールドのデータをMySQLデータベースの対応テーブル・カラムに簡単に書き込みや更新ができ、柔軟なデータ保存・管理が実現します。
- ビジネスプロセスの統合:データ統合により、デバイスデータをMySQLの豊富なエコシステムアプリケーションと連携可能にし、ERPやCRM、その他カスタムビジネスシステムとの統合を促進し、高度なビジネスプロセスや自動化を実現します。
- ランタイムメトリクス:各Sinkのランタイムメトリクス(総メッセージ数、成功/失敗数、現在のレートなど)を閲覧可能です。
柔軟なイベント処理、多様なメッセージ変換、柔軟なデータ操作、リアルタイムの監視・分析機能を通じて、効率的で信頼性が高くスケーラブルなIoTアプリケーションを構築し、ビジネスの意思決定や最適化に役立てられます。
はじめる前に
このセクションでは、EMQXダッシュボードでMySQLデータ統合を作成する前に必要な準備、MySQLサーバーのインストールやデータテーブルの作成について説明します。
前提条件
MySQLサーバーのインストール
Dockerを使ってMySQLサーバーをインストールし、Dockerイメージを起動します。
# MySQL Dockerイメージを起動し、パスワードをpublicに設定
docker run --name mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=public -d mysql
# コンテナにアクセス
docker exec -it mysql bash
# コンテナ内でMySQLサーバーに接続し、設定したパスワードを入力
mysql -u root -p
# データベースを作成し、選択
CREATE DATABASE emqx_data CHARACTER SET utf8mb4;
use emqx_data;データテーブルの作成
以下のSQL文を使い、MySQLデータベース内に
emqx_messagesテーブルを作成します。このテーブルは、メッセージのクライアントID、トピック、ペイロード、作成日時を保存します。sqlCREATE TABLE emqx_messages ( id INT AUTO_INCREMENT PRIMARY KEY, clientid VARCHAR(255), topic VARCHAR(255), payload TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
注意:バイナリペイロードが必要な場合は、カラムを"BLOB"型で宣言してください。
以下のSQL文を使い、クライアントID、イベントタイプ、作成日時を保存する
emqx_client_eventsテーブルを作成します。sqlCREATE TABLE emqx_client_events ( id INT AUTO_INCREMENT PRIMARY KEY, clientid VARCHAR(255), event VARCHAR(255), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
コネクターの作成
このセクションでは、SinkをMySQLサーバーに接続するためのコネクターの作成方法を説明します。
以下の手順は、EMQXとMySQLをローカルマシンで実行していることを前提としています。MySQLやEMQXがリモートで稼働している場合は、設定を適宜調整してください。
- EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。
- ページ右上のCreateをクリックします。
- Create ConnectorページでMySQLを選択し、Nextをクリックします。
- Configurationステップで以下の情報を設定します:
- Connector name:コネクター名を入力します。英数字の大文字・小文字の組み合わせで、例:
my_mysql。 - Server Host:
127.0.0.1:3306、またはMySQLサーバーがリモートの場合は実際のホスト名を入力。 - Database Name:
emqx_dataを入力。 - Username:
rootを入力。 - Password:
publicを入力。
- Connector name:コネクター名を入力します。英数字の大文字・小文字の組み合わせで、例:
- 高度な設定(任意):高度な設定を参照してください。
- Createをクリックする前に、Test ConnectivityをクリックしてコネクターがMySQLサーバーに接続可能かテストできます。
- ページ下部のCreateボタンをクリックしてコネクターの作成を完了します。ポップアップダイアログでは、Back to Connector Listをクリックしてコネクター一覧に戻るか、Create RuleをクリックしてMySQLへの転送やクライアントイベント記録を指定するルールを作成できます。詳細はメッセージ保存用MySQL Sink付きルール作成およびイベント記録用MySQL Sink付きルール作成を参照してください。
メッセージ保存用MySQL Sink付きルールの作成
このセクションでは、ダッシュボードでルールを作成し、ソースMQTTトピックt/#からのメッセージを処理し、設定したSinkを通じてMySQLのemqx_messagesテーブルに保存する方法を示します。
このデモは、EMQXとMySQLをローカルマシンで実行していることを前提としています。リモート環境の場合は設定を調整してください。
EMQXダッシュボードで、Integration -> Rulesをクリックします。
ページ右上のCreateをクリックします。
ルールIDに
my_ruleを入力し、SQL Editorに以下の文を設定します。これはトピックt/#配下のMQTTメッセージをMySQLに保存することを意味します。注意:独自のSQL構文を指定する場合は、Sinkが必要とするすべてのフィールドを
SELECT句に含めてください。sqlSELECT * FROM "t/#"TIP
初心者の方は、SQL Examplesをクリックし、Enable Testを有効にしてSQLルールを学習・テストできます。
- Add Actionボタンをクリックして、ルールによりトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理したデータをMySQLに送信します。
Type of Actionドロップダウンから
MySQLを選択します。ActionはデフォルトのCreate Actionのままにします。既に作成済みのSinkがあれば選択可能ですが、このデモでは新規Sinkを作成します。Sink名を入力します。英数字の大文字・小文字の組み合わせで指定してください。
Connectorドロップダウンから先ほど作成した
my_mysqlを選択します。新規コネクターを作成する場合はドロップダウン横のボタンをクリックしてください。設定パラメータはコネクター作成を参照してください。利用する機能に応じてSQL Templateを設定します:
バッチモードが無効の場合、MySQLはプリペアドステートメントを使用します。プレースホルダーを引用符で囲んだり、文末にセミコロンを付けたりしないでください。
重要なお知らせ
EMQX 6.3.1以降、バッチモードが有効な場合、EMQXはSink作成時に制限付きSQLパーサーを使用し、安全にテンプレートをレンダリングし、サポートされない構文を拒否します。
テンプレートは単一のMySQL
INSERT INTO ... VALUES文で、1行のデータを構成する必要があります。定数、文字列リテラル内のプレースホルダー、算術式、関数、条件式、ON DUPLICATE KEY UPDATEをサポートします。ON DUPLICATE KEY UPDATE内の代入式にプレースホルダーは使えません。SQLコメント、追加文、識別子内のプレースホルダーはサポートされません。EMQXは作成するすべての接続でMySQLの
ANSI_QUOTESおよびNO_BACKSLASH_ESCAPESSQLモードを無効にします。アップグレード前に、サポートされない構文を使っているテンプレートやこれらSQLモードに依存するテンプレートを見直してください。sqlINSERT INTO emqx_messages(clientid, topic, payload, created_at) VALUES( ${clientid}, ${topic}, ${payload}, FROM_UNIXTIME(${timestamp}/1000) )SQLテンプレート内でプレースホルダー変数が未定義の場合、SQL template上部のUndefined Vars as Nullスイッチを切り替えてルールエンジンの動作を定義できます:
無効(デフォルト):ルールエンジンは文字列
undefinedをデータベースに挿入します。有効:変数が未定義の場合、ルールエンジンは
NULLを挿入します。TIP
可能な限りこのオプションは有効にしてください。無効にするのは後方互換性確保のためのみです。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した際にトリガーされます。詳細はフォールバックアクションを参照してください。
高度な設定(任意):高度な設定を参照してください。
CreateボタンをクリックしてSink設定を完了します。新しいSinkがAction Outputsに追加されます。
Create Ruleページに戻り、設定内容を確認してCreateボタンをクリックしルールを生成します。
これでルールが正常に作成されました。Integration -> Rulesページで新規ルールを確認できます。**Actions(Sink)**タブをクリックすると新しいMySQL Sinkが表示されます。
また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/#配下のメッセージがMySQLに送信・保存されていることが確認できます。
イベント記録用MySQL Sink付きルールの作成
このセクションでは、クライアントのオンライン/オフライン状態を記録し、イベントデータをMySQLのemqx_client_eventsテーブルに保存するルールの作成方法を示します。
ルール作成手順はメッセージ保存用MySQL Sink付きルール作成とほぼ同様ですが、SQLルール構文とSQLテンプレートが異なります。
バッチモードが有効な場合、前述のSQLテンプレート制限がこのテンプレートにも適用されます。
オンライン/オフライン状態記録用ルールを作成するには、SQL Editorに以下の文を入力します:
SELECT
*
FROM
"$events/client_connected", "$events/client_disconnected"クライアントイベントデータをテーブルに挿入するには、以下のSQLテンプレートを使用します:
INSERT INTO emqx_client_events(clientid, event, created_at) VALUES (
${clientid},
${event},
FROM_UNIXTIME(${timestamp}/1000)
)ルールのテスト
MQTTXを使ってトピックt/1にメッセージを送信し、オンライン/オフラインイベントをトリガーします。
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello MySQL" }'2つのSinkの稼働状況を確認すると、新規の着信メッセージと送信メッセージがそれぞれ1件ずつ、イベントレコードが2件あるはずです。
emqx_messagesテーブルにデータが書き込まれているか確認します。
mysql> select * from emqx_messages;
+----+----------+-------+--------------------------+---------------------+
| id | clientid | topic | payload | created_at |
+----+----------+-------+--------------------------+---------------------+
| 1 | emqx_c | t/1 | { "msg": "hello MySQL" } | 2022-12-09 08:44:07 |
+----+----------+-------+--------------------------+---------------------+
1 row in set (0.01 sec)emqx_client_eventsテーブルにデータが書き込まれているか確認します。
mysql> select * from emqx_client_events;
+----+----------+---------------------+---------------------+
| id | clientid | event | created_at |
+----+----------+---------------------+---------------------+
| 1 | emqx_c | client.connected | 2022-12-09 08:44:07 |
| 2 | emqx_c | client.disconnected | 2022-12-09 08:44:07 |
+----+----------+---------------------+---------------------+
2 rows in set (0.00 sec)高度な設定
このセクションでは、MySQLコネクターおよびSinkの高度な設定オプションについて詳述します。ダッシュボードでコネクターやSinkを設定する際、Advanced Settingsに進み、以下のパラメータをニーズに合わせて調整してください。
| 項目 | 説明 | 推奨値 |
|---|---|---|
| Connection Pool Size | MySQLサービスとの接続プール内で同時に維持可能な接続数を指定します。このオプションはEMQXとMySQL間のアクティブ接続数を制限または増加させることで、アプリケーションのスケーラビリティやパフォーマンス管理に役立ちます。 注意:適切な接続プールサイズはシステムリソース、ネットワークレイテンシ、アプリケーションのワークロードなど複数要因に依存します。大きすぎるとリソース枯渇、小さすぎるとスループット制限となる可能性があります。 | 8 |
| Start Timeout | コネクターが自動起動したリソースが正常状態になるまで待機する最大時間(秒)を指定します。この設定により、MySQLなどの接続先リソースが完全に稼働し、データ処理準備が整うまでコネクターが処理を進めないようにします。 | 5 秒 |
| Buffer Pool Size | EMQXとMySQL間の出口方向(egress)Sinkでデータフロー管理に割り当てるバッファワーカー数を指定します。これらのワーカーはデータ送信前に一時的にデータを保持・処理します。入口方向(ingress)のみを扱うSinkではこの値を"0"に設定可能です。 | 16 |
| Request TTL | バッファに入ったリクエストが有効とみなされる最大期間(秒)を指定します。リクエストがこの期間を超えてバッファに残るか、MySQLからの応答やアックを受け取れなかった場合、リクエストは期限切れとみなされます。 | 45 秒 |
| Health Check Interval | コネクターがMySQLとの接続状態を自動的にヘルスチェックする間隔(秒)を指定します。 | 15 秒 |
| Max Buffer Queue Size | コネクター内の各バッファワーカーがバッファリング可能な最大バイト数を指定します。バッファワーカーはMySQL送信前にデータを一時保持し、データフローを効率的に処理します。システム性能やデータ転送要件に応じて調整してください。 | 256 MB |
| Max Batch Size | EMQXからMySQLへ単一転送操作で送信されるデータバッチの最大サイズを指定します。サイズ調整によりデータ転送の効率やパフォーマンスを最適化可能です。 "1"に設定すると、データレコードはバッチ化されず個別に送信されます。 | 1 |
| Query Mode | メッセージ送信要件に応じてasynchronousまたはsynchronousクエリモードを選択できます。非同期モードではMySQLへの書き込みがMQTTメッセージのパブリッシュ処理をブロックしませんが、クライアントがMySQL到着前にメッセージを受信する可能性があります。 | Async |
| Inflight Window | "in-flight query"とは開始されたがまだ応答やアックを受け取っていないクエリを指します。コネクターがMySQLと通信する際に同時に存在可能なin-flightクエリの最大数を制御します。 Query Modeが asyncの場合、このパラメータは特に重要です。同一MQTTクライアントからのメッセージを厳密に順序処理したい場合は、この値を1に設定してください。 | 100 |
さらに詳しく
以下のリンクから詳細を確認できます: