Skip to content

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

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

動作原理

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

  • 送信メッセージ:ローカルのトピックからリモート MQTT サービスの指定トピックへメッセージをパブリッシュする。
  • 受信メッセージ:リモート MQTT サービスのトピックをサブスクライブし、そのメッセージを現在のデプロイメントに転送する。

EMQX Cloud は同一接続上で複数のブリッジルールを設定可能で、それぞれ異なるトピックマッピングやメッセージ変換ルールを持たせることができ、メッセージルーティングに類似した機能を実現します。ブリッジング中はルールエンジンを介してメッセージのフィルタリング、拡充、変換を行い、転送前に処理を加えることも可能です。

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

EMQX Cloud-MQTT データ統合

特徴と利点

MQTT ブローカーのデータ統合は以下の特徴と利点を持ちます。

  • 広範な互換性:標準 MQTT プロトコルを用いて、AWS IoT Core や Azure IoT Hubs など主要な IoT プラットフォーム、さらにオープンソースや業界標準の MQTT ブローカーと統合可能です。
  • 柔軟なトピックマッピング:トピックにプレフィックスを追加したり、クライアントのコンテキスト情報(クライアントID、ユーザー名など)を用いて動的にトピックを構築し、カスタマイズされたメッセージルーティングを実現します。
  • 高性能:コネクションプーリングや共有サブスクリプションなどの最適化により、スループットを向上させ、レイテンシを低減します。
  • ペイロード変換:SQL ベースの処理により、メッセージペイロードの抽出、フィルタリング、拡充、変換を行い、転送前にメッセージを加工可能です。
  • メトリクス監視:メッセージ数、成功率・失敗率、スループットなどのリアルタイムメトリクスを提供し、統合の健全性とパフォーマンスを監視できます。

はじめる前に

前提条件

MQTT 接続情報の準備

MQTT ブローカーのデータ統合を作成する前に、リモート MQTT サービスの接続情報を入手してください。主な項目は以下の通りです。

  • MQTT サービスアドレス:ターゲット MQTT サービスのアドレスとポート(例:broker.emqx.io:1883)。
  • ユーザー名:接続に必要なユーザー名。認証不要の場合は空欄で構いません。
  • パスワード:接続に必要なパスワード。認証不要の場合は空欄で構いません。
  • プロトコルタイプ:ターゲットサービスが TLS を有効にしているか、TCP/TLS 上の MQTT を使用しているかを確認してください。EMQX Cloud の MQTT ブリッジは現時点で MQTT over WebSocket や MQTT over QUIC などのプロトコルはサポートしていません。
  • プロトコルバージョン:ターゲット MQTT サービスで使用されている MQTT バージョン。EMQX Cloud は MQTT 3.1、3.1.1、MQTT 5.0 をサポートしています。

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

クラスター モードに関する注意

EMQX Cloud がクラスター モードで稼働している場合やコネクションプールが有効な場合、同一のクライアントIDを用いて複数ノードが同じ MQTT サービスに接続すると、デバイスの競合が発生しやすくなります。そのため、MQTT メッセージブリッジでは固定のクライアントID設定は現在サポートしていません。

ネットワーク設定

データ統合を構成する前に、EMQX Cloudのデプロイメントを作成し、EMQX Cloudと対象サービス間のネットワーク接続を確立していることを確認してください。

  • Dedicated Flexデプロイメントの場合

    EMQX CloudのVPCと対象サービスのVPC間でVPCピアリング接続を作成します。ピアリング接続が確立されると、EMQX Cloudは対象サービスのプライベートIPアドレスを介してアクセス可能になります。

    パブリックIP経由でのアクセスが必要な場合は、NATゲートウェイを構成してアウトバウンド接続を有効にしてください。

  • BYOC(Bring Your Own Cloud)デプロイメントの場合

    BYOCデプロイメントが稼働しているVPCと対象サービスをホストするVPC間でVPCピアリング接続を作成します。ピアリングが確立されると、対象サービスのプライベートIPアドレスを介してアクセス可能になります。

    対象サービスにパブリックIP経由でアクセスする必要がある場合は、クラウドプロバイダーのコンソールを使用してBYOC VPCにNATゲートウェイを構成してください。

