Skip to content

他の MQTT サービスとのブリッジ

MQTT ブローカーのデータ統合は、EMQX に別の EMQX クラスターや他の MQTT サービスへのメッセージブリッジ接続機能を提供し、ネットワークやサービスを跨いだデータの相互作用と通信を可能にします。本ページでは、EMQX における MQTT メッセージブリッジの動作原理を紹介し、メッセージブリッジの作成および検証方法について実践的なガイドを提供します。

動作原理

ブリッジング中、EMQX はクライアントとしてターゲットサービスと MQTT 接続を確立し、パブリッシュ・サブスクライブモデルを通じて双方向のメッセージ送受信を実現します。

  • 送信メッセージ(Sink):ローカルのトピックからメッセージをパブリッシュし、リモート MQTT サービスの指定トピックへ送信します。
  • 受信メッセージ(Source):リモート MQTT サービスのトピックをサブスクライブし、そのメッセージをローカルの EMQX に転送します。

EMQX は同一接続で複数のブリッジルールを設定可能で、それぞれ異なるトピックマッピングやメッセージ変換ルールを持ち、メッセージルーティングに類似した機能を実装します。ブリッジング中はルールエンジンを通じてメッセージのフィルタリング、強化、変換処理も行えます。

下図は EMQX と他の MQTT サービス間のデータ統合の典型的なアーキテクチャを示しています。

EMQX Integration MQTT

特長と利点

MQTT ブローカーデータ統合は以下の特長と利点を備えています。

  • 広範な互換性:標準 MQTT プロトコルを使用しており、AWS IoT Core、Azure IoT Hubs などの各種 IoT プラットフォームや、オープンソースや他の業界 MQTT ブローカー、IoT プラットフォームとのブリッジを可能にします。これにより、多様なデバイスやプラットフォームとのシームレスな統合と通信が実現します。
  • 双方向データフロー:双方向のデータフローをサポートし、EMQX からリモート MQTT サービスへのメッセージパブリッシュと、リモート MQTT サービスからのメッセージサブスクライブおよびローカルパブリッシュの両方を実現します。この双方向通信により、異なるシステム間のデータ転送がより柔軟かつ制御可能になります。
  • 柔軟なトピックマッピング:MQTT のパブリッシュ・サブスクライブモデルに基づき、柔軟なトピックマッピングを実装しています。トピックへのプレフィックス追加や、クライアントのコンテキスト情報(クライアントID、ユーザー名など)を用いた動的トピック構築をサポートし、ニーズに応じたカスタマイズ処理やルーティングが可能です。
  • 高性能:接続プールや共有サブスクリプションなどのパフォーマンス最適化オプションを提供し、個々のブリッジクライアントの負荷を軽減、ブリッジレイテンシの低減とメッセージスループットの向上を実現します。これらの最適化により、システム全体のパフォーマンスとスケーラビリティが向上します。
  • ペイロード変換:SQL ルールを定義してメッセージペイロードの処理が可能です。メッセージ転送時にペイロードの抽出、フィルタリング、強化、変換などの操作を行えます。例えば、ペイロードからリアルタイムメトリクスを抽出し、変換・処理してからリモート MQTT ブローカーに届けることができます。
  • メトリクス監視:各 Sink/Source ごとにランタイムメトリクス監視を提供し、総メッセージ数、成功/失敗数、現在のレートなどをリアルタイムで確認でき、Sink/Source のパフォーマンスと健全性を監視・評価できます。

MQTT 接続情報の準備

前提条件

以下を把握していることを確認してください。

MQTT ブローカーデータ統合を作成する前に、リモート MQTT サービスの接続情報を取得する必要があります。ここでは EMQX の オンライン MQTT サーバー を例に説明します。

  • MQTT サービスアドレス:ターゲット MQTT サービスのアドレスとポート。例として broker.emqx.io:1883 を使用します。サポートされる形式は host:port[IPv6]:portmqtt://host:portmqtt://[IPv6]:portmqtts://host:portmqtts://[IPv6]:port です。ポート省略時はデフォルトの MQTT ポート 1883 が使用されます。その他の URI スキームは非対応です。mqttmqtts スキームはアドレス解析用であり、TLS 対応 MQTT リスナーに接続する場合はコネクター設定で TLS を別途構成してください。
  • ユーザー名:接続に必要なユーザー名。ターゲットサービスが認証不要の場合は空欄可。
  • パスワード:接続に必要なパスワード。認証不要の場合は空欄可。
  • プロトコルタイプ:ターゲットサービスが TLS を有効にしているか、TCP/TLS 上の MQTT プロトコルを使用しているかを確認します。EMQX の MQTT ブリッジは現時点で MQTT over WebSocket や MQTT over QUIC などのプロトコルはサポートしていません。
  • プロトコルバージョン:ターゲット MQTT サービスで使用されているプロトコルバージョン。EMQX は MQTT 3.1、3.1.1、および MQTT 5.0 をサポートしています。

