Skip to content

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 Databricksデータ統合

具体的なワークフローは以下の通りです:

  1. デバイスのEMQXへの接続:IoTデバイスがMQTTプロトコルで正常に接続するとオンラインイベントが発生します。このイベントにはデバイスID、送信元IPアドレスなどのプロパティ情報が含まれます。
  2. デバイスのメッセージパブリッシュと受信:デバイスは特定のトピックを通じてテレメトリやステータスデータをパブリッシュします。EMQXはこれらのメッセージを受信し、ルールエンジン内で比較処理を行います。
  3. ルールエンジンによるメッセージ処理:組み込みのルールエンジンは、トピックマッチングに基づき特定のソースからのメッセージやイベントを処理します。対応するルールとマッチし、データフォーマット変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などを行います。
  4. Amazon S3への書き込み:ルールがトリガーされると、Amazon S3 SinkがDatabricksワークスペースに紐づくS3バケットに処理済みデータを書き込みます。
  5. 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の概念: ​

  • ルールエンジン:MQTTメッセージからデータを抽出・変換するロジックを定義する方法を理解します。
  • データ統合:EMQXのコネクターおよびSinkの概念を理解します。

Databricksの概念: ​

  • ワークスペース:Databricksの全資産にアクセスする環境です。
  • 外部ロケーション:外部のS3パスをマッピングし、そこに保存されたデータをSQLで直接クエリ可能にするDatabricksの機能です。
  • ストレージ認証情報:外部ストレージの読み書き権限を付与するDatabricks内のアクセス認証情報です。

AWS MarketplaceでのDatabricksセットアップ ​

ここでは、AWS MarketplaceでDatabricksをサブスクライブする例を用いてデプロイ方法を説明します。

  1. AWS MarketplaceでDatabricksをサブスクライブします。Databricksアカウントとワークスペース作成の案内が表示されます。

  2. サブスクライブ後、ワークスペースを作成します。リージョンとストレージオプションを選択し、Createをクリックします。

    Databricksワークスペース作成

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

    Databricksワークスペース一覧

  3. ワークスペースを開き、Catalog -> External locationsへ移動し、EMQXがデータを書き込むS3パスを指す外部ロケーションを作成します。

    Databricks外部ロケーション

    Create locationをクリックし、Storage typeをS3に設定、URLにs3://databricks-workspace-stack-142ec-bucket/emqx-iot-data-newを入力し、Storage credentialを選択します。

    外部ロケーション作成

  4. S3バケットへの読み書き権限を持つIAMユーザーまたはロールのAWSアクセス認証情報(Access Key IDおよびSecret Access Key)を取得します。これらはEMQXコネクターの設定に使用します。

DatabricksワークスペースとS3バケットの設定が完了したら、EMQXでコネクターとSinkの作成準備が整います。

コネクターの作成 ​

Amazon S3 Sinkを追加する前に、対応するコネクターを作成します。

  1. ダッシュボードのIntegration -> Connectorsページへ移動します。
  2. 右上のCreateボタンをクリックします。
  3. コネクタータイプとしてAmazon S3を選択し、Nextをクリックします。
  4. コネクター名を入力します。名前は英数字で始まり、英数字、ハイフン、アンダースコアを含めることができます。例としてmy-databricksを入力します。
  5. 接続情報を入力します:
    • Host:DatabricksワークスペースがデプロイされているAWSリージョンのS3エンドポイントを、s3.{region}.amazonaws.com形式で入力します。
    • Port:443を入力します。
    • Access Key IDおよびSecret Access Key:AWS MarketplaceでのDatabricksセットアップで取得したAWSアクセス認証情報を入力します。
  6. 残りの設定はデフォルト値を使用します。
  7. Createをクリックする前に、Test ConnectivityをクリックしてEMQXがS3サービスに接続できることを確認できます。
  8. Createをクリックしてコネクターの設定を完了します。作成成功のダイアログが表示され、ルールを今すぐ作成するか尋ねられます。Create Ruleをクリックするとコネクターが事前選択された状態でルール作成画面に進みます。Back To Connector Listをクリックすると戻って後でルールを作成できます。

Amazon S3 Sinkを用いたルールの作成 ​