コネクターの作成

ここでは、EMQX のオンライン MQTT サーバーを例に、リモート MQTT サーバーとの接続設定方法を案内します。

データ統合のルールを作成する前に、MQTT サービスにアクセスするための MQTT ブローカーコネクターを作成する必要があります。コネクターは MQTT(Sink)および MQTT(Source)の両方で使用可能です。

  1. デプロイメントメニューから データ統合 を選択し、データ転送 カテゴリの中の MQTT (Sink) を選択します。既にコネクターを作成済みの場合は、新規コネクター を選択し、同様に MQTT (Sink) を選択します。

  2. コネクター名:システムが自動的にコネクター名を生成します。

  3. 接続情報を設定します。

    • MQTT ブローカー:TCP/TLS 上の MQTT のみサポート。ここでは broker.emqx.io:1883 を設定します。

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

    • ユーザー名パスワード:このサーバーは認証不要のため空欄で構いません。

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

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

    • 静的 ClientId エントリー:(上級者向け)Azure IoT Hubs などのサービスに接続する際に安定した接続を確保するため、コネクターに静的クライアントIDを設定できます。設定方法は静的 ClientId の設定を参照してください。

      TIP

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

  4. その他の設定はデフォルトのままにします。

  5. テスト ボタンをクリックし、MQTT ブローカーにアクセス可能であれば connector available のメッセージが返されます。

  6. 新規作成 ボタンをクリックして作成を完了します。

以降、このコネクターを基にデータブリッジルールを作成できます。

コネクションプールとクライアントID生成ルール

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

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

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

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

bash
myprefix:foo2bd61c44:1

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

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

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

  1. コネクター作成時に 詳細設定 をクリックします。

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

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

    • ノード名:クライアントIDを割り当てるノード名を指定します。例:emqx@10.0.0.1。現在のデプロイメントの EMQX ノード名はチケット提出で取得してください。
    • クライアントID:静的クライアントIDを入力します。例:device1。必要に応じて複数のクライアントIDを追加できます。
      • ユーザー名:(任意)このクライアントIDの認証に使用するユーザー名。
      • パスワード:(任意)このクライアントIDの認証に使用するパスワード。プラットフォームにより、デバイス固有のキーやシークレット、証明書などが該当します(例:Azure IoT Hubsの認証キー)。

    設定例

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

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

MQTT (Sink) でルールを作成する

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

  1. ルールエリアの 新規ルール をクリックするか、作成したコネクターの 操作 列にある新規ルールアイコンをクリックします。

  2. 利用したい機能に基づき、SQL エディターでルールを設定します。ここではクライアントが temp_hum/emqx トピックに温度・湿度メッセージを送信した際にエンジンをトリガーする例を示します。SQL は以下の通りです。

    sql
     SELECT
       topic,
       payload
     FROM
       "temp_hum/emqx"

    TIP

    初心者の方は SQL ExamplesTry It Out をクリックして、SQL ルールを学習・試行できます。

  3. 次へ をクリックしてアクションを追加します。

  4. コネクター ドロップダウンから先ほど作成したコネクターを選択します。

  5. EMQX Cloud から外部 MQTT サービスへメッセージをパブリッシュするための情報を設定します。

    • トピック:外部 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} と入力し、ペイロードのみを転送します。
  6. その他の設定はデフォルト値を使用し、確認 ボタンをクリックしてルール作成を完了します。

  7. 新規ルール作成成功 のポップアップで ルールに戻る をクリックし、データ統合の設定チェーンを完了します。

作成成功後、ルール作成ページに戻ります。ルール リストに新規作成したルールが表示され、操作 リストにはデータ転送アクションが表示されます。

ルールのテスト

MQTTX を使って温湿度データの送信をシミュレーションすることを推奨しますが、他の任意のクライアントでも構いません。

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

    bash
    mqttx sub -t pub/# -q 1 -h broker.emqx.io -v
  2. MQTTX でデプロイメントに接続し、以下のトピックにメッセージを送信します。

    • トピック:temp_hum/emqx

    • クライアントID:test_client

    • ペイロード:

      json
      {
        "temp": "27.5",
        "hum": "41.8"
      }
  3. MQTTX で pub/temp_hum/emqx トピックをサブスクライブし、メッセージを受信できれば、EMQX Cloud から外部 MQTT サービスへのメッセージ転送が成功していることを示します。

    bash
    [2024-3-21] [10:43:13] › topic: pub/temp_hum/emqx
    payload:
    { "temp": "27.5", "hum": "41.8"}

