Skip to content

TDengineへのMQTTデータ取り込み ​

TDengineは、IoTおよび産業用IoT(IIoT)シナリオ向けに設計・最適化されたビッグデータプラットフォームです。中核には高性能な時系列データベースがあり、クラスター指向のアーキテクチャ、クラウドネイティブ設計、ミニマリスティックなアプローチが特徴です。EMQXはTDengineとの統合をサポートしており、多数のデバイスやデータコレクターからの大量データの送信、保存、分析、配信を可能にします。これにより、ビジネス運用状態のリアルタイム監視や早期警告を提供し、リアルタイムのビジネスインサイトを実現します。

本ページでは、EMQXとTDengine間のデータ統合について包括的に紹介し、データ統合の作成および検証方法を実践的に説明します。

動作概要 ​

TDengineデータ統合はEMQXの組み込み機能です。組み込みのルールエンジンコンポーネントにより、EMQXからTDengineへのデータ取り込みが簡素化され、複雑なコーディングが不要になります。EMQXはルールエンジンとSinkを通じてデバイスデータをTDengineに転送します。TDengineデータ統合により、MQTTメッセージやクライアントイベントをTDengineに保存できます。さらに、TDengine内のデータ更新や削除はイベントによってトリガー可能であり、デバイスのオンライン状態や過去のオンライン/オフラインイベントの記録が可能です。

以下の図は、産業用IoTにおけるEMQXとTDengineのデータ統合の典型的なアーキテクチャを示しています。

EMQX Integration TDengine

産業用エネルギー消費管理シナリオを例に、ワークフローは以下の通りです。

  1. メッセージのパブリッシュと受信:産業用デバイスはMQTTプロトコルを通じてEMQXに正常に接続し、定期的にエネルギー消費データをパブリッシュします。このデータには生産ライン識別子やエネルギー消費値が含まれます。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
  2. ルールエンジンによるメッセージ処理:組み込みのルールエンジンは、トピックマッチングに基づいて特定のソースからのメッセージを処理します。メッセージが到着するとルールエンジンを通過し、対応するルールとマッチングしてメッセージデータを処理します。これにはデータ形式の変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などが含まれます。
  3. TDengineへのデータ取り込み:ルールエンジンで定義されたルールが、メッセージをTDengineに書き込む操作をトリガーします。TDengine SinkはSQLテンプレートを提供し、特定のメッセージフィールドをTDengineの対応するテーブルやカラムに柔軟に書き込むデータ形式を定義可能です。

エネルギー消費データがTDengineに書き込まれた後、標準SQLと強力な時系列拡張機能を用いてリアルタイムにデータ分析が可能となり、多数のサードパーティのバッチ分析、リアルタイム分析、レポートツール、AI/MLツール、可視化ツールとシームレスに統合できます。例えば:

  • Grafanaなどの可視化ツールに接続し、エネルギー消費データのチャートを生成・表示。
  • ERPやPower BIなどのアプリケーションシステムに接続し、生産分析や生産計画の調整を実施。
  • ビジネスシステムに接続し、リアルタイムのエネルギー使用分析を行い、データ駆動型のエネルギー管理を支援。

特長と利点 ​

TDengineデータ統合は、以下の特長と利点をビジネスにもたらします。

  • 効率的なデータ処理:EMQXは多数のIoTデバイス接続とメッセージスループットを効率的に処理可能です。TDengineはデータの書き込み、保存、クエリに優れており、IoTシナリオのデータ処理ニーズをシステムに過負荷をかけずに満たします。
  • メッセージ変換:メッセージはEMQXのルール内で豊富な処理や変換を経てからTDengineに書き込まれます。
  • クラスターとスケーラビリティ:EMQXとTDengineはクラスター機能をサポートし、クラウドネイティブアーキテクチャ上に構築されているため、クラウドプラットフォームの弾力的なストレージ、計算、ネットワークリソースを最大限に活用し、ビジネスの成長に応じて柔軟な水平スケールが可能です。
  • 高度なクエリ機能:TDengineはタイムスタンプデータの効率的なクエリと分析のために最適化された関数、演算子、インデックス技術を提供し、IoT時系列データから正確なインサイトを抽出可能です。

はじめる前に ​

このセクションでは、TDengineデータ統合の作成を開始する前に必要な準備、TDengineサーバーのセットアップやデータテーブルの作成方法について説明します。

前提条件 ​

TDengineの起動とデータベース作成 ​

TDengineを起動またはTDengineサービスに接続し、データベースを作成するには以下の2つの方法があります。

TDengineでのデータテーブル作成 ​

メッセージ保存とステータス記録用に、TDengineデータベース内に2つのデータテーブルを作成する必要があります。

  1. 以下のSQL文を使用して、t_mqtt_msgテーブルを作成します。このテーブルは各メッセージのクライアントID、トピック、ペイロード、作成時間を保存します。
