他の MQTT サービスとのブリッジ
MQTT ブローカーのデータ統合は、EMQX に別の EMQX クラスターや他の MQTT サービスとメッセージブリッジで接続する機能を提供し、ネットワークやサービスを跨いだデータの相互作用と通信を可能にします。本ページでは、EMQX における MQTT メッセージブリッジの動作原理を紹介し、メッセージブリッジの作成と検証の実践的な手順を説明します。
動作原理
ブリッジング中、EMQX はターゲットサービスとクライアントとして MQTT 接続を確立し、パブリッシュ・サブスクライブモデルを通じて双方向のメッセージ送受信を実現します。
- 送信メッセージ(Sink):ローカルのトピックからメッセージをパブリッシュし、リモート MQTT サービスの指定トピックに送信します。
- 受信メッセージ(Source):リモート MQTT サービスのトピックをサブスクライブし、そのメッセージを EMQX ローカルに転送します。
EMQX は同一接続上で複数のブリッジルールを設定可能で、それぞれ異なるトピックマッピングやメッセージ変換ルールを持ち、メッセージルーティングに類似した機能を実装できます。ブリッジング中はルールエンジンを介してメッセージのフィルタリング、拡充、変換処理も行えます。
以下の図は、EMQX と他の 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]:port、mqtt://host:port、mqtt://[IPv6]:port、mqtts://host:port、mqtts://[IPv6]:portです。ポートが省略された場合、EMQX はデフォルトの MQTT ポート1883を使用します。その他の URI スキームはサポートされません。mqttとmqttsスキームはアドレス解析のためのみ使用され、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 サーバーとの接続を設定する手順を案内します。
ダッシュボードの Integration -> Connector ページに移動します。
ページ右上の Create をクリックします。
コネクタータイプ一覧から MQTT Broker を選択し、Next をクリックします。
コネクターの name を入力します。英数字の組み合わせで、例として
my_mqtt_bridgeとします。接続情報を設定します:
MQTT Broker:TCP/TLS 上の MQTT のみサポートします。MQTT サービスアドレスを入力します。例:
broker.emqx.io:1883、[::1]:1883、mqtt://broker.emqx.io:1883。TLS 対応の MQTT リスナーの場合はコネクター設定で TLS を別途構成してください。ClientID Prefix:空欄でも構いません。実際の運用ではクライアント ID プレフィックスを指定するとクライアント管理が容易になります。EMQX はクライアント ID プレフィックスと接続プールのサイズに基づき自動的にクライアント ID を生成します。詳細は 接続プールとクライアント ID 生成ルール を参照してください。
Username と Password: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 を自動生成します。
[Client ID Prefix]:{Connector Name}{8桁のランダム文字列}:{接続プール内の接続シーケンス番号}例として、クライアント ID プレフィックスが myprefix、コネクター名が foo の場合、実際のクライアント ID は以下のようになります。
myprefix:foo2bd61c44:1バージョン 5.4.1 以降、EMQX は MQTT クライアント ID の長さを 23 バイトに制限しています。クライアント ID がこれを超える場合はハッシュ値に置き換えられます。プレフィックスやコネクター名が長すぎるとユーザー体験が悪化する可能性があります。
この問題に対処するため、バージョン 5.7.1 以降、EMQX は以下のルールを実装しています。
- プレフィックスなし:動作は変わらず、長い(23 バイト超)クライアント ID は 23 バイトにハッシュ化されます。
- プレフィックスあり:
- プレフィックスが最大 19 バイトまで:プレフィックスは保持され、残りのクライアント ID は 4 バイトのハッシュに変換され、全体の長さは 23 バイト以内に収まります。
- プレフィックスが 20 バイト以上:EMQX は設定されたプレフィックスを使用し、クライアント ID の短縮は行いません。
静的クライアント ID の設定
統合で使用できるクライアント ID が有限の場合、コネクター設定時に個々のノードに静的クライアント ID セットを割り当てることが可能です。静的クライアント ID を設定するには、EMQX クラスターの各ノードに対してクライアント ID のリストを提供します。各クライアント ID に対して対応するユーザー名とパスワードを指定できます。これは Azure IoT Hubs のように各デバイス(クライアント ID)に固有の認証情報が必要なシナリオに特に有用です。
静的クライアント ID を設定する手順は以下の通りです。
Static ClientId Entries セクションで Add ボタンをクリックし、新しい静的クライアント ID エントリーを追加します。必要に応じて複数のノードに対して複数のエントリーを追加できます。
各エントリーに以下の項目を入力します。
- Node Name:クライアント ID を割り当てるノード名。例:
emqx@10.0.0.1 - Client ID:静的クライアント ID。例:
device1。必要に応じてノードに複数のクライアント ID を追加できます(Add ボタンで追加)。- Username:(任意)このクライアント ID に対応する認証用ユーザー名。
- Password:(任意)このクライアント ID に対応する認証用パスワード。プラットフォームによりデバイス固有のキー、シークレット、証明書などが使われます(例:Azure IoT Hubs の認証キー)。
設定例:
Node Client ID Username (任意) Password (任意) emqx@10.0.0.1clientid1username1secret1clientid3emqx@10.0.0.2clientid2username2emqx@10.0.0.3clientid4clientid5- Node Name:クライアント ID を割り当てるノード名。例:
設定ファイルでノードごとに static_clientids パラメータを個別に定義することも可能です。
静的クライアント ID が設定されている場合、これらのクライアント ID を使った MQTT 接続のみが開始され、pool_size や clientid_prefix など動的クライアント ID に関する設定は無効になります。
MQTT ブローカー Sink を使ったルールの作成
ここでは、リモート MQTT サービスに転送するデータを指定するルールの作成方法を示します。
ダッシュボードの Integration -> Rules ページに移動します。
ページ右上の Create をクリックします。
ルール ID に
my_ruleを入力します。SQL Editor に、
t/#トピックからの MQTT メッセージをリモート MQTT サーバーに保存するルールを入力します。ルール SQL は以下の通りです。sqlSELECT * FROM "t/#"MQTTv5 プロトコルを使用している場合、
pub_propsフィールドがルール SQL に含まれていれば、パブリッシュプロパティはそのままリモートブローカーに転送されます。上記のSELECT * FROM t/#の例が該当します。さらにユーザープロパティを動的に追加したい場合は、ルールが出力する
pub_propsフィールドに含めることができます。例えば、以下のルールはペイロードから取得したキーと値をユーザープロパティとして追加します。sqlSELECT *, map_put(concat('User-Property.', payload.extra_key), payload.extra_value, pub_props) as pub_props FROM 't/#'Action Type ドロップダウンから
MQTT Brokerを選択してアクションを追加します。Action ドロップダウンはデフォルトのCreate Actionのままにします。この例では新しい Sink を作成し、ルールに追加します。Sink の名前と説明をフォームに入力します。
Connector ドロップダウンから先ほど作成した
my_mqtt_bridgeコネクターを選択します。あるいは、ドロップダウン横の作成ボタンをクリックして新しいコネクターを作成し、コネクターの作成 の設定パラメータを利用できます。EMQX から外部 MQTT サービスにメッセージをパブリッシュするための Sink 情報を設定します。
- Topic:外部 MQTT サービスにパブリッシュするトピック。
${var}プレースホルダーをサポートします。ここではpub/${topic}と入力し、元のトピックにpub/プレフィックスを付けて転送します。例えば元のメッセージトピックがt/1の場合、外部 MQTT サービスに転送されるトピックはpub/t/1になります。 - QoS:メッセージパブリッシュの QoS。ドロップダウンから
0、1、2、${qos}のいずれかを選択可能で、他のフィールドから QoS を設定するプレースホルダーも使用できます。ここでは元メッセージの QoS に従うため${qos}を選択します。 - Retain:メッセージをリテインとしてパブリッシュするかどうか。
true、false、${flags.retain}を選択可能で、他のフィールドからリテインフラグを設定するプレースホルダーも使用できます。ここでは元メッセージのリテインフラグに従うため${flags.retain}を選択します。 - Payload:転送メッセージのペイロード生成に使うテンプレート。デフォルトは空欄でルール出力結果をそのまま転送します。ここでは
${payload}と入力し、ペイロードのみを転送します。
- Topic:外部 MQTT サービスにパブリッシュするトピック。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のために、1つ以上のフォールバックアクションを定義できます。詳細は フォールバックアクション を参照してください。
その他の設定はデフォルトのままにして、Create ボタンをクリックし Sink の作成を完了します。作成後はルール作成ページに戻り、新しい Sink がルールのアクション出力に追加されます。
ルール作成ページ下部の Create ボタンをクリックしてルール作成を完了します。
これでルールが正常に作成されました。Integration -> Rules ページで新規作成したルールを確認できます。Actions(Sink) タブをクリックすると新しい MQTT ブローカー Sink が表示されます。
また、Integration -> Flow Designer をクリックするとトポロジーを確認できます。トポロジーは t/# トピックのメッセージがルール my_rule で処理された後、リモート MQTT ブローカーに送信される様子を視覚的に表現しています。
MQTT ブローカー Sink を使ったルールのテスト
MQTTX CLI を使って、EMQX の t/# トピックから外部 MQTT サービスの pub/${topic} トピックへメッセージをブリッジするルールをテストできます。EMQX の t/1 トピックにメッセージをパブリッシュすると、外部 MQTT サービスの pub/t/1 トピックに転送されるはずです。
外部 MQTT サービスで
pub/#トピックをサブスクライブします。bashmqttx sub -t pub/# -q 1 -h broker.emqx.io -vMQTTX で
t/1トピックにメッセージをパブリッシュします。bashmqttx pub -t t/1 -m "hello world" -rMQTTX で
pub/t/1トピックをサブスクライブすると、EMQX から外部 MQTT サービスにメッセージが正常に転送されたことが確認できます。bash[2024-1-31] [16:43:13] › topic: pub/t/1 payload: hello worldステップ 1 を繰り返すと、
pub/t/1トピックのリテインメッセージも MQTTX で確認できます。bash[2024-1-31] [16:44:29] › topic: pub/t/1 payload: hello world retain: true
MQTT ブローカー Source を使ったルールの作成
ここでは、リモート MQTT サービスからローカル EMQX へデータを転送するルールの作成方法を示します。MQTT Source とメッセージ再パブリッシュアクションを作成し、リモート MQTT サービスから EMQX へのサブスクライブと、サブスクライブしたデータの転送を実現します。
MQTT ブローカー Source の作成とルールへの追加
ダッシュボードの Integration -> Rules ページに移動します。
ページ右上の Create をクリックします。
ルール ID に
my_rule_sourceを入力します。ルールのトリガーソース(データ入力)を設定します。ページ右側の Data Inputs タブでデフォルトの Message タイプの入力を削除し、Add Input をクリックして MQTT Source を作成します。
Add Input ポップアップで、Input Type ドロップダウンから
MQTT Brokerを選択し、Source ドロップダウンはデフォルトのCreate Sourceのままにします。この例では新しい Source を作成しルールに追加します。Source の名前と説明をフォームに入力します。
先ほど作成した
my_mqtt_bridgeコネクターをドロップダウンから選択します。あるいは、ドロップダウン横の作成ボタンをクリックして新しいコネクターを作成し、コネクターの作成 の設定パラメータを利用できます。Source 情報を設定し、外部 MQTT サービスから EMQX へのサブスクライブを完了させます。
Topic:サブスクライブするトピック。
+と#のワイルドカードをサポートします。TIP
EMQX がクラスター モードで稼働している場合やコネクターが接続プールを使用している場合、重複メッセージを避けるため共有サブスクリプションを使用する必要があります。
ここでは
$share/1/f/#と入力し、f/#トピックにマッチするすべてのメッセージをサブスクライブします。QoS:サブスクライブの QoS。ドロップダウンから
0または1を選択します。No Local:同一コネクターでメッセージをパブリッシュし、かつこの Source が同じトピックをサブスクライブしている場合に、リモート MQTT ブローカーがそのメッセージを Source に再転送するのを防止します。デフォルトは無効で、MQTT 5.0 接続時のみ有効です。
Retain As Published:上流 MQTT ブローカーから転送されたメッセージの元の
retainフラグを保持します。有効にしない場合、上流ブローカーはretainフラグをクリアします。デフォルトは有効で、MQTT 5.0 接続時のみ有効です。
その他の設定はデフォルトのままにして、Create ボタンをクリックし Source の作成を完了します。Source はルールのデータ入力に追加されます。ルール SQL は以下のように変更されます。
sqlSELECT * FROM "$bridges/mqtt:my_source"ルール SQL は MQTT Source から以下のフィールドを抽出でき、SQL を調整してデータ処理が可能です。デフォルトの SQL で十分です。
フィールド名 説明 topic 発信元メッセージのトピック server 接続された Source のサーバーアドレス retain 上流 MQTT ブローカーから受信したメッセージの retainフラグ。Retain As Published が無効の場合はfalseqos メッセージの QoS(サービス品質) pub_props MQTT 5.0 メッセージプロパティオブジェクト。ユーザープロパティペアやその他属性を含む pub_props.User-Property-Pairs キーと値のペアを含むユーザープロパティ配列。例: {"key":"foo", "value":"bar"}pub_props.User-Property キーと値のペアを含むユーザープロパティオブジェクト。例: {"foo":"bar"}pub_props.* その他のメッセージプロパティのキーと値のペア。例: Content-Type: JSONpayload メッセージ内容 message_received_at メッセージ受信タイムスタンプ(ミリ秒) id メッセージ ID dup メッセージが重複かどうか
これで MQTT Source の作成は完了しましたが、サブスクライブしたデータは直接 EMQX ローカルにパブリッシュされません。次に、Source がサブスクライブしたメッセージを EMQX ローカルに転送するためのメッセージ再パブリッシュアクションを作成します。
メッセージ再パブリッシュアクションの作成
ルール作成ページの右側 Action Outputs タブに切り替え、Add Action ボタンをクリックします。Type of Action ドロップダウンから
Republishアクションを選択します。メッセージ再パブリッシュの設定を行います。
- Topic:
sub/${topic}と入力し、元のトピックにsub/プレフィックスを付けて転送します。例えば元のメッセージトピックがf/1の場合、EMQX に転送されるトピックはsub/f/1になります。 - QoS:
0、1、2、${qos}から選択可能で、他のフィールドから QoS を設定するプレースホルダーも使用できます。ここでは元メッセージの QoS に従うため${qos}を選択します。 - Retain:メッセージをリテインとしてパブリッシュするかどうかを
trueまたはfalseで選択します。プレースホルダーも使用可能です。ここではfalseを選択します。- データソースが MQTT Source のため、
${flags.retain}オプションは適用されません。 ${retain}を入力すると MQTT Source が受信したリテインフラグに従います。MQTT 5.0 接続の場合は Source で Retain As Published を有効にしてください。
- データソースが MQTT Source のため、
- Payload:転送メッセージのペイロード生成に使用します。デフォルトは空欄でルール出力結果を転送します。ここでは
${payload}と入力し、ペイロードのみを転送します。
- Topic:
Add ボタンをクリックしてアクション作成を完了します。ルール作成ページに戻り、新しいアクションが Action Outputs タブに追加されます。
ルール作成ページ下部の Create ボタンをクリックしてルール作成を完了します。
これでルールが正常に作成されました。Integration -> Rules ページで新規作成したルールを確認できます。Source タブをクリックすると新しい MQTT Source が表示されます。
また、Integration -> Flow Designer をクリックするとトポロジーを確認できます。トポロジーでは MQTT Source からのメッセージが再パブリッシュアクションを通じて sub/${topic} に転送される様子が示されています。
MQTT ブローカー Source を使ったルールのテスト
MQTTX CLI を使って、外部 MQTT サービスの f/# トピックから EMQX の sub/${topic} トピックへメッセージをブリッジするルールをテストできます。外部 MQTT サービスの f/1 トピックにメッセージをパブリッシュすると、EMQX の sub/f/1 トピックに転送されるはずです。
EMQX の
sub/#トピックをサブスクライブします。bashmqttx sub -t sub/# -q 1 -vMQTTX で外部 MQTT サービスの
f/1トピックにメッセージをパブリッシュします。bashmqttx pub -t f/1 -m "I'm from broker.emqx.io" -r -h broker.emqx.ioMQTTX で
sub/f/1トピックにパブリッシュされたメッセージを確認できれば、外部 MQTT サービスから EMQX へのメッセージ転送が成功したことを示します。bash[2024-1-31] [16:49:22] › topic: sub/f/1 payload: I'm from broker.emqx.io