MQTT (Source) でルールを作成する

このセクションでは、リモート MQTT サービスから現在のデプロイメントへデータを転送するルールの作成方法を示します。リモート MQTT サービスからのサブスクライブを実現するために MQTT (Source) と、サブスクライブしたデータを転送するためのメッセージ再パブリッシュアクションを両方作成する必要があります。

MQTT (Source) コネクターと MQTT (Sink) コネクターの作成方法は同じです。詳細はコネクターの作成を参照してください。

  1. コネクター一覧の 操作 列にある新規ルールアイコンをクリックするか、ルール一覧の 新規ルール をクリックして 新規ルール ページに入ります。

  2. ソースデータルールでは、まず入力アクションを設定します。ルール編集ページで自動的に入力アクション設定がポップアップするか、パネル右側の アクション(Source) -> 新規アクション を選択し、MQTT (Source) を選択して 次へ をクリックします。

    • 先ほど作成したコネクターを選択します。
    • 共有サブスクリプショントピック:MQTT Source がリモート MQTT サービスからメッセージをサブスクライブするための共有サブスクリプショントピックを入力します。メッセージ重複を防ぐため必須です。トピックは +# のワイルドカードをサポートします。例として $share/cloud_bridge/temp_hum/emqx と入力すると、共有サブスクリプショングループ cloud_bridge を通じてリモートの temp_hum/emqx トピックにパブリッシュされたメッセージをサブスクライブします。
    • QoS:リモートサブスクリプションのメッセージ QoS を選択します。
    • No Local:同じコネクターでパブリッシュしたメッセージがこのソースでサブスクライブされるトピックに送られる場合に、リモートブローカーがそれらのメッセージをこのソースに再送しないようにするオプションです。コネクターが MQTT 5.0 を使用している場合に有効です。
    • Retain As Published:MQTT 5.0 のリモートサブスクリプションで Retain As Published オプションを設定します。有効にすると、上流ブローカーは転送メッセージの retain フラグを保持し、無効にするとクリアします。コネクターが MQTT 5.0 を使用している場合に有効です。

    確認 ボタンをクリックして設定を完了します。

  3. SQL エディター のデータソースフィールドが更新されます。

    sql
    SELECT
      *
    FROM
      "$bridges/mqtt:source-d1f51e81"
  4. 次へ をクリックして出力アクションの作成を開始します。

  5. 新規出力アクションで 再パブリッシュ を選択します。

  6. 出力アクション情報を設定します。

    • トピック:転送先の MQTT トピック。${var} 形式のプレースホルダーをサポートします。ここでは sub/${topic} と入力し、元のトピックに sub/ プレフィックスを付けて転送します。例えば元のメッセージトピックが t/1 の場合、転送先トピックは sub/t/1 になります。
    • QoS:メッセージパブリッシュの QoS。012${qos} から選択可能で、他のフィールドから QoS を設定するプレースホルダーも使用できます。ここでは元のメッセージの QoS に従うため ${qos} を選択します。
    • Retain:メッセージをリテインとしてパブリッシュするかどうか。truefalse${flags.retain} から選択可能で、他のフィールドからリテインフラグを設定するプレースホルダーも使用できます。ここでは元のメッセージのリテインフラグに従うため ${flags.retain} を選択します。
    • メッセージテンプレート:転送メッセージのペイロード生成テンプレート。デフォルトは空欄でルール出力結果を転送します。ここでは ${payload} と入力し、ペイロードのみを転送します。
  7. その他の設定はデフォルト値を使用し、確認 ボタンをクリックして出力アクションの作成を完了します。

作成成功後、ルール作成ページに戻ります。現在、再パブリッシュアクションは アクション(Source) に表示されません。必要に応じてルール編集ボタンをクリックすると、ルール設定の下部に再パブリッシュ出力アクションが表示されます。