データ統合は EMQX や他の標準 MQTT サーバーとの高い互換性とサポートを提供します。その他の種類の MQTT サービスに接続する場合は、該当サービスのドキュメントを参照して接続情報を取得してください。一般的に多くの IoT プラットフォームは標準 MQTT アクセス方法を提供しており、それに基づきデバイス情報を上記の MQTT 接続情報に変換できます。

注意

EMQX がクラスター運用中、または接続プールが有効な場合、同一クライアントIDを用いて複数ノードが同一 MQTT サービスに接続すると通常デバイス競合が発生します。そのため、MQTT メッセージブリッジは現在固定クライアントIDの設定をサポートしていません。

コネクターの作成

ここではリモート MQTT サーバーとの接続を設定する手順を案内します。

  1. ダッシュボードの Integration -> Connector ページに移動します。

  2. ページ右上の Create をクリックします。

  3. コネクタータイプ一覧から MQTT Broker を選択し、Next をクリックします。

  4. コネクター名を入力します。英数字の組み合わせで、例として my_mqtt_bridge とします。

  5. 接続情報を設定します:

    • MQTT Broker:TCP/TLS 上の MQTT のみサポート。MQTT サービスアドレスを入力します。例:broker.emqx.io:1883[::1]:1883mqtt://broker.emqx.io:1883。TLS 対応 MQTT リスナーの場合はコネクター設定で TLS を別途構成してください。

    • ClientID Prefix:空欄可。実際の運用ではクライアント ID プレフィックスを指定するとクライアント管理が容易になります。EMQX はクライアント ID プレフィックスと接続プールサイズに基づき自動でクライアント ID を生成します。詳細は 接続プールとクライアント ID 生成ルール を参照してください。

    • UsernamePassword:MQTT ブローカーが認証を要求する場合はクライアント ID に対応するユーザー名とパスワードを入力します。認証不要の場合(パブリックブローカーなど)は空欄可。

    • Keepalive:希望のキープアライブ間隔を指定します。

    • MQTT Version:ブローカー接続に適したバージョンを選択します。

    • Static ClientId Entries:Azure IoT Hubs など安定した接続が必要なサービスに接続する際に、コネクターに静的クライアント ID を設定できます。詳細は 静的クライアント ID の設定 を参照してください。

      TIP

      静的クライアント ID が定義されている場合、明示的に静的クライアント ID が割り当てられた EMQX ノードのみが MQTT 接続を開始します。

その他の設定はデフォルトのままにして、Create ボタンをクリックしてコネクター作成を完了します。作成したコネクターは Sink と Source の両方で利用可能です。次に、このコネクターを基にデータブリッジルールを作成できます。

接続プールとクライアント ID 生成ルール

EMQX は複数のクライアントが同時にブリッジ先 MQTT サービスに接続できるようにします。コネクター作成時に MQTT クライアント接続プールを設定し、そのサイズを指定してプール内のクライアント接続数を決定します。接続プールはサーバーリソースを最大限活用し、メッセージスループットと同時接続性能を向上させ、高負荷・高同時接続シナリオに対応します。

MQTT プロトコルでは MQTT サーバーに接続するクライアントは一意のクライアント ID を持つ必要があり、EMQX はクラスター展開可能なため、MQTT ブリッジの各クライアントに一意のクライアント ID を割り当てます。EMQX は以下のパターンでクライアント ID を自動生成します。

bash
[Client ID Prefix]:{Connector Name}{8桁のランダム文字列}:{プール内の接続シーケンス番号}

例えば、クライアント ID プレフィックスが myprefix、コネクター名が foo の場合、実際のクライアント ID は以下のようになります。

bash
myprefix:foo2bd61c44:1

バージョン 5.4.1 以降、EMQX は MQTT クライアント ID の長さを 23 バイトに制限しています。これを超える場合、ハッシュ値に置き換えられます。プレフィックスやコネクター名が長すぎるとユーザー体験が悪化する可能性があります。

