Skip to content

BigQueryへのMQTTデータ取り込み

BigQueryは、大量のリレーショナル構造化データ向けのエンタープライズデータウェアハウスです。大規模かつアドホックなSQLベースの分析とレポーティングに最適化されており、組織の洞察を得るのに最適です。EMQXは、MQTTデータのリアルタイム抽出、処理、分析のためにBigQueryとのシームレスな統合をサポートしています。

本ページでは、EMQXとBigQuery間のデータ統合について包括的に解説し、データ統合の作成と検証に関する実践的な手順を提供します。

動作の仕組み

BigQueryデータ統合は、EMQXの標準機能として提供されており、ユーザーがMQTTデータストリームをGoogle Cloudとシームレスに統合し、IoTアプリケーション開発における豊富なサービスと機能を活用できるよう設計されています。

bigquery_architecture

EMQXはルールエンジンとSinkを介してMQTTデータをBigQueryに転送します。全体の流れは以下の通りです。

  1. IoTデバイスがメッセージをパブリッシュ:デバイスは特定のトピックを通じてテレメトリやステータスデータをパブリッシュし、ルールエンジンをトリガーします。
  2. ルールエンジンがメッセージを処理:組み込みのルールエンジンを用いて、特定の送信元からのMQTTメッセージをトピックマッチングに基づいて処理します。ルールエンジンは対応するルールをマッチングし、データ形式の変換、特定情報のフィルタリング、コンテキスト情報の付加などの処理を行います。
  3. BigQueryへのブリッジング:ルールがメッセージのBigQuery転送アクションをトリガーし、データプロパティ、オーダーキー、MQTTトピックとBigQueryトピックのマッピングを簡単に設定できます。これにより、データ統合におけるより豊かなコンテキスト情報と順序保証が提供され、柔軟なIoTデータ処理が可能になります。

特長と利点

EMQXとBigQueryの統合は、MQTTデータのための堅牢でスケーラブルかつリアルタイムなデータパイプラインを提供します。以下の特長と利点により、IoT分析やデータ駆動型の意思決定を簡素化します。

  • リアルタイムデータ取り込み:低レイテンシでEMQXからBigQueryへMQTTメッセージをシームレスにストリームします。即時処理と分析が必要なタイムセンシティブなアプリケーションに対応します。
  • 柔軟なデータマッピング:MQTTトピックやメッセージペイロードをBigQueryのテーブルやフィールドにカスタマイズしてマッピング可能です。
  • スケーラブルでサーバレスな分析:BigQueryのフルマネージドかつサーバレスなアーキテクチャを活用し、大規模なIoTデータの分析を実現します。
  • Google Cloudエコシステムとの簡単な統合:Data Studio、Looker、AI PlatformなどGoogle Cloudのネイティブサービスと連携し、可視化や機械学習を簡単に実現。データ収集から洞察生成までのエンドツーエンドパイプライン構築を容易にします。

はじめる前に

このセクションでは、BigQueryデータ統合を作成する前に必要な準備について説明します。

前提条件

GCPでのサービスアカウントキーの作成

Service Account JSON認証を使用する場合は、Google Cloudでサービスアカウントを作成し、JSON形式のキーを生成してください。

  1. GCPアカウントでサービスアカウントを作成します。サービスアカウントには、使用するデータセットやテーブルへのアクセス権限が必要です。例えば、「BigQuery Data Editor」ロールを付与して必要なデータセットやテーブルの読み書きを許可するか、少なくともデータへの読み書き権限を持たせてください。

  2. 作成したサービスアカウントのメールアドレスをクリックします。

  3. Keyタブをクリックし、Add keyのドロップダウンからCreate new keyを選択してサービスアカウントキーを作成し、JSON形式でダウンロードします。

    TIP

    ダウンロードしたサービスアカウントキーは後でEMQXのBigQuery認証に使用するため、安全に保管してください。

    サービスアカウントキー

GCPでのWorkload Identity Federationの設定

Workload Identity Federation(WIF)を使うと、長期間有効なサービスアカウントキーを使わずにEMQXがGCPリソースにアクセスできます。EMQXは外部IDプロバイダー(例:Microsoft Azure)からのトークンをGCPのSecurity Token Service経由で一時的なGCPトークンに交換し、そのトークンでGCPサービスアカウントをなりすまします。トークンの更新は自動で行われます。