sql
   CREATE TABLE t_mqtt_msg (
       ts timestamp,
       msgid NCHAR(64),
       mqtt_topic NCHAR(255),
       qos TINYINT,
       payload BINARY(1024),
       arrived timestamp
     );
  1. 以下のSQL文を使用して、emqx_client_eventsテーブルを作成します。このテーブルは各イベントのクライアントID、イベントタイプ、作成時間を保存します。
sql
     CREATE TABLE emqx_client_events (
         ts timestamp,
         clientid VARCHAR(255),
         event VARCHAR(255)
       );

コネクターの作成 ​

このセクションでは、SinkをTDengineサーバーに接続するためのコネクター作成方法を説明します。

  1. EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。

  2. 画面右上のCreateをクリックします。

  3. Create ConnectorページでTDengineを選択し、Nextをクリックします。

  4. Configurationステップで、接続先に応じて以下の情報を設定します。

  5. 高度な設定(任意):詳細はSinkの機能を参照してください。

  6. Createをクリックする前に、Test ConnectivityをクリックしてコネクターがTDengineサーバーに接続できるかテストできます。

  7. ページ下部のCreateボタンをクリックしてコネクター作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create Ruleをクリックしてルール作成に進み、Sinkを使ってTDengineに転送するデータやクライアントイベントの記録を指定します。詳細はメッセージ保存用のTDengine Sinkルール作成およびイベント記録用のTDengine Sinkルール作成を参照してください。

メッセージ保存用のTDengine Sinkルール作成 ​

このセクションでは、ダッシュボードでMQTTトピックt/#からのメッセージを処理し、処理済みデータを設定済みのSinkを通じてTDengineのt_mqtt_msgテーブルに保存するルール作成方法を説明します。

  1. EMQXダッシュボードでIntegration -> Rulesをクリックします。

  2. 画面右上のCreateをクリックします。

  3. ルールIDにmy_ruleを入力し、SQL Editorでメッセージ保存用のルールを作成します。例えば、以下のステートメントはトピックt/#配下のMQTTメッセージをTDengineに保存します。

    注意:独自のSQL構文を指定する場合は、Sinkで必要な全てのフィールドをSELECT句に含めていることを確認してください。

    sql
      SELECT
        *,
        now_timestamp('millisecond')  as ts
      FROM
        "t/#"

    TIP

    初心者の方はSQL ExamplesやEnable TestをクリックしてSQLルールを学習・テストしてください。

    • Add Actionボタンをクリックし、ルール発動時にトリガーされるアクションを定義します。このアクションによりEMQXはルールで処理したデータをTDengineに送信します。
  4. Type of ActionドロップダウンからTDengineを選択します。ActionはデフォルトのCreate Actionのままにします。既に作成済みのTDengine Sinkがあれば選択可能です。本デモでは新規Sinkを作成します。

  5. Sinkの名前を入力します。英数字の組み合わせで指定してください。

  6. Connectorドロップダウンから先ほど作成したmy_tdengineを選択します。隣のボタンから新規コネクター作成も可能です。設定パラメータはコネクター作成を参照してください。

  7. SinkのSQL Templateを設定します。以下のSQLを使ってデータ挿入を完了できます。CSVファイルによるバッチ設定もサポートしています。詳細はバッチ設定を参照してください。

    重要なお知らせ

    EMQX 6.3.1以降、Sink作成時にSQLテンプレートを解析し、SQLコンテキストに基づいてプレースホルダー値をエスケープし、サポートされない構文を拒否します。テンプレートは単一のTDengine INSERT文でなければなりません。複数テーブルのVALUES挿入やUSING ... TAGS句をサポートします。ターゲットテーブル名は生の名前またはバッククォート付きでプレースホルダーを含めることができ、例:test_${clientid}。EMQXは完全なターゲットを1つの引用識別子としてレンダリングします。SQLコメント、FILE入力、追加文はサポートされません。

    この検証により、以前のバージョンで受け入れられたテンプレートが拒否される場合があります。アップグレード前に互換性のないテンプレートを修正してください。

    TIP

    EMQX 6.3.0では文字列プレースホルダー値を手動で引用符で囲む必要がありましたが、6.3.1以降はプレースホルダーを完全な値または文字列リテラル内で使用でき、EMQXがSQLコンテキストに応じて値をエスケープします。

    sql
    INSERT INTO t_mqtt_msg(ts, msgid, mqtt_topic, qos, payload, arrived) 
        VALUES (${ts}, '${id}', '${topic}', ${qos}, '${payload}', ${timestamp})

    SQLテンプレート内でプレースホルダー変数が未定義の場合、SQL template上のUndefined Vars as Nullスイッチでルールエンジンの動作を定義可能です。

    • Disabled(デフォルト):ルールエンジンは未定義変数として文字列undefinedをデータベースに挿入可能。

    • Enabled:未定義変数の場合、ルールエンジンはNULLを挿入。

      TIP

      可能な限りこのオプションは有効にしてください。無効化は後方互換性確保のためのみ推奨されます。

  8. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義可能です。詳細はフォールバックアクションを参照してください。

  9. 高度な設定(任意):必要に応じてsyncまたはasyncクエリモードを選択可能です。詳細はSinkの機能を参照してください。

  10. Createをクリックする前に、Test ConnectivityでSinkがTDengineに接続できるかテスト可能です。

  11. CreateボタンをクリックしてSink設定を完了します。新しいSinkがAction Outputsに追加されます。

  12. Create Ruleページに戻り、設定内容を確認後、Createをクリックしてルールを作成します。