この問題に対応し、バージョン 5.7.1 以降は以下のルールを実装しています。

  • プレフィックスなし:動作は変わらず、長さ 23 バイトを超えるクライアント ID はハッシュ化されます。
  • プレフィックスあり
    • プレフィックスが最大 19 バイト:プレフィックスは保持され、残りのクライアント ID 部分は 4 バイトのハッシュに変換され、全体の長さは 23 バイト以内に収まります。
    • プレフィックスが 20 バイト以上:EMQX は設定されたプレフィックスを使用し、クライアント ID の短縮は行いません。

静的クライアント ID の設定

統合で使用するクライアント ID が限られている場合、コネクター設定時に個別ノードに静的クライアント ID セットを割り当てることが可能です。静的クライアント ID を設定するには、EMQX クラスター内の各ノードに対してクライアント ID のリストを提供します。各クライアント ID に対して対応するユーザー名とパスワードを指定できます。これは Azure IoT Hubs のように、各デバイス(クライアント ID)ごとに固有の認証情報が必要なシナリオで特に有用です。

静的クライアント ID を設定する手順は以下の通りです。

  1. Static ClientId Entries セクションで Add ボタンをクリックし、新しい静的クライアント ID エントリーを追加します。必要に応じて複数ノード分のエントリーを追加可能です。

  2. 各エントリーに以下の項目を入力します。

    • Node Name:クライアント ID を割り当てるノード名。例:emqx@10.0.0.1
    • Client ID:静的クライアント ID。例:device1。必要に応じてノードごとに複数のクライアント ID を追加可能です。
      • Username:(任意)認証用のユーザー名。
      • Password:(任意)認証用のパスワード。プラットフォームによりデバイス固有のキー、シークレット、証明書などが該当します(例:Azure IoT Hubs の認証キー)。

    設定例

    ノード名クライアント IDユーザー名(任意)パスワード(任意)
    emqx@10.0.0.1clientid1username1secret1
    clientid3
    emqx@10.0.0.2clientid2username2
    emqx@10.0.0.3clientid4
    clientid5

設定ファイルで各ノードごとに static_clientids パラメータを定義することも可能です。

静的クライアント ID を設定した場合、これらのクライアント ID を使用した MQTT 接続のみが開始され、pool_sizeclientid_prefix など動的クライアント ID に関する設定は無効になります。

MQTT Broker Sink を使ったルールの作成

ここでは、リモート MQTT サービスに転送するデータを指定するルールの作成方法を示します。

  1. ダッシュボードの Integration -> Rules ページに移動します。

  2. ページ右上の Create をクリックします。

  3. ルール ID に my_rule と入力します。

  4. SQL Editor に、t/# トピックからの MQTT メッセージをリモート MQTT サーバーに保存するルールを入力します。SQL は以下の通りです。

    sql
    SELECT
      *
    FROM
      "t/#"

    MQTTv5 プロトコルを使用する場合、pub_props フィールドをルール SQL に含めると、パブリッシュプロパティがそのままリモートブローカーに転送されます。上記の SELECT * FROM t/# はこの例に該当します。

    ユーザープロパティを動的に追加したい場合は、ルール出力の pub_props フィールドに含めることが可能です。例えば以下のルールは、ペイロードからキーと値を取得してユーザープロパティを追加します。

    sql
    SELECT
      *,
      map_put(concat('User-Property.', payload.extra_key), payload.extra_value, pub_props) as pub_props
    FROM
      't/#'
  5. Action Type ドロップダウンから MQTT Broker を選択し、Action はデフォルトの Create Action のままにします。この操作で新しい Sink が作成され、ルールに追加されます。

  6. フォームに Sink の名前と説明を入力します。

  7. Connector ドロップダウンから先ほど作成した my_mqtt_bridge コネクターを選択します。新しいコネクターを作成する場合は、ドロップダウン横の作成ボタンをクリックし、コネクターの作成 の設定パラメータを利用してください。

  8. Sink の情報を設定し、EMQX から外部 MQTT サービスへメッセージをパブリッシュします。

    • Topic:外部 MQTT サービスにパブリッシュするトピック。${var} プレースホルダーをサポートします。ここでは pub/${topic} と入力し、元のトピックに pub/ プレフィックスを付けて転送します。例えば元のメッセージトピックが t/1 の場合、外部 MQTT サービスに転送されるトピックは pub/t/1 となります。
    • QoS:メッセージパブリッシュの QoS。ドロップダウンから 012${qos} のいずれかを選択可能で、他のフィールドから QoS を設定するプレースホルダーも使用できます。ここでは元メッセージの QoS に従うため ${qos} を選択します。
    • Retain:メッセージをリテインとしてパブリッシュするかどうかを truefalse${flags.retain} から選択、または他のフィールドから設定するプレースホルダーを入力できます。ここでは元メッセージのリテインフラグに従うため ${flags.retain} を選択します。
    • Payload:転送メッセージのペイロード生成に使うテンプレート。デフォルトは空欄でルール出力結果を転送します。ここでは ${payload} と入力し、ペイロードのみを転送します。
  9. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。詳細は フォールバックアクション を参照してください。

  10. その他の設定はデフォルトのままにして、Create ボタンをクリックし Sink の作成を完了します。作成後はルール作成ページに戻り、新しい Sink がルールのアクション出力に追加されます。

  11. ルール作成ページの下部にある Create ボタンをクリックし、ルール作成を完了します。

