DatabricksへのMQTTデータ取り込み
Databricksは、Apache Sparkを基盤とした統合データ分析プラットフォームであり、大規模なデータエンジニアリング、機械学習、協調分析を目的としています。EMQXは、Databricksが管理するAmazon S3バケットにMQTTデータを書き込み、Databricksが外部ロケーションを通じて直接クエリできる形でDatabricksと連携します。
本ページでは、EMQXとDatabricks間のデータ統合について詳細に解説し、コネクターおよびSinkの作成方法を実践的に案内します。
動作概要
EMQXにおけるDatabricksデータ統合は、Amazon S3連携を基盤としています。EMQXはMQTTデータをDatabricksが管理するS3バケットに書き込み、Databricksは外部ロケーションを介してこのバケットにアクセスし、保存されたデータに対して直接SQLクエリを実行します。

具体的なワークフローは以下の通りです:
- デバイスのEMQXへの接続:IoTデバイスがMQTTプロトコルで正常に接続するとオンラインイベントが発生します。このイベントにはデバイスID、送信元IPアドレスなどのプロパティ情報が含まれます。
- デバイスのメッセージパブリッシュと受信:デバイスは特定のトピックを通じてテレメトリやステータスデータをパブリッシュします。EMQXはこれらのメッセージを受信し、ルールエンジン内で比較処理を行います。
- ルールエンジンによるメッセージ処理:組み込みのルールエンジンは、トピックマッチングに基づき特定のソースからのメッセージやイベントを処理します。対応するルールとマッチし、データフォーマット変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などを行います。
- Amazon S3への書き込み:ルールがトリガーされると、Amazon S3 SinkがDatabricksワークスペースに紐づくS3バケットに処理済みデータを書き込みます。
- DatabricksによるS3からの読み込み:Databricksは外部ロケーションを通じてS3バケット内のデータを直接クエリし、リアルタイム分析や機械学習ワークフローを実現します。
特長とメリット
EMQXのDatabricksデータ統合を利用することで、以下の特長と利点が得られます:
- メッセージ変換:EMQXルール内でメッセージに対する高度な処理や変換を行い、S3への書き込み前にデータを整形可能です。これにより後続の保存や分析が容易になります。
- 柔軟なデータ操作:Amazon S3 Sinkを用いて、Databricks管理のS3バケットに特定フィールドを動的なオブジェクトキー設定で書き込めるため、柔軟なデータストレージが可能です。
- 統合分析プラットフォーム:EMQXとDatabricksを連携させることで、IoTデータを即座にDatabricksワークスペース内のSQL分析、機械学習、データエンジニアリングパイプラインに活用できます。
- 低コストの長期保存:基盤ストレージとしてS3を活用することで、高可用性かつ信頼性の高いコスト効率に優れたデータストアを実現し、大規模IoTワークロードに適しています。
はじめる前に
このセクションでは、EMQXでDatabricks向けAmazon S3コネクターおよびSinkを作成する前の準備について説明します。
前提条件
以下の内容を理解していることを確認してください:
EMQXの概念:
Databricksの概念:
- ワークスペース:Databricksの全資産にアクセスする環境です。
- 外部ロケーション:外部のS3パスをマッピングし、そこに保存されたデータをSQLで直接クエリ可能にするDatabricksの機能です。
- ストレージ認証情報:外部ストレージの読み書き権限を付与するDatabricks内のアクセス認証情報です。
AWS MarketplaceでのDatabricksセットアップ
ここでは、AWS MarketplaceでDatabricksをサブスクライブする例を用いてデプロイ方法を説明します。
AWS MarketplaceでDatabricksをサブスクライブします。Databricksアカウントとワークスペース作成の案内が表示されます。
サブスクライブ後、ワークスペースを作成します。リージョンとストレージオプションを選択し、Createをクリックします。

ワークスペース作成後、Workspaces一覧に表示されます。ワークスペース用に自動プロビジョニングされたS3バケット名(例:
databricks-workspace-stack-142ec-bucket)を控えておきます。このバケットにEMQXからのMQTTデータを保存します。
ワークスペースを開き、Catalog -> External locationsへ移動し、EMQXがデータを書き込むS3パスを指す外部ロケーションを作成します。