これでTDengine Sink用ルールの作成が完了しました。Integration -> Rulesページで新規ルールを確認できます。**Actions(Sink)**タブで新しいTDengine Sinkを確認可能です。

また、Integration -> Flow Designerでトポロジーを確認すると、トピックt/#配下のメッセージがルールmy_ruleで解析され、TDengineに送信・保存されていることが分かります。

バッチ設定 ​

TDengineでは1つのデータエントリに数百のデータポイントを含むことがあり、SQL文の作成が困難です。これに対応するため、EMQXはSQLのバッチ設定機能を提供しています。

SQLテンプレート編集時に、バッチ設定機能を使ってCSVファイルから挿入用フィールドをインポートできます。

  1. SQL Template下のBatch Settingボタンをクリックし、Import Batch Settingポップアップを開きます。

  2. 指示に従いバッチ設定テンプレートファイルをダウンロードし、テンプレート内のフィールドのキー・値ペアを入力します。デフォルトのテンプレート内容は以下の通りです。

    FieldValueChar Value備考(任意)
    tsnowFALSE例
    msgid${id}TRUE
    mqtt_topic${topic}TRUE
    qos${qos}FALSE
    temp${payload.temp}FALSE
    hum${payload.hum}FALSE
    status${payload.status}FALSE
    • Field:フィールドキー。定数または${var}形式のプレースホルダーをサポート。
    • Value:フィールド値。定数または${var}形式のプレースホルダーをサポート。SQLでは文字列は引用符で囲む必要がありますが、テンプレートファイル内では不要。文字列かどうかはChar Value列で指定。
    • Char Value:フィールドが文字列型かどうかを指定。SQL生成時に引用符を付加。文字列型ならTRUEまたは1、そうでなければFALSEまたは0を入力。
    • 備考:CSVファイル内の注釈用で、EMQXへのインポート対象外。

    CSVファイルのバッチ設定データは2048行を超えないようにしてください。

  3. 入力済みテンプレートファイルを保存し、Import Batch SettingポップアップにアップロードしてImportをクリックしバッチ設定を完了します。

  4. インポート後、SQL Template内でテーブル名設定やSQLコードの整形などをさらに調整可能です。

イベント記録用のTDengine Sinkルール作成 ​

このセクションでは、クライアントのオンライン/オフライン状態を記録し、イベントデータを設定済みSinkを通じてTDengineのemqx_client_eventsテーブルに保存するルール作成方法を説明します。

ルール作成手順はメッセージ保存用のTDengine Sinkルール作成とほぼ同様で、SQLルール構文とSQLテンプレートのみ異なります。

オンライン/オフライン状態記録用のSQLルール構文は以下の通りです。

sql
SELECT
      *,
      now_timestamp('millisecond')  as ts
    FROM 
      "$events/client_connected", "$events/client_disconnected"

SinkのSQLテンプレートは以下の通りです。

上記のSQLテンプレート制限はこのテンプレートにも適用されます。

以下のテンプレートはシングルクォート内の文字列プレースホルダーを使用しています。SQL文の末尾にセミコロン(;)を付けないでください。

sql
INSERT INTO emqx_client_events(ts, clientid, event) VALUES (
      ${ts},
      '${clientid}',
      '${event}'
    )

ルールのテスト ​

MQTTXを使ってトピックt/1にメッセージを送信し、オンライン/オフラインイベントをトリガーします。

bash
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello TDengine" }'

2つのSinkの稼働状況を確認すると、1件の新規受信メッセージと1件の新規送信メッセージ、2件のイベントレコードがあるはずです。

t_mqtt_msgデータテーブルにデータが書き込まれているか確認します。

bash
taos> select * from t_mqtt_msg;
           ts            |             msgid              |           mqtt_topic           | qos  |            payload             |         arrived         |
==============================================================================================================================================================
 2023-02-13 06:10:53.787 | 0005F48EB5A83865F440000014F... | t/1                            |    0 | { "msg": "hello TDengine" }    | 2023-02-13 06:10:53.787 |
Query OK, 1 row(s) in set (0.002968s)

emqx_client_eventsテーブル:

bash
taos> select * from emqx_client_events;
           ts            |            clientid            |             event              |
============================================================================================
 2023-02-13 06:10:53.777 | emqx_c                         | client.connected               |
 2023-02-13 06:10:53.791 | emqx_c                         | client.disconnected            |
Query OK, 2 row(s) in set (0.002327s)