これでルールの作成が完了しました。Integration -> Rules ページで新規ルールを確認できます。Actions(Sink) タブをクリックすると新しい MQTT Broker Sink が表示されます。

また、Integration -> Flow Designer をクリックするとトポロジーを確認できます。トポロジーは、t/# トピックのメッセージがルール my_rule によって処理され、リモート MQTT ブローカーに送信される流れを視覚的に表現しています。

MQTT Broker Sink を使ったルールのテスト

MQTTX CLI を使い、EMQX の t/# トピックから外部 MQTT サービスの pub/${topic} トピックへメッセージをブリッジするルールをテストできます。EMQX の t/1 トピックにメッセージをパブリッシュすると、外部 MQTT サービスの pub/t/1 トピックに転送されるはずです。

  1. 外部 MQTT サービスで pub/# トピックをサブスクライブします。

    bash
    mqttx sub -t pub/# -q 1 -h broker.emqx.io -v
  2. MQTTX で t/1 トピックにメッセージをパブリッシュします。

    bash
    mqttx pub -t t/1 -m "hello world" -r
  3. MQTTX で pub/t/1 トピックをサブスクライブし、メッセージを受信できれば、EMQX から外部 MQTT サービスへのメッセージ転送が成功したことを示します。

    bash
    [2024-1-31] [16:43:13] › topic: pub/t/1
    payload: hello world
  4. ステップ1を繰り返すと、MQTTX で pub/t/1 トピックのリテインメッセージを受信できます。

    bash
    [2024-1-31] [16:44:29] › topic: pub/t/1
    payload: hello world
    retain: true

MQTT Broker Source を使ったルールの作成

ここでは、リモート MQTT サービスからローカル EMQX へデータを転送するルールの作成方法を示します。MQTT Source とメッセージ再パブリッシュアクションを作成し、リモート MQTT サービスから EMQX へのサブスクライブと転送を実現します。

MQTT Broker Source の作成とルールへの追加

  1. ダッシュボードの Integration -> Rules ページに移動します。

  2. ページ右上の Create をクリックします。

  3. ルール ID に my_rule_source と入力します。

  4. ルールのトリガーソース(データ入力)を設定します。ページ右側の Data Inputs タブで、デフォルトの Message タイプの入力を削除し、Add Input をクリックして MQTT Source を作成します。

  5. Add Input ポップアップで、Input Type ドロップダウンから MQTT Broker を選択し、Source ドロップダウンはデフォルトの Create Source のままにします。この操作で新しい Source が作成され、ルールに追加されます。

  6. フォームに Source の名前と説明を入力します。

  7. 先ほど作成した my_mqtt_bridge コネクターをドロップダウンから選択します。新しいコネクターを作成する場合は、ドロップダウン横の作成ボタンをクリックし、コネクターの作成 の設定パラメータを利用してください。

  8. Source の情報を設定し、外部 MQTT サービスから EMQX へのサブスクライブを完了します。

    • Topic:サブスクライブするトピック。+# ワイルドカードをサポートします。

      TIP

      EMQX がクラスター運用中、またはコネクターが接続プール設定の場合、重複メッセージを避けるため共有サブスクリプションを使用する必要があります。

      ここでは $share/1/f/# と入力し、f/# トピックにマッチするすべてのメッセージをサブスクライブします。

    • QoS:サブスクライブの QoS。ドロップダウンから 0 または 1 を選択します。

    • No Local:同一コネクターでパブリッシュしたメッセージがリモート MQTT ブローカーから再度 Source に転送されるのを防止したい場合に有効化します。デフォルトは無効で、コネクターが MQTT 5.0 を使用している場合にのみ有効です。

    • Retain As Published:上流 MQTT ブローカーから転送されるメッセージの元のリテインフラグを保持します。有効化しない場合、上流 MQTT ブローカーはリテインフラグをクリアします。デフォルトは有効で、コネクターが MQTT 5.0 を使用している場合にのみ有効です。

  9. その他の設定はデフォルトのままにして、Create ボタンをクリックし Source 作成を完了します。Source はルールのデータ入力に追加されます。ルール SQL は以下のように変更されます。

    sql
    SELECT
      *
    FROM
      "$bridges/mqtt:my_source"

    ルール SQL は MQTT Source から以下のフィールドを抽出でき、データ処理用に SQL を調整可能です。デフォルトの SQL は本例に十分です。

    フィールド名説明
    topic発信元メッセージのトピック
    server接続された Source のサーバーアドレス
    retain上流 MQTT ブローカーから受信したメッセージの retain フラグ。Retain As Published 無効時は false
    qosメッセージの QoS(サービス品質)
    pub_propsMQTT 5.0 メッセージプロパティオブジェクト(ユーザープロパティペアなどを含む)
    pub_props.User-Property-Pairsキー・バリューのペアを含むユーザープロパティ配列例:{"key":"foo", "value":"bar"}
    pub_props.User-Propertyキー・バリューのペアを含むユーザープロパティオブジェクト例:{"foo":"bar"}
    pub_props.*その他のメッセージプロパティのキー・バリュー例:Content-Type: JSON
    payloadメッセージ内容
    message_received_atメッセージ受信タイムスタンプ(ミリ秒)
    idメッセージ ID
    dupメッセージが重複かどうか

