GCP Pub/Sub 統合
Google Cloud Pub/Sub は、非常に高い信頼性とスケーラビリティを実現するために設計された非同期メッセージングサービスです。EMQX Cloud は、MQTT データのリアルタイム抽出、処理、分析のために Google Cloud Pub/Sub とのシームレスな統合をサポートしています。Cloud Functions、App Engine、Cloud Run、Kubernetes Engine、Compute Engine などのさまざまな Google Cloud サービスへデータをプッシュすることが可能です。また、Google Cloud から MQTT へデータを配信することもでき、ユーザーが GCP 上で迅速に IoT アプリケーションを構築するのに役立ちます。
本ページでは、EMQX Cloud と GCP Pub/Sub 間のデータ統合について包括的に紹介し、プロデューサー(Sink)およびコンシューマー(Source)統合の作成と検証に関する実践的な手順を提供します。
動作概要
GCP Pub/Sub データ統合は、EMQX Cloud の標準機能として提供されており、MQTT データストリームを Google Cloud とシームレスに連携させ、IoT アプリケーション開発のための豊富なサービスと機能を活用できるよう支援します。

EMQX Cloud は、ルールエンジンと Sink を介して MQTT データを GCP Pub/Sub に転送します。GCP Pub/Sub のプロデューサー役割の例を挙げると、全体の流れは以下の通りです。
- IoT デバイスがメッセージをパブリッシュ: デバイスは特定のトピックを通じてテレメトリやステータスデータをパブリッシュし、ルールエンジンをトリガーします。
- ルールエンジンがメッセージを処理: 組み込みのルールエンジンは、特定のソースからの MQTT メッセージをトピックマッチングに基づいて処理します。ルールエンジンは該当するルールをマッチングし、データフォーマットの変換、特定情報のフィルタリング、コンテキスト情報の付加などの処理を行います。
- GCP Pub/Sub へのブリッジング: ルールはメッセージを GCP Pub/Sub に転送するアクションをトリガーし、データプロパティ、オーダーキー、MQTT トピックと GCP Pub/Sub トピックのマッピングを簡単に設定できます。これにより、より豊かなコンテキスト情報と順序保証を伴うデータ統合が可能となり、柔軟な IoT データ処理を実現します。
MQTT メッセージデータが GCP Pub/Sub に書き込まれた後は、以下のような柔軟なアプリケーション開発が可能です。
- リアルタイムデータ処理と分析: Dataflow、BigQuery、Pub/Sub のストリーミング機能など、強力な Google Cloud のデータ処理・分析ツールを活用し、メッセージデータのリアルタイム処理・分析を行い、価値ある洞察や意思決定支援を得られます。
- イベント駆動型機能: Cloud Functions や Cloud Run などの Google Cloud イベント処理をトリガーし、動的かつ柔軟な機能トリガーと処理を実現します。
- データの保存と共有: Cloud Storage や Firestore などの Google Cloud ストレージサービスにメッセージデータを送信し、大量データの安全な保存・管理を行います。これにより、他の Google Cloud サービスと連携してデータの共有や分析を行い、多様なビジネスニーズに対応できます。
特徴と利点
GCP Pub/Sub とのデータ統合は、以下のような特徴と利点を提供します。
- 堅牢なメッセージングサービス: EMQX と GCP Pub/Sub はともに高可用性とスケーラビリティを備え、大規模なメッセージストリームの確実な受信、配信、処理を保証します。IoT データの順序管理、メッセージの品質保証、パーシステンス(永続化)をサポートし、信頼性の高いメッセージ伝送と処理を実現します。
- 柔軟なルールエンジン: 組み込みのルールエンジンにより、特定のソースメッセージやイベントをトピックマッチングに基づいて処理できます。メッセージのデータフォーマット変換、特定情報のフィルタリング、コンテキスト情報の付加などが可能です。これと GCP Pub/Sub を組み合わせることで、さらなる処理や分析が行えます。
- 豊富なコンテキスト情報: GCP Pub/Sub データ統合を通じて、メッセージにより豊かなコンテキスト情報を付加できます。クライアント属性を Pub/Sub 属性やソーティングキーにマッピングすることができ、後続のアプリケーション開発やデータ処理においてより精密な分析や処理を支援します。
まとめると、EMQX Cloud と GCP Pub/Sub の統合により、高信頼性かつスケーラブルなメッセージ配信が可能となり、データ分析や統合のための豊富なツールとサービスを活用できます。これにより、堅牢な IoT アプリケーションの構築やイベント駆動型の柔軟なビジネスロジックの実装が可能となります。
はじめる前に
このセクションでは、GCP Pub/Sub データ統合の作成を開始する前に必要な準備について説明します。
前提条件
ネットワークの設定
開始前に、EMQX Cloud 上にデプロイメント(EMQX クラスター)を作成し、ネットワークを設定する必要があります。
- Dedicated Flex デプロイメントユーザーの場合: VPC が Google Cloud Platform (GCP) 上にある場合、VPC ピアリング接続を確立せずに直接データ転送が可能です。他のクラウドプラットフォーム上の VPC の場合は、NAT ゲートウェイ を設定し、パブリック IP を介してターゲットコネクターにアクセスしてください。
- BYOC (Bring Your Own Cloud) デプロイメントユーザーの場合: BYOC がデプロイされている VPC とターゲットコネクターが存在する VPC 間でピアリング接続を確立してください。ピアリング接続作成後は、内部ネットワーク IP 経由でターゲットコネクターにアクセス可能です。パブリック IP 経由でリソースにアクセスする必要がある場合は、パブリッククラウドコンソールで BYOC がデプロイされている VPC に対して NAT ゲートウェイを設定してください。
GCP でサービスアカウントキーを作成する
GCP PubSub サービスを利用するには、サービスアカウントとサービスアカウントキーを作成する必要があります。
GCP アカウントでサービスアカウントを作成します。サービスアカウントには、対象トピックに対する Pub/Sub Editor 権限が付与されていることを確認してください。
作成したサービスアカウントのメールアドレスをクリックし、「キー」タブを開きます。「キーを追加」ドロップダウンリストから 新しいキーを作成 を選択し、そのアカウント用のサービスアカウントキーを JSON 形式で作成・ダウンロードします。
GCP でトピックを作成・管理する
EMQX で GCP Pub/Sub データ統合を設定する前に、トピックを作成し、基本的な管理操作に慣れておく必要があります。
Google Cloud コンソールで、Pub/Sub -> トピック ページに移動します。詳細な手順は トピックの作成と管理 を参照してください。
TIP
サービスアカウントには、そのトピックに対してパブリッシュ権限が必要です。
トピック ID フィールドにトピックの ID を入力し、トピックを作成 をクリックします。
サブスクリプション ページに移動し、リストの中から作成したトピック ID をクリックします。トピックに対するサブスクリプションを作成します。
- 配信タイプは Pull を選択します。
- メッセージ保持期間は 7 日を選択します。
詳細は GCP Pub/Sub サブスクリプション を参照してください。
サブスクリプション ID -> メッセージ -> Pull をクリックすると、トピックに送信されたメッセージを確認できます。
MQTT データを GCP Pub/Sub にストリームする
このセクションでは、EMQX Cloud デプロイメントから GCP Pub/Sub へ MQTT メッセージを転送するプロデューサー(Sink)コネクターの作成方法を示します。内容はコネクター作成、ルール作成、ルールのテストを含みます。
コネクターを作成する
データ統合ルールを作成する前に、Google PubSub コネクターを作成してサーバーにアクセスできるようにします。
- デプロイメントに移動し、左側ナビゲーションメニューから データ統合 をクリックします。
- 初めてコネクターを作成する場合は、データ転送 カテゴリの下にある Google PubSub を選択します。既にコネクターを作成済みの場合は、新しいコネクター を選択し、続いて データ転送 カテゴリの下の Google PubSub を選択します。
- 新しいコネクター ページで以下を設定します。
- コネクター名: システムが自動的に生成します。
- GCP サービスアカウント認証情報: GCP でサービスアカウントキーを作成するでエクスポートしたサービスアカウント認証情報の JSON 全文を貼り付けるか、ファイル選択 をクリックして JSON ファイルをインポートします。
- その他の設定はデフォルトのままか、ビジネスニーズに応じて設定してください。
- テスト をクリックして接続を検証します。Google PubSub サービスにアクセス可能であれば、成功メッセージが表示されます。
- 新規作成 をクリックしてコネクターの設定を完了します。作成成功ダイアログが表示され、ルールを今すぐ作成するか尋ねられます。新しいルール をクリックするとルール作成画面に進み、コネクターに戻る をクリックすると後でルールを作成できます。
ルールを作成する
次に、書き込むデータを指定するルールを作成し、処理済みデータを GCP PubSub に転送するアクションをルールに追加します。
前のステップで 新しいルール をクリックした場合は、ルール編集ページが自動的に開きます。そうでない場合は、コネクターの 操作 列にある 新しいルール アイコンをクリックするか、ルール セクションの 新しいルール をクリックします。
SQL エディターにルールマッチング用の SQL 文を入力します。以下の例は、
temp_hum/emqxトピックに送信されたメッセージから報告時刻up_timestamp、クライアント ID、メッセージ本文(ペイロード)を読み取り、温度と湿度を抽出します。sqlSELECT timestamp as up_timestamp, clientid as client_id, payload.temp as temp, payload.hum as hum FROM "temp_hum/emqx"Try It Out を使ってデータ入力をシミュレートし、結果をテストできます。
次へ をクリックしてアクションを追加します。
コネクター ドロップダウンボックスから先ほど作成したコネクターを選択します。
以下の情報を設定します。
アクション名: システムが自動生成するか、任意で命名可能です。
GCP PubSub トピック: 「GCP でトピックを作成・管理する」で作成したトピック ID xxx を入力します。
ペイロードテンプレート: 空欄のままにするかテンプレートを定義します。
- 空欄の場合、MQTT メッセージの可視入力(clientid、topic、payload など)を JSON 形式でエンコードします。
- 定義したテンプレートを使う場合、
${variable_name}形式のプレースホルダーが MQTT コンテキストの対応する値で置換されます。例えば、${topic}は MQTT メッセージのトピックがmy/topicならそれに置換されます。
本例では、以下の GCP Pub/Sub トピックとメッセージテンプレートを使用できます。
text# GCP Pub/Sub メッセージテンプレート {"up_timestamp": ${up_timestamp}, "client_id": ${client_id}, "temp": ${temp}, "hum": ${hum}}属性テンプレートおよびオーダーキー(Ordering Key)テンプレート(任意): 同様に、送信メッセージの属性やオーダーキーのフォーマット用テンプレートを定義可能です。
- 属性では、キーと値の両方に
${variable_name}形式のプレースホルダーを使用でき、MQTT コンテキストから値を抽出します。キーのテンプレートが空文字列に解決された場合、そのキーは GCP PubSub 送信メッセージから省略されます。 - オーダーキーでは
${variable_name}形式のプレースホルダーを使用可能で、空文字列に解決された場合は GCP PubSub の送信メッセージにorderingKeyフィールドが設定されません。
- 属性では、キーと値の両方に
高度な設定(任意): その他の設定はデフォルトのままか、ビジネスニーズに応じて設定してください。
確定 をクリックしてルール作成を完了します。
新規ルール作成成功 ポップアップで ルールに戻る をクリックして終了します。
作成成功後、ルール リストに新しいルールが表示されます。操作 (Sink) セクションで関連するアクションを確認できます。
ルールのテスト
温湿度データの報告をシミュレートするために MQTTX の使用を推奨しますが、他のクライアントでも構いません。
MQTTX を使ってデプロイメントに接続し、以下のトピックにメッセージを送信します。
トピック:
temp_hum/emqxペイロード:
json{ "temp": "27.5", "hum": "41.8" }
GCP Pub/Sub -> サブスクリプションに移動し、MESSAGES タブをクリックします。メッセージが確認できるはずです。
コンソールで運用データを確認します。ルールリストのルール ID をクリックすると、ルールの統計情報とそのルールに紐づくすべてのアクションの統計が表示されます。
GCP Pub/Sub からメッセージをコンシュームする
このセクションでは、EMQX Cloud デプロイメントが GCP Pub/Sub からメッセージをコンシュームし、設定したデータ統合を通じて MQTT トピックに再パブリッシュする方法を示します。
TIP
GCP PubSub コンシューマーは、EMQX バージョン 5.10.3 以降を実行する Dedicated Flex デプロイメントで利用可能です。
GCP PubSub コンシューマーコネクターを作成する
コンシューマールールを追加する前に、EMQX Cloud デプロイメントと GCP Pub/Sub 間の接続を確立するために GCP PubSub コンシューマーコネクター を作成します。
デプロイメントに移動し、左側ナビゲーションメニューから データ統合 をクリックします。
初めてコネクターを作成する場合は、データ入力 カテゴリの下にある Google PubSub Consumer を探します。既にコネクターを作成済みの場合は、+ 新しいコネクター をクリックし、続いて データ入力 カテゴリの下の Google PubSub Consumer を選択します。