このセクションでは、ソースMQTTトピックt/#からのメッセージを処理し、処理結果をDatabricks管理のS3バケットに書き込むルール作成手順を示します。

  1. 前のステップでCreate Ruleをクリックした場合、Add Actionパネルが自動で開き、Type of ActionがAmazon S3、コネクターが事前選択されています。ステップ5へ進んでください。そうでなければ、ダッシュボードのIntegration -> Rulesページに移動し、右上のCreateをクリックします。

  2. ルールIDを入力し、SQLエディターに以下のルールSQLを入力します:

    sql
    SELECT
      *
    FROM
        "t/#"

    TIP

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

  3. 右側の**+ Add Actionをクリックし、Add ActionパネルでType of Action**ドロップダウンからAmazon S3を選択し、ActionはデフォルトのCreate Actionのままにします。

  4. Connectorsドロップダウンから先ほど作成したmy-databricksコネクターを選択します。ドロップダウン横の作成ボタンをクリックすると、ポップアップで新規コネクターを素早く作成できます。必要な設定パラメータはコネクターの作成を参照してください。

  5. Sinkの名前と任意で説明を入力します。

  6. Bucketにdatabricks-workspace-stack-142ec-bucketを入力します。このフィールドは${var}形式のプレースホルダーもサポートしますが、対応するバケットがS3に存在することを確認してください。

  7. アップロードするオブジェクトのアクセス権限を指定するため、必要に応じてACLを選択します。

  8. Upload Methodを選択します:

    • Direct Upload:ルールがトリガーされるたびに、事前設定したオブジェクトキーと内容に従ってデータを直接S3にアップロードします。バイナリや大きなテキストデータの保存に適しています。
    • Aggregated Upload:複数のルールトリガー結果を1つのファイル(例:CSVファイル)にまとめてS3にアップロードします。構造化データの保存に適し、ファイル数を減らし書き込み効率を向上させます。

    選択した方法により設定パラメータが異なります。以下のタブから該当する方法の設定を行ってください。

  9. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。詳細はフォールバックアクションを参照してください。

  10. Advanced Settingsを展開し、必要に応じて詳細設定を行います(任意)。詳細は高度な設定を参照してください。

  11. 残りの設定はデフォルト値を使用します。Createをクリックする前に、Test ConnectivityをクリックしてSinkがS3サービスに接続できることを確認できます。

  12. CreateをクリックしてSinkの作成を完了します。作成成功後、ルール作成画面に戻り、新しいSinkがルールアクションに追加されます。

  13. ルール作成画面でSaveをクリックし、ルール作成全体を完了します。

これでルールの作成が完了しました。Rulesページで新規作成したルールを確認でき、**Actions (Sink)**タブで新しいAmazon S3 Sinkも確認できます。

ルールのテスト ​

MQTTXを使ってトピックt/1にメッセージをパブリッシュします:

bash
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello Databricks" }'

数件のメッセージ送信後、DatabricksワークスペースでWorkspaceを右クリックし、Create -> Notebookを選択して新規ノートブックを作成します。

ノートブック作成

ノートブック内で外部ロケーションに対してSQLクエリを実行し、データが正常に取り込まれていることを確認します:

sql
SELECT * FROM json.`s3://databricks-workspace-stack-142ec-bucket/emqx-iot-data-new/`

Databricksクエリ結果

高度な設定 ​

このセクションでは、Amazon S3 Sinkの高度な設定オプションについて説明します。ダッシュボードのSink設定画面でAdvanced Settingsを展開し、用途に応じて以下のパラメータを調整できます。

フィールド名説明デフォルト値
Buffer Pool SizeEMQXとS3間のデータフローを管理するバッファワーカープロセスの数を指定します。16
Request TTLバッファに入ったリクエストが有効とみなされる最大時間(秒)を指定します。45
Health Check IntervalSinkがS3との接続状態を自動でヘルスチェックする間隔(秒)を指定します。15秒
Health Check Interval Jitter複数ノードが同時にヘルスチェックを開始する確率を減らすため、基本間隔に加える一様ランダム遅延(ミリ秒)です。0ミリ秒
Health Check TimeoutコネクターがS3との接続状態を自動ヘルスチェックする際のタイムアウト時間を指定します。60秒
Max Buffer Queue SizeS3 Sinkの各バッファワーカープロセスがバッファリング可能な最大バイト数を指定します。256MB
Query Modeメッセージ送信を最適化するため、synchronous(同期)またはasynchronous(非同期)リクエストモードを選択できます。Asynchronous
In-flight WindowSinkがS3と通信中に同時に存在可能なインフライトキューリクエストの最大数を制御します。100
Min Part Size集計完了後のパートアップロード時の最小チャンクサイズを指定します。5MB
Max Part Sizeパートアップロード時の最大チャンクサイズを指定します。5GB