WIFを利用するには、コネクター作成前にGCPプロジェクトで以下を完了してください。

  1. Google CloudコンソールでIAM & Admin -> Workload Identity Federationに移動し、ワークロードアイデンティティプールを作成し、Pool IDProject Numberを控えます。

  2. プールにプロバイダーを追加し、Provider IDを控えます。OIDC認証の場合は、外部IDプロバイダーからOAuth 2.0クライアント認証情報(クライアントID、クライアントシークレット、トークンエンドポイントURI)を取得してください。

  3. ワークロードアイデンティティプールに、BigQueryデータセットやテーブルにアクセス可能なGCPサービスアカウントをなりすます権限を付与します。コネクター設定時にサービスアカウントのメールアドレスが必要です。

    TIP

    詳細な手順はWorkload Identity Federationの設定を参照してください。

例:Microsoft Azure(Entra ID)

Microsoft Entra IDでAPIを公開するアプリケーションを登録し、クライアントシークレットを作成します。コネクター設定時に以下の値を使用します。

コネクター項目
Endpoint URIhttps://login.microsoftonline.com/<tenant-id>/oauth2/v2.0/token
OAuth Client IDアプリケーション(クライアント)ID、形式は api://<application-id>
OAuth Client Secretアプリケーション用に生成したクライアントシークレット
OAuth Request Scopeapi://<application-id>/.default

補足

scopeはアプリケーションのaudience(aud)と完全に一致させる必要があります。そうしないとGCP STSとのトークン交換が失敗します。詳細はMicrosoftのOAuth 2.0クライアント認証フローを参照してください。

サービスアカウントにWIFプールへのアクセス権を付与する際は、Subject値にアプリケーションIDではなくObject IDを使用してください。Object IDはAzureポータルのアプリケーション概要ページのEnterprise applicationsに表示されます。

Attached Service Accountの前提条件