MQTT Source の作成は完了しましたが、サブスクライブしたデータはローカル EMQX に直接パブリッシュされません。次に、Source がサブスクライブしたメッセージをローカル EMQX に転送するためのメッセージ再パブリッシュアクションを作成します。

再パブリッシュアクションの作成

  1. ルール作成ページの右側 Action Outputs タブに切り替え、Add Action ボタンをクリックします。Type of Action ドロップダウンから Republish アクションを選択します。

  2. メッセージ再パブリッシュの設定を行います。

    • Topicsub/${topic} と入力し、元のトピックに sub/ プレフィックスを追加して転送します。例えば元のメッセージトピックが f/1 の場合、EMQX に転送されるトピックは sub/f/1 となります。
    • QoS012${qos} から選択可能。プレースホルダーで他のフィールドから設定も可能です。ここでは元メッセージの QoS に従うため ${qos} を選択します。
    • Retain:メッセージをリテインとしてパブリッシュするかどうかを true または false から選択可能。プレースホルダーで他のフィールドから設定も可能です。ここでは false を選択します。
      • データソースが MQTT Source のため、${flags.retain} は使用できません。
      • ${retain} と入力すると MQTT Source が受信したリテインフラグに従います。MQTT 5.0 接続の場合は Source で Retain As Published を有効にして上流のリテインフラグを保持してください。
    • Payload:転送メッセージのペイロード生成に使用。デフォルトは空欄でルール出力結果を転送します。ここでは ${payload} と入力し、ペイロードのみを転送します。
  3. Add ボタンをクリックしてアクション作成を完了します。ルール作成ページに戻り、新しいアクションが Action Outputs タブに追加されます。

  4. ルール作成ページ下部の Create ボタンをクリックし、ルール作成を完了します。

これでルールの作成が完了しました。Integration -> Rules ページで新規ルールを確認できます。Source タブをクリックすると新しい MQTT Source が表示されます。

また、Integration -> Flow Designer をクリックするとトポロジーを確認できます。トポロジーでは MQTT Source からのメッセージが再パブリッシュアクションを通じて sub/${topic} に転送される流れを視覚的に把握できます。

MQTT Broker Source を使ったルールのテスト

MQTTX CLI を使い、外部 MQTT サービスの f/# トピックから EMQX の sub/${topic} トピックへメッセージをブリッジするルールをテストできます。外部 MQTT サービスの f/1 トピックにメッセージをパブリッシュすると、EMQX の sub/f/1 トピックに転送されるはずです。

  1. EMQX の sub/# トピックをサブスクライブします。

    bash
    mqttx sub -t sub/# -q 1 -v
  2. MQTTX で外部 MQTT サービスの f/1 トピックにメッセージをパブリッシュします。

    bash
    mqttx pub -t f/1 -m "I'm from broker.emqx.io" -r -h broker.emqx.io
  3. MQTTX で sub/f/1 トピックにメッセージが届けば、外部 MQTT サービスから EMQX へのメッセージ転送が成功したことを示します。

    bash
    [2024-1-31] [16:49:22] › topic: sub/f/1
    payload: I'm from broker.emqx.io