Create locationをクリックし、Storage typeを
S3に設定、URLにs3://databricks-workspace-stack-142ec-bucket/emqx-iot-data-newを入力し、Storage credentialを選択します。
S3バケットへの読み書き権限を持つIAMユーザーまたはロールのAWSアクセス認証情報(Access Key IDおよびSecret Access Key)を取得します。これらはEMQXコネクターの設定に使用します。
DatabricksワークスペースとS3バケットの設定が完了したら、EMQXでコネクターとSinkの作成準備が整います。
コネクターの作成
Amazon S3 Sinkを追加する前に、対応するコネクターを作成します。
- ダッシュボードのIntegration -> Connectorsページへ移動します。
- 右上のCreateボタンをクリックします。
- コネクタータイプとしてAmazon S3を選択し、Nextをクリックします。
- コネクター名を入力します。名前は英数字で始まり、英数字、ハイフン、アンダースコアを含めることができます。例として
my-databricksを入力します。 - 接続情報を入力します:
- Host:DatabricksワークスペースがデプロイされているAWSリージョンのS3エンドポイントを、
s3.{region}.amazonaws.com形式で入力します。 - Port:
443を入力します。 - Access Key IDおよびSecret Access Key:AWS MarketplaceでのDatabricksセットアップで取得したAWSアクセス認証情報を入力します。
- Host:DatabricksワークスペースがデプロイされているAWSリージョンのS3エンドポイントを、
- 残りの設定はデフォルト値を使用します。
- Createをクリックする前に、Test ConnectivityをクリックしてEMQXがS3サービスに接続できることを確認できます。
- Createをクリックしてコネクターの設定を完了します。作成成功のダイアログが表示され、ルールを今すぐ作成するか尋ねられます。Create Ruleをクリックするとコネクターが事前選択された状態でルール作成画面に進みます。Back To Connector Listをクリックすると戻って後でルールを作成できます。
Amazon S3 Sinkを用いたルールの作成
このセクションでは、ソースMQTTトピックt/#からのメッセージを処理し、処理結果をDatabricks管理のS3バケットに書き込むルール作成手順を示します。
前のステップでCreate Ruleをクリックした場合、Add Actionパネルが自動で開き、Type of Actionが
Amazon S3、コネクターが事前選択されています。ステップ5へ進んでください。そうでなければ、ダッシュボードのIntegration -> Rulesページに移動し、右上のCreateをクリックします。ルールIDを入力し、SQLエディターに以下のルールSQLを入力します:
sqlSELECT * FROM "t/#"TIP
初心者の方は、SQL Examplesをクリックし、Enable Testを有効にしてSQLルールの学習とテストを行うことができます。
右側の**+ Add Actionをクリックし、Add ActionパネルでType of Action**ドロップダウンから
Amazon S3を選択し、ActionはデフォルトのCreate Actionのままにします。Connectorsドロップダウンから先ほど作成した
my-databricksコネクターを選択します。ドロップダウン横の作成ボタンをクリックすると、ポップアップで新規コネクターを素早く作成できます。必要な設定パラメータはコネクターの作成を参照してください。Sinkの名前と任意で説明を入力します。
Bucketに
databricks-workspace-stack-142ec-bucketを入力します。このフィールドは${var}形式のプレースホルダーもサポートしますが、対応するバケットがS3に存在することを確認してください。アップロードするオブジェクトのアクセス権限を指定するため、必要に応じてACLを選択します。
Upload Methodを選択します:
- Direct Upload:ルールがトリガーされるたびに、事前設定したオブジェクトキーと内容に従ってデータを直接S3にアップロードします。バイナリや大きなテキストデータの保存に適しています。
- Aggregated Upload:複数のルールトリガー結果を1つのファイル(例:CSVファイル)にまとめてS3にアップロードします。構造化データの保存に適し、ファイル数を減らし書き込み効率を向上させます。
選択した方法により設定パラメータが異なります。以下のタブから該当する方法の設定を行ってください。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。詳細はフォールバックアクションを参照してください。
Advanced Settingsを展開し、必要に応じて詳細設定を行います(任意)。詳細は高度な設定を参照してください。
残りの設定はデフォルト値を使用します。Createをクリックする前に、Test ConnectivityをクリックしてSinkがS3サービスに接続できることを確認できます。
CreateをクリックしてSinkの作成を完了します。作成成功後、ルール作成画面に戻り、新しいSinkがルールアクションに追加されます。
ルール作成画面でSaveをクリックし、ルール作成全体を完了します。
これでルールの作成が完了しました。Rulesページで新規作成したルールを確認でき、**Actions (Sink)**タブで新しいAmazon S3 Sinkも確認できます。
ルールのテスト
MQTTXを使ってトピックt/1にメッセージをパブリッシュします:
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello Databricks" }'数件のメッセージ送信後、DatabricksワークスペースでWorkspaceを右クリックし、Create -> Notebookを選択して新規ノートブックを作成します。

ノートブック内で外部ロケーションに対してSQLクエリを実行し、データが正常に取り込まれていることを確認します:
SELECT * FROM json.`s3://databricks-workspace-stack-142ec-bucket/emqx-iot-data-new/`
高度な設定
このセクションでは、Amazon S3 Sinkの高度な設定オプションについて説明します。ダッシュボードのSink設定画面でAdvanced Settingsを展開し、用途に応じて以下のパラメータを調整できます。
| フィールド名 | 説明 | デフォルト値 |
|---|---|---|
| Buffer Pool Size | EMQXとS3間のデータフローを管理するバッファワーカープロセスの数を指定します。 | 16 |
| Request TTL | バッファに入ったリクエストが有効とみなされる最大時間(秒)を指定します。 | 45 |
| Health Check Interval | SinkがS3との接続状態を自動でヘルスチェックする間隔(秒)を指定します。 | 15秒 |
| Health Check Interval Jitter | 複数ノードが同時にヘルスチェックを開始する確率を減らすため、基本間隔に加える一様ランダム遅延(ミリ秒)です。 | 0ミリ秒 |
| Health Check Timeout | コネクターがS3との接続状態を自動ヘルスチェックする際のタイムアウト時間を指定します。 | 60秒 |
| Max Buffer Queue Size | S3 Sinkの各バッファワーカープロセスがバッファリング可能な最大バイト数を指定します。 | 256MB |
| Query Mode | メッセージ送信を最適化するため、synchronous(同期)またはasynchronous(非同期)リクエストモードを選択できます。 | Asynchronous |
| In-flight Window | SinkがS3と通信中に同時に存在可能なインフライトキューリクエストの最大数を制御します。 | 100 |
| Min Part Size | 集計完了後のパートアップロード時の最小チャンクサイズを指定します。 | 5MB |
| Max Part Size | パートアップロード時の最大チャンクサイズを指定します。 | 5GB |