Skip to content

MQTTデータをHTTPサーバーに取り込む

EMQX CloudのHTTPサービスデータ統合は、EMQXと外部HTTPサービスを迅速に連携させる方法を提供します。リクエストメソッドやリクエストデータ形式の柔軟な設定をサポートし、HTTPSによる安全な通信や認証機構も備えています。クライアントのメッセージやイベントデータをリアルタイムに効率的かつ柔軟に送信でき、IoTデバイスの状態通知、アラート通知、データ統合などのシナリオを実現します。

本ページでは、HTTPサービスデータ統合の機能について詳しく紹介するとともに、HTTPサーバーコネクターの作成、ルールの作成、ルールのテストなど、データ統合の実践的な手順を案内します。

動作の仕組み

HTTPサービスデータ統合はEMQX Cloudに標準搭載された機能で、外部サービスとの連携を簡単に設定できます。HTTPサービスを利用することで、ユーザーは好みのプログラミング言語やフレームワークでカスタムかつ柔軟な複雑なデータ処理ロジックを実装できます。

http frame

EMQX Cloudは、設定されたデータ統合を通じてデバイスのイベントやデータをHTTPサービスに転送します。ワークフローは以下の通りです。

  1. デバイスがEMQX Cloudに接続:IoTデバイスが正常に接続すると、デバイスIDや送信元IPアドレスなどの属性を含むオンラインイベントが発生します。
  2. デバイスがメッセージをパブリッシュ:デバイスはテレメトリや状態データをEMQX Cloudに報告し、MQTTプロトコルを通じて特定のトピックにメッセージをパブリッシュしてルールをトリガーします。
  3. ルールエンジンがメッセージを処理:組み込みのルールエンジンはトピックマッチングに基づき特定のソースからのメッセージやイベントを処理します。ルールエンジンは対応するルールをマッチングし、データ形式の変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などを行います。
  4. HTTPサービスに送信:ルールのトリガーによりメッセージがHTTPサービスのイベントに送信されます。ユーザーはルール処理結果からデータを抽出し、リクエストヘッダーやボディ、URLを動的に構築でき、外部サービスとの柔軟なデータ連携を実現します。

イベントやメッセージデータがHTTPサービスに送信された後は、以下のように柔軟に処理できます。

  • デバイスの状態更新やイベント記録を実装し、データに基づいたデバイス管理システムを開発する。
  • メッセージデータをデータベースに書き込み、軽量なデータストレージ機能を実現する。
  • ルールSQLでフィルタリングされた異常データに対して、HTTPサービス経由でアラーム通知システムを直接呼び出し、デバイス異常監視を行う。

特長とメリット

EMQX CloudのHTTPサービス統合を利用することで、以下のメリットをビジネスにもたらします。

  • より多くの下流システムへのデータ連携:HTTPサービスを通じてMQTTデータを分析プラットフォームやクラウドサービスなど多様な外部システムに容易に統合し、マルチシステムのデータ分配を実現します。
  • リアルタイムな応答と業務プロセスのトリガー:HTTPサービスにより外部システムがMQTTデータをリアルタイムに受信し、業務プロセスを迅速にトリガーできます。例えば、アラームデータ受信による業務ワークフローの起動などです。
  • カスタムデータ処理:外部システム側で受信データの二次処理が可能で、より複雑な業務ロジックをEMQXの機能に制限されずに実装できます。
  • 疎結合な連携:HTTPサービスはシンプルなHTTPインターフェースを利用し、システム間の疎結合な連携手法を提供します。

まとめると、HTTPサービスはリアルタイムかつ柔軟でカスタマイズ可能なデータ統合機能を提供し、多様で豊富なアプリケーション開発のニーズに応えます。

はじめる前に

このセクションでは、EMQX CloudでHTTPサービスデータ統合を作成するための準備作業を紹介します。

前提条件

ネットワーク設定

データ統合を構成する前に、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ゲートウェイを構成してください。

簡単なHTTPサービスのセットアップ

以下の例を使って簡単なWebサーバーを作成します。

python
from http.server import HTTPServer, BaseHTTPRequestHandler

class SimpleHTTPRequestHandler(BaseHTTPRequestHandler):
    def do_GET(self):
       self.send_response(200)
       self.end_headers()
       self.wfile.write(b'Hello, world!')

    def do_POST(self):
       content_length = int(self.headers['Content-Length'])
       body = self.rfile.read(content_length)
       print("Received POST request with body: " + str(body))
       self.send_response(201)
       self.end_headers()

httpd = HTTPServer(('0.0.0.0', 8080), SimpleHTTPRequestHandler)
httpd.serve_forever()

HTTPサーバーコネクターの作成

データ統合ルールを作成する前に、HTTPサービスにアクセスするためのHTTPサーバーコネクターを作成する必要があります。

  1. デプロイメントに移動し、左側のナビゲーションメニューからデータ統合をクリックします。
  2. 初めてコネクターを作成する場合は、Webサービスカテゴリの中からHTTPサーバーを選択します。すでにコネクターを作成済みの場合は、新規コネクターを選択し、続いてWebサービスカテゴリのHTTPサーバーを選択します。
  3. 新規コネクターページで以下の項目を設定します。
    • コネクター名:システムが自動で生成する名前を使うか、自分で名前を付けられます。例としてmy_httpserverを使用できます。
    • URL:WebサービスのURLを入力します。ネットワーク経由で正常にアクセスできることを確認してください。
    • その他の設定はデフォルト値で構いません。必要に応じてHTTPリクエストヘッダーのキーと値を設定できます。
    • OAuth2クライアント認証情報:対象のHTTPサーバーがOAuth2認証を必要とする場合、このオプションを有効にするとEMQXがアクセストークンを取得し、コネクター経由のリクエストにトークンを追加します。この機能はEMQX v6.1.4以降のデプロイメントで利用可能です。詳細はOAuth2クライアント認証情報の設定を参照してください。
  4. テストボタンを押して接続をテストします。Webサービスにアクセス可能であれば成功メッセージが表示されます。
  5. 新規作成ボタンを押して作成を完了します。