ルールのテスト

作成したルールは、外部 MQTT サービスの temp_hum/emqx トピックからメッセージを受け取り、現在のデプロイメントの sub/${topic} トピックに転送するよう設定されています。ソース設定では共有サブスクリプショントピックに $share/cloud_bridge/temp_hum/emqx を指定しています。したがって、外部 MQTT サービスの temp_hum/emqx トピックにメッセージをパブリッシュすると、この MQTT Source が受信し、sub/temp_hum/emqx トピックに転送します。

以下は MQTTX を使って外部 MQTT サービスにメッセージを送信し、転送されたトピックをサブスクライブしてメッセージを受信する手順です。

  1. MQTTX で現在のデプロイメントのトピック sub/# をサブスクライブします。

  2. MQTTX で外部 MQTT サービスの temp_hum/emqx トピックにメッセージをパブリッシュします。

    json
    {
       "temp": 55,
       "hum": 32
    }
  3. MQTTX は現在のデプロイメントの sub/temp_hum/emqx トピックでメッセージを受信します。

    json
    {
       "temp": 55,
       "hum": 32
    }

トラブルシューティング

このセクションでは、MQTT ブローカーのデータ統合で MQTT コネクターを使用する際に発生しやすい問題とその解決策を紹介します。

MQTT コネクター画面に「クラスター内ノード間の不整合」が表示される

説明

クラスターのノード間で MQTT ブローカーコネクターの状態が不整合に見える問題です。

考えられる原因

  1. リモート MQTT ブローカーの接続制限

    リモート MQTT ブローカーが接続数の上限に達しているか、制限を課しており、EMQX クラスターのノードが正常に接続できていない可能性があります。

  2. VPC ピアリング接続の制限

    VPC ピアリング接続を利用している場合、セキュリティグループのルールにより EMQX ノードからリモートブローカーへのアクセスが制限されている可能性があります。

解決策

  • リモートブローカーの制限の場合

    MQTT コネクター設定でコネクションプールサイズを現在の半分に減らして再試行してください。

  • VPC 制限の場合

    EMQX がデプロイされている VPC サブネット(通常は 10.x.x.0/24)からリモートブローカーへのアクセスを許可するセキュリティグループルールを確認してください。

  • 問題が解決しない場合は、EMQ サポートにサポートチケットを提出してください。

MQTT コネクター設定の更新・保存時に「Axios Error: timeout」が発生する

説明

共有サブスクリプションを使用する MQTT コネクターの設定を更新・保存する際にタイムアウトエラーが発生します。

考えられる原因

ソースコネクターに複数のコネクションプールとサブスクリプショントピックが含まれる場合、更新時にすべての共有サブスクリプション接続を再確立する必要があります。

設定が大規模(例:2ノードで計16プール、クライアントごとに複数トピック)だと、リモートブローカーは再確立時にサブスクリプション確認応答(SUBACK)を順次処理します。

共有サブスクリプションのパフォーマンスレイテンシにより初期化処理が長時間かかり、コンソールや API リクエストが30秒のタイムアウトを超えて timeout エラーとなる場合があります。

解決策

  1. データソースごとにコネクターを分割し、各コネクターが2~4ソースを処理するようにします。
  2. 各コネクターのコネクションプールサイズを 2 に設定し、再初期化性能を向上させます。

MQTT アクションのデータが転送されず、メッセージが期限切れで破棄される

説明

一部の MQTT アクションデータが正常に転送されず、監視でメッセージが期限切れで破棄されたことが確認されます。

考えられる原因

リモートブローカーが現在のスループットに対応できないか、コネクションプールサイズが不足していると、メッセージがバッファキューに蓄積されます。デフォルトではこのキュー内のメッセージは TTL(有効期限)が45秒に設定されており、期限切れになると破棄されます。

解決策

  1. リモートブローカーが現在の TPS(トランザクション数/秒)に対応可能か確認してください。
  2. スループット要件に応じてコネクションプールサイズを増やしてください。
  3. アクション設定の リクエストタイムアウト 値を調整し、メッセージの早期期限切れを防止してください。