新しいコネクター ページで以下を設定します。
- コネクター名: システムが自動的に生成します。
- GCP サービスアカウント認証情報: GCP でサービスアカウントキーを作成するでエクスポートしたサービスアカウント認証情報の JSON 全文を貼り付けるか、ファイル選択 をクリックして JSON ファイルをインポートします。
- その他の設定はデフォルトのままか、ビジネスニーズに応じて設定してください。
テスト をクリックして接続を検証します。GCP Pub/Sub サービスにアクセス可能であれば、成功メッセージが表示されます。
新規作成 をクリックしてコネクターの設定を完了します。作成成功ダイアログが表示され、ルールを今すぐ作成するか尋ねられます。新しいルール をクリックするとルール作成画面に進み、コネクターに戻る をクリックすると後でルールを作成できます。
ルールを作成する
次に、データソースを指定し、コンシュームしたメッセージを MQTT トピックに転送する出力アクションを追加するルールを作成します。
前のステップで 新しいルール をクリックした場合は、ルール編集ページが自動的に開きます。そうでない場合は、コネクターの 操作 列にある 新しいルール アイコンをクリックするか、ルール セクションの 新しいルール をクリックします。
ルール編集ページで自動的にアクションソース設定パネルが開きます。ソースタイプとして Google PubSub Consumer を選択し、次へ をクリックします。