OAuth2クライアント認証情報の設定

EMQX v6.1.4以降のデプロイメントでは、HTTPサーバーコネクターがOAuth2クライアント認証情報をサポートします。このオプションを有効にすると、EMQXは設定されたトークンエンドポイントからアクセストークンを取得・キャッシュし、自動で更新します。EMQXが対象HTTPサーバーにデータを送信する際、Authorization: Bearer <access_token>ヘッダーをリクエストに追加します。

HTTPサーバーコネクターの作成または編集時にOAuth2クライアント認証情報を有効にし、以下の項目を設定してください。

項目説明
トークンエンドポイント必須。アクセストークンをリクエストするOAuth2認可サーバーのエンドポイント。URLはHTTPまたはHTTPSで、ユーザー情報を含まない必要があります。
クライアントID必須。アクセストークン取得に使用するOAuth2クライアントID。
クライアントシークレット必須。アクセストークン取得に使用するOAuth2クライアントシークレット。
スコープ任意。アクセストークンに要求するOAuth2スコープ。
トークンリクエストタイムアウトトークンエンドポイントへのHTTPリクエストのタイムアウト。デフォルトは5秒です。
トークンエンドポイントTLSトークンエンドポイントへのTLS接続を有効にします。この設定は対象HTTPサーバーへのTLS接続を制御するTLS有効化とは独立しています。

EMQXはapplication/x-www-form-urlencodedのコンテンツタイプでPOSTリクエストをトークンエンドポイントに送信します。リクエストボディにはgrant_typeclient_idclient_secret、および任意のscopeを含みます。トークンエンドポイントは200レスポンスでJSONボディを返し、access_tokenを含む必要があります。token_typeexpires_inも返せます。token_typeがある場合はBearerでなければならず、expires_inは正の整数である必要があります。

重要なお知らせ

  • OAuth2を有効にした場合、HTTPサーバーコネクターまたはそのアクションにAuthorizationヘッダーを設定しないでください。EMQXは自動生成されるBearer認証ヘッダーと競合するため設定を拒否します。
  • トークンエンドポイントはクライアントIDとクライアントシークレットをリクエストボディのフォームフィールドとして受け入れる必要があります。HTTP Basic認証のAuthorizationヘッダーによる認証はサポートしていません。

ルールの作成

次に、書き込むデータを指定するルールを作成し、処理済みデータをHTTPサーバーに転送するアクションをルールに追加します。

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

  2. SQLエディターにルールマッチングのSQL文を入力します。以下のSQL例は、temp_hum/emqxトピックに送信されたメッセージから、メッセージの報告時間、クライアントID、メッセージボディ(ペイロード)の温度と湿度を読み取ります。

    sql
    SELECT 
    
    timestamp as up_timestamp, clientid as client_id, payload.temp as temp, payload.hum as hum
    
    FROM
    
    "temp_hum/emqx"

    Try It Out機能でデータ入力をシミュレートし、結果をテストできます。

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

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

  5. 以下の情報を設定します。

    • アクション名:システムが自動生成する名前を使うか、自分で名前を付けられます。

    • メソッド:HTTPリクエストメソッドとしてPOSTを選択します。

    • URLパス:アクション専用のリクエストパスを設定できます。これはコネクターのURL設定に追加され、完全なURLを構成します。テンプレート変数も利用可能です。まずルールSQLで関連変数を定義します(例:select clientid as device_id)。その後、HTTPアクションのURLパスに${device_id}のように変数を組み込めます。

    • ヘッダー:アクション固有のリクエストヘッダーを設定するか、コネクターの設定ヘッダーをそのまま使用できます。

    • ボディ:以下のメッセージコンテンツテンプレートのように、ルールから出力されたフィールドをアクションのリクエストボディに設定します。

      json
      {"up_timestamp": ${up_timestamp}, "client_id": ${client_id}, "temp": ${temp}, "hum": ${hum}}
  6. 確定ボタンをクリックしてルール作成を完了します。

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

ルールのテスト

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

  1. MQTTXでデプロイメントに接続し、以下のトピックにメッセージを送信します。

    • トピック:temp_hum/emqx

    • ペイロード:

      json
      {
        "temp": "27.5",
        "hum": "41.8"
      }
  2. メッセージがHTTPサービスに転送されているか確認します。

    bash
    py server.py
    
    Received POST request with body: b'[\n "temp": "27.5",\n "hum": "41.8"\n)127.0.0.1 - -[18/Dec/2023 14:50:44]"POST  HTTP/1.1" 201 -
  3. コンソールで運用データを確認します。ルール一覧のルールIDをクリックすると、ルールの統計情報やそのルールに属する全アクションの統計情報が表示されます。