Attached Service Account認証を使用する場合、EMQXはGCP Compute Engineインスタンス上で実行され、そのインスタンスにサービスアカウントがアタッチされている必要があります。インスタンスのOAuthアクセススコープがBigQueryへのアクセスを許可していることを確認してください。Googleはcloud-platformスコープ(https://www.googleapis.com/auth/cloud-platform)の使用を推奨し、サービスアカウントの権限はIAMロールで制限することを推奨しています。サービスアカウントは対象のBigQueryデータセットおよびテーブルへのアクセス権限を持っている必要があります。詳細はGoogle Cloudのサービスアカウントを参照してください。

対象のBigQueryデータセットとテーブルは、Compute Engineインスタンスに関連付けられたGCPプロジェクト内に存在する必要があります。EMQXクラスターの場合、すべてのノードがこれらの要件を満たし、同じプロジェクトのCompute Engineインスタンス上で実行されている必要があります。

コネクター起動時に、EMQXは自動的にインスタンスメタデータエンドポイントからGCPプロジェクトIDとアクセストークンを取得します。サービスアカウントキーのアップロードは不要です。

GCPでのデータセットおよびテーブルの作成と管理

EMQXでBigQueryデータ統合を設定する前に、GCPで必要なデータセットとテーブルを作成してください。

  1. Google CloudコンソールでBigQuery -> Studioページに移動します。詳細な手順はデータのロードとクエリクイックスタートガイドを参照してください。

    TIP

    使用するサービスアカウントは、対象テーブルに対する書き込み権限を持っている必要があります。

  2. Explorerペインでケバブメニュー(⋮)をクリックし、Create datasetを選択します。データセット名を定義し、Create datasetをクリックします。

  3. データセット作成後、Explorerペインでデータセットを選択し、(+) Create tableをクリックします。

    • ソースは「Empty Table」を選択します。

    • テーブル名を入力します。

    • テーブルスキーマを定義します。例えば、Edit as textトグルをクリックし、以下のスキーマ定義をテキストフィールドに貼り付けます。

      clientid:string,payload:bytes,topic:string,publish_received_at:timestamp
    • Create tableをクリックして設定を完了します。

  4. EMQXが書き込み可能なように権限を設定します。

    • データセットを選択し、Shareをクリックします。

    • サービスアカウントのメールアドレスをプリンシパルとして追加します。

    • 以下のような適切なロールを割り当てます。

      • データセットに対して「BigQuery Data Viewer」(読み取りアクセス)

      • テーブルに対して「Editor」(読み書きアクセス)

  5. テーブル作成後、クエリを実行して確認できます。

    • テーブルをクリックし、Queryをクリックします。

    • 以下のようなSQL文を実行してテーブルにアクセスできることを確認します。

      sql
      SELECT * FROM `my_project.my_dataset.my_tab` LIMIT 1000

BigQueryコネクターの作成

BigQuery Producer Sinkアクションを追加する前に、EMQXとBigQuery間の接続を確立するためにBigQueryコネクターを作成する必要があります。

  1. EMQXダッシュボードでIntegration -> Connectorをクリックします。

  2. ページ右上のCreateをクリックし、コネクター選択ページでBigQueryを選択してNextをクリックします。

  3. 名前と説明を入力します(例:my_bigquery)。名前はBigQuery Sinkとコネクターを紐付けるために使用され、クラスター内で一意である必要があります。

  4. Authenticationリストから以下のいずれかの認証方法を選択し、対応するフィールドを設定します。

    • Service Account JSON:前述のサービスアカウントキーの作成でエクスポートしたJSON形式のサービスアカウント認証情報をアップロードします。

    • Workload Identity Federation (WIF):以下のフィールドを入力します。この方法はサービスアカウントJSONファイルを使用しません。前提条件はWorkload Identity Federationの設定を参照してください。

      • GCP Project ID:コネクターがアクセスするリソースのプロジェクトID。

      • GCP Project Number:コネクターがアクセスするリソースのプロジェクト番号。

      • Service Account Email:なりすますサービスアカウントのメールアドレス。

      • Workload Identity Pool ID:WIFトークン交換で使用するワークロードアイデンティティプールのID。

      • Workload Identity Provider ID:WIFトークン交換で使用するワークロードアイデンティティプロバイダーのID。

      • Initial Token Configuration:認証情報タイプを選択し、対応するフィールドを入力します。現在はOIDC with Client Credentials Grant Typeのみサポートされています。

        • Endpoint URI:OIDCプロバイダーのOAuthトークンエンドポイントURI。

        • OAuth Client ID:OAuthサーバーにトークンを要求するためのクライアントID。

        • OAuth Client Secret:OAuthサーバーにトークンを要求するためのクライアントシークレット。

        • OAuth Request Scope:OAuthアクセストークン要求時に必要な場合のscope

    • Attached Service Account:追加のフィールドは不要です。EMQXはインスタンスメタデータエンドポイントから自動的にGCPプロジェクトIDとアクセストークンを取得します。前提条件はAttached Service Accountの前提条件を参照してください。

  5. Createをクリックする前に、Test ConnectivityをクリックしてコネクターがBigQueryサーバーに接続できるかテストできます。

  6. ページ下部のCreateボタンをクリックしてコネクターの作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてBigQueryに転送するデータを指定するルールを作成できます。詳細はBigQuery Sink付きルールの作成を参照してください。

BigQuery Sink付きルールの作成

このセクションでは、BigQueryに保存するデータを指定するルールの作成方法を示します。

  1. EMQXダッシュボードでIntegration -> Rulesをクリックします。

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

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

  4. SQL Editorでルールを設定します。例えば、トピックt/bqのMQTTメッセージをBigQueryに保存したい場合、以下のSQLを使用できます。

    注意:独自のSQLを指定する場合は、Sinkのペイロードテンプレートで必要なすべてのフィールドがSELECT句に含まれていることを確認してください。

    sql
    SELECT
      clientid,
      topic,
      base64_encode(payload) AS payload,
      timestamp/1000 AS publish_received_at
    FROM
      "t/bq"

    注意

    BigQueryテーブルのカラムであるフィールドのみを選択してください。そうしないとBigQueryは不明なフィールドとして認識しません。

    TIP

    初心者の方はSQL Examplesをクリックし、Enable Testを有効にしてSQLルールの学習とテストが可能です。

  5. Add Actionボタンをクリックして、ルールでトリガーされるアクションを定義します。Type of ActionドロップダウンリストからBigQueryを選択し、ルールで処理されたデータをBigQueryに送信するようにします。

  6. ActionドロップダウンはCreate Actionのままにするか、既存のBigQuery Sinkを選択できます。本デモでは新しいSinkを作成してルールに追加します。

  7. NameフィールドにSinkの名前を入力します。名前は英数字の組み合わせにしてください。

  8. Connectorドロップダウンから先ほど作成したmy_bigqueryを選択します。ドロップダウン横のボタンから新しいコネクターを作成することも可能です。設定パラメータの詳細はコネクターの作成を参照してください。

  9. 以下のBigQueryリソースパラメータを設定します。

    • Project ID(任意):対象のデータセットとテーブルが存在するGCPプロジェクトのIDを入力します。指定すると、選択したコネクターの認証設定から抽出されたプロジェクトIDを上書きし、このSinkにのみ適用されます。空欄の場合は認証設定から取得したプロジェクトIDが使用されます。

    • DatasetおよびTableGCPでのデータセットおよびテーブルの作成と管理で作成したデータセット名とテーブル名をそれぞれ入力します。

  10. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はフォールバックアクションを参照してください。

  11. 詳細設定(任意):必要に応じて詳細設定オプションを構成します。詳細は詳細設定を参照してください。

  12. Createをクリックする前に、Test ConnectivityをクリックしてコネクターがBigQueryサーバーに接続可能かテストできます。

  13. CreateボタンをクリックしてSinkの設定を完了すると、新しいSinkがAction Outputsタブに表示されます。

  14. Create Ruleページに戻り、Createをクリックしてルールを作成します。

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

また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/bqのメッセージがルールmy_ruleで解析されBigQueryに送信・保存されていることが確認できます。

ルールのテスト

  1. MQTTクライアントMQTTXを使ってトピックt/bqにメッセージを送信します。

    bash
    mqttx pub -i emqx_c -t t/bq -m '{ "msg": "hello BigQuery" }'
  2. Sinkの稼働状況を確認し、新規の受信メッセージと送信メッセージが1件ずつあることを確認します。

  3. GCPのBigQuery -> Studioに移動し、テーブルをクリックしてQueryをクリックします。クエリを実行するとメッセージが確認できます。

詳細設定

このセクションでは、BigQuery Producer Sinkの詳細設定オプションについて説明します。ダッシュボードでSinkを設定する際に、Advanced Settingsを展開して以下のパラメータをニーズに応じて調整できます。

フィールド名説明デフォルト値
Buffer Pool SizeEMQXとBigQuery間のデータフローを管理するバッファワーカープロセスの数を指定します。これらのワーカーはデータを一時的に格納・処理し、ターゲットサービスへの送信を最適化しスムーズなデータ伝送を保証します。16
Request TTLバッファに入ったリクエストが有効とみなされる最大期間(秒)を指定します。このタイマーはリクエストがバッファに入った時点からカウントされ、TTLを超えてバッファに残るか、BigQueryからの応答やアックがタイムリーに得られない場合、リクエストは期限切れとみなされます。45
Health Check IntervalSinkがBigQueryとの接続の自動ヘルスチェックを行う間隔(秒)を指定します。15
Health Check Interval Jitter複数のノードが同時にヘルスチェックを開始するのを防ぐために、基本のヘルスチェック間隔に加える一様ランダム遅延です。複数のアクションやソースが同じコネクターを共有する場合、ジッターを有効にするとヘルスチェックが少しずつ異なるタイミングで実行されます。0ミリ秒
Health Check TimeoutコネクターがBigQueryとの接続の自動ヘルスチェックを行う際のタイムアウト時間(秒)を指定します。60
Max Buffer Queue SizeBigQuery Sinkの各バッファワーカーがバッファリングできる最大バイト数を指定します。バッファワーカーはデータを一時的に格納し、BigQueryへの送信を効率化します。システム性能やデータ伝送要件に応じて調整してください。256
Query Modesynchronousまたはasynchronousのリクエストモードを選択し、メッセージ送信を最適化します。非同期モードではBigQueryへの書き込みがMQTTメッセージのパブリッシュ処理をブロックしません。ただし、クライアントがBigQuery到達前にメッセージを受信する可能性があります。Async
Batch SizeEMQXからBigQueryへ一度に送信するデータバッチの最大サイズを指定します。サイズを調整することでデータ転送の効率と性能を最適化できます。Batch Size1の場合は、データレコードがバッチ化されず個別に送信されます。1000
Inflight Window「インフライトキューリクエスト」とは、送信済みだがまだ応答やアックを受け取っていないリクエストのことです。この設定はSinkがBigQueryと通信中に同時に存在可能なインフライトキューリクエストの最大数を制御します。Request Modeasynchronousの場合に特に重要です。同一MQTTクライアントからのメッセージを厳密に順序処理する必要がある場合は、この値を1に設定してください。100