ソースを設定します。
- コネクター: 先ほど作成した GCP PubSub コンシューマーコネクターを選択します。
- GCP PubSub トピック: コンシューム対象の GCP Pub/Sub トピック名を入力します。例:
my-iot-topic - 一度にプルする最大メッセージ数: 1 回のリクエストでプルする最大メッセージ数を設定します。デフォルト値を使うか、スループットに応じて調整してください。
- その他の設定はデフォルトのままか、ビジネスニーズに応じて設定してください。
確定 をクリックしてソース設定を完了します。
SQL エディターがデータソースに合わせて自動更新されます。必要に応じて
SELECTフィールドを調整できます。例:sqlSELECT * FROM "$bridges/gcp_pubsub_consumer:<source-name>"GCP Pub/Sub メッセージから利用可能なフィールドは以下の通りです。
フィールド名 説明 message_idGCP Pub/Sub によって割り当てられたメッセージ ID。 publish_timeメッセージがパブリッシュされたタイムスタンプ。 topicメッセージが読み取られた GCP Pub/Sub トピック。 valueメッセージのペイロード。 attributesメッセージに付与されたキー・バリュー属性。 ordering_key設定されている場合のメッセージのオーダーキー。 次へ をクリックして出力アクションを追加します。
アクションタイプとして 再パブリッシュ (Republish) を選択し、以下を設定します。
- トピック: パブリッシュ先の MQTT トピック。例:
gcp/messages。${topic}のようなプレースホルダーを使って動的にトピックを決定可能です。 - QoS:
0、1、2のいずれかを明示的に選択します。動的に QoS を設定したい場合は、ルール SQL にqosフィールドを追加し、アクションで参照してください。 - Retain:
trueまたはfalseを明示的に選択します。動的に Retain を設定したい場合は、ルール SQL に対応するフィールドを追加し、アクションで参照してください。 - メッセージテンプレート: 空欄にするとルール出力のすべてのフィールドを転送します。
${.}と入力するとすべてのフィールドを含め、${value}と入力するとメッセージペイロードのみを転送します。
- トピック: パブリッシュ先の MQTT トピック。例:
確定 をクリックしてルール作成を完了します。
新規ルール作成成功 ポップアップで ルールに戻る をクリックして終了します。
作成成功後、ルール リストに新しいルールが表示されます。操作 (Source) セクションで関連するアクションを確認できます。再パブリッシュアクションはルール編集ボタンをクリックして確認可能です。

ルールのテスト
GCP Pub/Sub トピックにメッセージをパブリッシュし、デプロイメント内の転送先 MQTT トピックをサブスクライブしてルールを検証できます。
MQTTX(または任意の MQTT クライアント)を使い、再パブリッシュアクションで設定した MQTT トピック(例:
gcp/messages)をサブスクライブします。GCP コンソールまたは
gcloudCLI を使って GCP Pub/Sub トピックにメッセージをパブリッシュします。bashgcloud pubsub topics publish my-iot-topic --message='{"temp": 27.5, "hum": 41.8}'MQTT クライアントがサブスクライブしたトピックでメッセージを受信することを確認します。
コンソールでルール統計を確認します。ルール リストのルール ID をクリックすると、ルール実行統計とそのルールに関連付けられたすべてのアクションの統計が表示されます。