Bigtable に MQTT データを取り込む
Cloud Bigtable は、Google Cloud 上のフルマネージドのワイドカラム型 NoSQL データベースサービスです。時系列データ、テレメトリストレージ、イベントレコード、高スループットの IoT データ取り込みなど、大規模かつ低レイテンシのワークロード向けに設計されています。
EMQX はルールエンジンと Bigtable Sink を通じて Bigtable との連携をサポートしています。MQTT メッセージをルール SQL で処理し、選択したフィールドを Bigtable の行キーやセルのミューテーションにマッピングして、処理済みデータをリアルタイムで Bigtable テーブルに書き込みます。
このページでは、Bigtable データ統合の仕組みを紹介し、EMQX ダッシュボードでの連携作成およびテストのワークフローを説明します。
仕組み
Bigtable データ統合は EMQX の標準機能です。MQTT データを Google Cloud にストリームし、デバイスのテレメトリやイベントデータを Bigtable に保存して、後のクエリ、分析、下流処理に活用できます。

EMQX はルールエンジンと Sink を介して MQTT データを Bigtable に転送します。全体の流れは以下の通りです。
- IoT デバイスがメッセージをパブリッシュ: デバイスがテレメトリ、状態、イベントデータを MQTT トピックにパブリッシュします。
- ルールエンジンがメッセージを処理: ルールエンジンがトピックで MQTT メッセージをマッチングし、SQL で Bigtable に必要なフィールドを抽出または変換します。
- Bigtable への書き込み: Bigtable Sink が各ルール出力レコードを行ミューテーションとして Bigtable テーブルに書き込みます。設定された行キーや
set_cellミューテーションフィールドを使用します。下流のアプリケーションやサービスは保存されたデータをクエリや処理に利用し、低レイテンシアプリケーション、時系列クエリ、分析、AI/ML パイプラインに活用できます。
特長とメリット
EMQX と Bigtable の統合により、以下のメリットがあります。
- 高スループットの IoT データ取り込み: 大規模なテレメトリやイベントワークロードに対して MQTT メッセージを Bigtable に書き込み可能。
- 柔軟なフィールドマッピング: ルール SQL で Bigtable の行キー、カラムファミリー、カラム修飾子、タイムスタンプ、セル値として使用するフィールドを明示的に選択・エイリアス可能。
- バッチおよび非同期書き込み: バッチモードや非同期リクエストモードを利用して書き込みスループットを向上させ、MQTT メッセージのパブリッシュへの影響を軽減。
- Google Cloud 連携: MQTT データを Bigtable に保存し、他の Google Cloud サービスと連携して分析、処理、アプリケーション開発に活用可能。
はじめる前に
Bigtable データ統合を作成する前に必要な準備について説明します。
前提条件
- EMQX データ統合のルールに関する知識
- データ統合に関する知識
- Bigtable が有効な Google Cloud プロジェクト
- Bigtable インスタンス、テーブル、および少なくとも1つのカラムファミリー
- 使用予定の認証方式に必要な認証情報:
- サービスアカウント JSON: サービスアカウントキーの JSON ファイル
- Workload Identity Federation (WIF): ワークロードアイデンティティプール、プロバイダー、プロジェクト ID、プロジェクト番号、サービスアカウントメール、外部 ID プロバイダーからの OAuth 2.0 クライアント認証情報
- Attached Service Account: GCP Compute Engine 上で動作する EMQX デプロイメントで、Attached Service Account の前提条件を満たすもの
GCP でサービスアカウントキーを作成する
サービスアカウント JSON 認証を使用する場合は、Google Cloud でサービスアカウントを作成し、JSON 形式のキーを生成します。
GCP アカウントでサービスアカウントを作成します。
サービスアカウントに Bigtable インスタンスおよびテーブルへの書き込み権限を付与します。例えば、対象テーブルのデータ読み書き操作を許可する Bigtable ロールを割り当てます。
作成したサービスアカウントのメールアドレスをクリックします。
Keys タブをクリックし、Add key のドロップダウンリストから Create new key を選択し、JSON 形式のキーをダウンロードします。
TIP
サービスアカウントキーは安全に保管してください。Bigtable コネクター作成時に使用します。
GCP で Workload Identity Federation を設定する
Workload Identity Federation (WIF) を使うと、長期間有効なサービスアカウントキーを使わずに EMQX が GCP リソースにアクセスできます。EMQX は Microsoft Azure などの外部 ID プロバイダーからトークンを取得し、GCP Security Token Service を通じて一時的な GCP トークンに交換し、そのトークンで GCP サービスアカウントをなりすまします。トークンの更新は自動で行われます。
WIF を使うには、コネクター作成前に GCP プロジェクトで以下を完了してください。
- Google Cloud コンソールで IAM & Admin -> Workload Identity Federation に移動し、ワークロードアイデンティティプールを作成し、Pool ID と Project Number を控えます。
- プールにプロバイダーを追加し、Provider ID を控えます。OIDC ベースの認証の場合、外部 ID プロバイダーから OAuth 2.0 クライアント認証情報(クライアント ID、クライアントシークレット、トークンエンドポイント URI、リクエストスコープ)を取得します。
- ワークロードアイデンティティプールに Bigtable インスタンスとテーブルにアクセスできる GCP サービスアカウントのなりすまし権限を付与し、サービスアカウントメールを控えます。
- Bigtable リソースを含むプロジェクトの GCP プロジェクト ID を控えます。
TIP
詳細は Workload Identity Federation の設定 を参照してください。
例: Microsoft Azure (Entra ID)
Microsoft Entra ID で API を公開するアプリケーションを登録し、クライアントシークレットを作成します。コネクター設定時に以下の値を使用します。
| コネクター項目 | 値 |
|---|---|
| OAuth Token Endpoint URI | https://login.microsoftonline.com/<tenant-id>/oauth2/v2.0/token |
| OAuth Client ID | api://<application-id> 形式のアプリケーション(クライアント)ID |
| OAuth Client Secret | アプリケーション用に生成したクライアントシークレット |
| OAuth Request Scope | api://<application-id>/.default |
::: note
OAuth Request Scope はアプリケーションのオーディエンス(aud)と完全に一致する必要があります。そうしないと GCP STS とのトークン交換に失敗します。WIF プールにサービスアカウントアクセスを付与する際は、アプリケーション ID ではなく Object ID をサブジェクト識別子として使用してください。Object ID は Azure ポータルのアプリケーション概要ページの Enterprise applications で確認できます。
:::
Attached Service Account の前提条件
Attached Service Account 認証を使う場合、EMQX は GCP Compute Engine インスタンス上で動作し、サービスアカウントがアタッチされている必要があります。インスタンスの OAuth アクセススコープが Bigtable へのアクセスを許可していることを確認してください。Google は cloud-platform スコープ(https://www.googleapis.com/auth/cloud-platform)の使用を推奨し、サービスアカウントの権限は IAM ロールで制限することを推奨しています。サービスアカウントには対象の Bigtable インスタンスおよびテーブルへのアクセス権限が必要です。詳細は Google Cloud ドキュメントのサービスアカウントを参照してください。
対象の Bigtable インスタンスおよびテーブルは Compute Engine インスタンスに関連付けられた GCP プロジェクト内に存在する必要があります。EMQX クラスターの場合、すべてのノードがこれらの要件を満たし、そのプロジェクトの Compute Engine インスタンス上で動作している必要があります。
コネクター起動時に EMQX はインスタンスメタデータエンドポイントから GCP プロジェクト ID とアクセストークンを自動取得します。サービスアカウントキーのアップロードは不要です。
GCP で Bigtable リソースを作成・管理する
EMQX で Bigtable データ統合を設定する前に、Google Cloud で対象の Bigtable リソースを作成してください。
Google Cloud コンソールで Bigtable ページに移動します。
Bigtable インスタンスを作成または選択します。インスタンス作成時の Instance name は Google Cloud コンソール上の表示名としてのみ使用されます。
EMQX MQTT Messagesのような読みやすい名前を入力可能です。Instance ID は後で EMQX で使用する値で、emqxinstのようなシンプルで一意な識別子にしてください。テーブルを作成し、テーブル ID(例:
mqtt_messages)を控えます。テーブルに少なくとも1つのカラムファミリーを作成します(例:
cf)。TIP
EMQX は Google Cloud コンソールのインスタンス表示名や
projects/<project-id>/instances/<instance-id>のような完全修飾リソース名ではなく、Instance ID と Table ID を使用します。
Bigtable コネクターを作成する
Bigtable Sink アクションを追加する前に、EMQX と Bigtable 間の接続を確立するための Bigtable コネクターを作成します。
- EMQX ダッシュボードで Integration -> Connectors に移動します。
- 画面右上の Create をクリックし、Bigtable を選択して Next をクリックします。
- コネクター名と説明を入力します(例:
my_bigtable)。名前は Bigtable Sink とコネクターを関連付けるために使用され、クラスター内で一意である必要があります。 - 認証オプションを設定します。
- Authentication: EMQX が GCP に認証する方法を選択します。
- Service Account JSON: GCP でサービスアカウントキーを作成するでエクスポートした JSON ファイルを GCP Service Account Credentials にアップロードします。Select file をクリックして JSON ファイルを選択可能です。
- Workload Identity Federation (WIF): 以下の項目を入力します。この方法ではサービスアカウント JSON ファイルは不要です。前提条件は GCP で Workload Identity Federation を設定するを参照してください。
- GCP Project ID: コネクターがアクセスするリソースの GCP プロジェクト ID
- GCP Project Number: コネクターがアクセスするリソースの GCP プロジェクト番号
- Service Account Email: なりすますサービスアカウントのメールアドレス
- Workload Identity Pool ID: WIF トークン交換に使用するワークロードアイデンティティプール ID
- Workload Identity Provider ID: WIF トークン交換に使用するワークロードアイデンティティプロバイダー ID
- Credential Type: 外部 ID プロバイダーが使用する認証情報タイプ。現在は OIDC クライアント認証情報をサポート。選択後、以下を入力します。
- OAuth Client ID: OAuth サーバーからトークンを要求するためのクライアント ID
- OAuth Client Secret: OAuth サーバーからトークンを要求するためのクライアントシークレット
- OAuth Token Endpoint URI: OIDC プロバイダーの OAuth トークンエンドポイント URI
- OAuth Request Scope: OAuth サーバーからアクセストークンを要求する際の
scope。プロバイダーが必要な場合に入力。 - OAuth Request Audience: OAuth サーバーからアクセストークンを要求する際の
audience。プロバイダーが必要な場合に入力。
- Attached Service Account: 追加の項目は不要です。EMQX がインスタンスメタデータエンドポイントから GCP プロジェクト ID とアクセストークンを自動取得します。前提条件は Attached Service Account の前提条件を参照してください。
- Enable TLS: デプロイ環境で TLS が必要な場合は有効にします。
- Advanced Settings: 詳細接続オプションを設定する場合はこのセクションを展開します。
- Authentication: EMQX が GCP に認証する方法を選択します。
- Create をクリックする前に、Test Connectivity をクリックして EMQX が Bigtable に接続できるか検証可能です。
- Create ボタンをクリックしてコネクターの作成を完了します。作成成功ダイアログが表示され、ルールを今すぐ作成するか尋ねられます。Create Rule をクリックするとコネクターが事前選択された状態でルール作成画面に進みます。Back To Connector List をクリックすると戻って後でルールを作成できます。
Bigtable Sink を使ったルールの作成
このセクションでは、MQTT メッセージを Bigtable に書き込むルールの作成方法を示します。
前のステップで Create Rule をクリックした場合、Add Action パネルが自動で開き、Type of Action が
Bigtable、コネクターが事前選択されています。まずアクションを設定するためにステップ5に進み、アクション作成後にルールページに戻ってルール ID と SQL を設定してください。そうでない場合は、ダッシュボードの Integration -> Rules ページに移動し、右上の Create をクリックします。ルール ID に
my_ruleを入力します。SQL Editor にルール SQL を入力します。Bigtable Sink は Sink で設定したフィールド名を使ってルール出力から値を参照するため、SQL では Bigtable ミューテーションに必要なすべてのフィールドを明示的に選択し、エイリアスを付ける必要があります。
例:
sqlSELECT clientid AS rk, 'cf' AS fn, '' AS cq, payload AS v, publish_received_at * 1000 AS tm FROM "t/bigtable"この例では、
rkが Bigtable の行キーとして使われます。fnがカラムファミリー名として使われます。cqがカラム修飾子として使われます。tmがマイクロ秒単位のタイムスタンプとして使われます。vがセル値として使われます。
TIP
Bigtable Sink のフィールドはルール出力フィールドを指すキー名であり、テンプレート式ではありません。SQL で必要なキーが選択されていない場合、そのメッセージの Bigtable ミューテーションを構築できません。
Add Action をクリックし、Add Action パネルで Type of Action ドロップダウンから
Bigtableを選択します。Action は
Create Actionのままにするか、既存の Bigtable Sink を選択します。コネクター成功ダイアログからルール作成した場合は、Type of Action がすでにBigtable、コネクターが事前選択されていることを確認してください。Name に Sink 名を入力します。Description に説明を入力することも可能です。
Connectors で Bigtable コネクターを作成するで作成した Bigtable コネクターを選択します。未選択の場合はこのパネルからプラスアイコンをクリックして新規作成できます。
Bigtable アクションパラメーターを設定します。
フィールド 説明 例 Instance ID Bigtable インスタンス識別子。完全修飾の projects/.../instances/...ではなくシンプルな ID を使用。emqxinstTable ID Bigtable テーブル識別子。完全修飾の projects/.../instances/.../tables/...ではなくシンプルな ID を使用。mqtt_messagesRow Key メッセージの行キーを含むキー名。 rkMutations 受信メッセージに適用するセルミューテーションのリスト。Add をクリックしてミューテーションを追加。 - Mutation Type ミューテーション操作タイプ。現在は Set Cell ミューテーションをサポート。 Set CellColumn Family ミューテーションのカラムファミリーを含むキー名。 fnColumn Qualifier ミューテーションのカラム修飾子を含むキー名。 cqTimestamp (microseconds) ミューテーションのタイムスタンプ(マイクロ秒単位)を含むキー名。 tmValue ミューテーションの値を含むキー名。 vメッセージ配信失敗時の信頼性向上のために Fallback Actions を設定することも可能です。詳細は Fallback Actions を参照してください。
必要に応じて Advanced Settings を設定します。詳細は Advanced Settings を参照してください。
Create をクリックする前に、Test Connectivity をクリックして Sink が Bigtable に接続できるか検証可能です。
Create をクリックして Sink の設定を完了します。
Create Rule ページに戻り、Create をクリックしてルールを作成します。
ルールのテスト
MQTTX を使ってトピック
t/bigtableにメッセージをパブリッシュします。bashmqttx pub -i emqx_c -t t/bigtable -m '{ "msg": "hello Bigtable" }'ルールと Sink のメトリクスを確認します。マッチ数と成功数が増加しているはずです。
Google Cloud で対象の Bigtable テーブルをクエリし、以下の内容で行が書き込まれていることを確認します。
- 行キー: MQTT クライアント ID(例:
emqx_c) - カラムファミリー:
cf - カラム修飾子: 空文字列
- セル値: MQTT ペイロード
- 行キー: MQTT クライアント ID(例:
詳細設定
このセクションでは、Bigtable コネクターおよび Sink の一般的な詳細設定について説明します。
コネクター詳細設定
| フィールド | 説明 | デフォルト値 |
|---|---|---|
| Connection Pool Size | Bigtable 接続プールの接続数。 | 8 |
| Connect Timeout | Bigtable への接続確立タイムアウト。 | 5s |
| Start Timeout | コネクター起動タイムアウト。 | 5s |
| Health Check Interval | Bigtable 接続のヘルスチェック間隔。 | 15s |
| Health Check Timeout | コネクターのヘルスチェックタイムアウト。 | 60s |
Sink 詳細設定
| フィールド | 説明 | デフォルト値 |
|---|---|---|
| Buffer Pool Size | Bigtable へのデータ処理・送信に使うバッファワーカープロセス数。 | 16 |
| Dispatch Strategy | バッファワーカーへのリクエスト振り分け戦略。デフォルトは MQTT クライアント ID ごとに振り分け。 | Per Client ID |
| Request TTL | バッファに入ってから有効な最大時間。期限切れになると送信やアック前でも期限切れとみなす。 | 45s |
| Health Check Interval | Bigtable 接続のヘルスチェック間隔。 | 15s |
| Health Check Interval Jitter | ヘルスチェック間隔に加えるランダムジッター。 | 0ms |
| Health Check Timeout | コネクターのヘルスチェックタイムアウト。 | 60s |
| Max Buffer Queue Size | 各バッファワーカーの最大バッファキューサイズ。 | 256MB |
| Batch Size | 1 バッチあたりの最大書き込みレコード数。1 に設定するとバッチ処理を無効化。 | 1000 |
| Query Mode | リクエストモード。非同期モードでは Bigtable への書き込みが MQTT メッセージパブリッシュをブロックしない。 | Async |
| Inflight Window | 非同期モードでの最大インフライトリクエスト数。同一 MQTT クライアントからのメッセージに厳密な順序が必要な場合は 1 に設定。 | 100 |
高スループット環境では、Connection Pool Size、Buffer Pool Size、Dispatch Strategy、Batch Size、Inflight Window を想定されるクラスターのワークロードに応じて調整してください。例えば、クラスター全体で 2 分間に約 11,000,000 メッセージ、5,000~10,000 MQTT 接続のワークロードの場合は、本番運用前に代表的なベンチマークで設定を検証してください。