Azure Blob Storage に MQTT データを取り込む
Azure Blob Storage は、マイクロソフトが提供するクラウドベースのオブジェクトストレージソリューションで、大量の非構造化データを扱うために設計されています。非構造化データとは、特定のデータモデルやフォーマットに従わないデータタイプのことで、テキストファイルやバイナリデータなどが該当します。EMQX は MQTT メッセージを効率的に Blob Storage コンテナに保存でき、IoT データの保存に柔軟なソリューションを提供します。
本ページでは、EMQX と Azure Blob Storage 間のデータ統合について詳しく紹介し、ルールおよび Sink の作成方法について実践的なガイダンスを提供します。
動作の仕組み
EMQX における Azure Blob Storage データ統合は、すぐに利用可能な機能であり、複雑なビジネス開発にも簡単に設定できます。典型的な IoT アプリケーションでは、EMQX がデバイスの接続とメッセージ伝送を担う IoT プラットフォームとして機能し、Azure Blob Storage はメッセージデータの保存を担当するデータストレージプラットフォームとして利用されます。

EMQX はルールエンジンと Sink を利用してデバイスのイベントやデータを Azure Blob Storage に転送します。アプリケーションは Azure Blob Storage からデータを読み取り、さらなるデータ活用を行えます。具体的なワークフローは以下の通りです。
- デバイスの EMQX への接続:IoT デバイスは MQTT プロトコルで正常に接続されるとオンラインイベントをトリガーします。このイベントにはデバイスID、送信元IPアドレスなどのプロパティ情報が含まれます。
- デバイスのメッセージパブリッシュと受信:デバイスは特定のトピックを通じてテレメトリやステータスデータをパブリッシュします。EMQX はメッセージを受信し、ルールエンジン内で比較処理を行います。
- ルールエンジンによるメッセージ処理:組み込みのルールエンジンはトピックマッチングに基づき特定のソースからのメッセージやイベントを処理します。対応するルールにマッチしたメッセージやイベントに対し、データフォーマットの変換、特定情報のフィルタリング、コンテキスト情報の付加などを行います。
- Azure Blob Storage への書き込み:ルールはメッセージをストレージコンテナに書き込むアクションをトリガーします。Azure Blob Storage Sink を利用して、処理結果からデータを抽出し Blob Storage に送信します。メッセージはテキストまたはバイナリ形式で保存でき、複数行の構造化データはメッセージ内容や Sink 設定に応じて CSV、JSON Lines、Parquet ファイルにまとめて保存可能です。
イベントやメッセージデータがストレージコンテナに書き込まれた後は、Azure Blob Storage に接続してデータを読み取り、以下のような柔軟なアプリケーション開発に活用できます。
- データアーカイブ:デバイスメッセージを Azure Blob Storage のオブジェクトとして長期保存し、コンプライアンス要件やビジネスニーズに対応。
- データ分析:ストレージコンテナからデータを分析サービス(例:Snowflake)に取り込み、予知保全やデバイス効率評価などのデータ分析に利用。
特長とメリット
EMQX の Azure Blob Storage データ統合を利用することで、以下の特長と利点がビジネスにもたらされます。
- メッセージ変換:メッセージは Azure Blob Storage に書き込む前に EMQX ルール内で高度な処理や変換が可能で、後続の保存や利用を容易にします。
- 柔軟なデータ操作:Azure Blob Storage Sink により、特定のデータフィールドを Azure Blob Storage コンテナに簡単に書き込み可能で、コンテナやオブジェクトキーの動的設定にも対応し柔軟なデータ保存を実現します。
- 統合されたビジネスプロセス:Azure Blob Storage Sink はデバイスデータを Azure Blob Storage の豊富なエコシステムアプリケーションと組み合わせることができ、データ分析やアーカイブなど多様なビジネスシナリオを実現します。
- 低コストの長期保存:データベースと比較して、Azure Blob Storage は高可用性かつ信頼性の高いコスト効率の良いオブジェクトストレージサービスを提供し、長期保存に適しています。
これらの特長により、効率的で信頼性が高くスケーラブルな IoT アプリケーションの構築が可能となり、ビジネスの意思決定や最適化に役立ちます。
はじめる前に
このセクションでは、EMQX で Azure Blob Storage Sink を作成する前に必要な準備について説明します。
前提条件
Azure Storage でコンテナを作成する
Azure Storage にアクセスするには Azure サブスクリプションが必要です。まだお持ちでない場合は、無料アカウントを作成してください。
Azure Storage へのすべてのアクセスはストレージアカウントを通じて行われます。このクイックスタートでは、Azure ポータル、Azure PowerShell、または Azure CLI を使ってストレージアカウントを作成します。ストレージアカウント作成の詳細はストレージアカウントの作成を参照してください。
Azure ポータルでコンテナを作成するには、新しく作成したストレージアカウントに移動します。ストレージアカウントの左メニューの「データストレージ」セクションまでスクロールし、「コンテナ」を選択します。+ コンテナ ボタンを押し、新しいコンテナ名に
iot-dataと入力し、作成 をクリックしてコンテナを作成します。
ストレージアカウントの セキュリティ+ネットワーク -> アクセスキー に移動し、キー をコピーします。EMQX で Sink を設定する際にこのキーが必要です。

コネクターを作成する
Azure Blob Storage Sink を追加する前に、対応するコネクターを作成する必要があります。
- ダッシュボードの Integration -> Connector ページに移動します。
- 右上の 作成 ボタンをクリックします。
- コネクタータイプとして Azure Blob Storage を選択し、次へ進みます。
- コネクター名を入力します。英数字の組み合わせで、ここでは
my-azureと入力します。 - 接続情報を入力します。
- アカウント名:ストレージアカウント名
- アカウントキー:前のステップで取得したストレージアカウントキー
- 作成 をクリックする前に、接続テスト をクリックしてコネクターが Azure Storage に接続できるか確認できます。
- 下部の 作成 ボタンをクリックしてコネクターの作成を完了します。
これでコネクターの作成が完了しました。次に、Azure Storage サービスに書き込むデータを指定するためのルールと Sink を作成します。
Azure Blob Storage Sink を使ったルールの作成
このセクションでは、EMQX でソース MQTT トピック t/# からのメッセージを処理し、処理結果を設定済みの Sink を通じて Azure Storage の iot-data コンテナに書き込むルールの作成方法を示します。
ダッシュボードの Integration -> Rules ページに移動します。
右上の 作成 ボタンをクリックします。
ルールID に
my_ruleを入力し、SQLエディターに以下のルールSQLを入力します。sqlSELECT * FROM "t/#"TIP
SQL に不慣れな場合は、SQL Examples と Enable Debug をクリックしてルールSQLの学習やテストができます。
アクションを追加し、Action Type ドロップダウンリストから
Azure Blob Storageを選択します。アクションのドロップダウンはデフォルトのcreate actionのままにするか、既存の Azure Blob Storage アクションを選択します。ここでは新しい Sink を作成してルールに追加します。Sink の名前と説明を入力します。
コネクターのドロップダウンから先ほど作成した
my-azureコネクターを選択します。ドロップダウン横の作成ボタンをクリックすると、ポップアップで新しいコネクターを素早く作成することも可能です。必要な設定パラメータはコネクターの作成を参照してください。Container に
iot-dataと入力します。Upload Method を選択します。2つの方法の違いは以下の通りです。
- Direct Upload:ルールがトリガーされるたびに、設定済みのオブジェクトキーと内容に従ってデータを直接 Azure Blob Storage にアップロードします。バイナリや大きなテキストデータの保存に適していますが、多数のファイルが生成される可能性があります。
- Aggregated Upload:複数のルールトリガーの結果を1つのファイル(例:CSVファイル)にまとめて Azure Blob Storage にアップロードします。構造化データの保存に適し、ファイル数を減らし書き込み効率を向上させます。
設定パラメータは選択した方法により異なります。以下から該当する方法を選択して設定してください。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。プライマリ Sink がメッセージ処理に失敗した場合にこれらのアクションがトリガーされます。詳細はフォールバックアクションを参照してください。
詳細設定を展開し、必要に応じて高度な設定オプションを構成します(任意)。詳細は詳細設定を参照してください。
残りの設定はデフォルト値を使用し、作成 ボタンをクリックして Sink の作成を完了します。作成成功後、ページはルール作成画面に戻り、新しい Sink がルールアクションに追加されます。
ルール作成ページに戻り、作成 ボタンをクリックしてルール作成全体を完了します。
これでルールの作成が完了しました。Rules ページで新規作成したルールを確認でき、Actions (Sink) タブで新しい Azure Blob Storage Sink を確認できます。
また、Integration -> Flow Designer をクリックするとトポロジーを表示できます。トポロジーでは、トピック t/# のメッセージがルール my_rule によって解析され、Azure Storage コンテナに書き込まれる様子が視覚的に示されます。
Parquet フォーマットオプション
Aggregation Type が parquet に設定されている場合、EMQX は集約されたルール結果を Apache Parquet 形式で保存します。Parquet はカラム指向の圧縮ファイル形式で、分析ワークロードに最適化されています。
このセクションでは Parquet 出力フォーマットのすべての設定可能なオプションを説明します。
Parquet スキーマ(Avro)
このオプションは MQTT メッセージのフィールドを Parquet ファイルのカラムにどのようにマッピングするかを定義します。EMQX は Apache Avro スキーマ仕様を用いて Parquet データの構造を記述します。
以下のいずれかを選択できます。
スキーマレジストリに存在する Avro スキーマ:EMQX の スキーマレジストリ で管理されている既存の Avro スキーマ を使用します。
このオプションを選択した場合は、シリアライズに使用するスキーマを特定するために スキーマ名 を指定する必要があります。
TIP
スキーマを中央管理し、複数システム間で一貫したスキーマ進化を行いたい場合にこのオプションを使用してください。
Avro スキーマを直接定義:EMQX 内でスキーマ JSON 構造を スキーマ定義 フィールドに直接入力して Avro スキーマを定義します。
例:
json{ "type": "record", "name": "MessageRecord", "fields": [ {"name": "clientid", "type": "string"}, {"name": "timestamp", "type": "long"}, {"name": "payload", "type": "string"} ] }
TIP
フィールド名とデータ型がルールSQLの戻り値と一致していることを確認してください。不正確または不足しているフィールドは、Parquet への書き込み時にシリアライズエラーの原因となります。
Parquet デフォルト圧縮
このオプションは、各行グループ内の Parquet データページに適用される圧縮アルゴリズムを指定します。圧縮によりストレージ容量を削減し、データクエリ時の I/O 効率を向上させます。
サポートされる値:
| 値 | 説明 |
|---|---|
snappy(デフォルト) | 高速な圧縮・解凍と良好な圧縮率をバランスよく提供。ほとんどのケースで推奨されます。 |
zstd | 中程度のCPU使用率でより高い圧縮率を提供。大規模分析データや長期保存に最適。 |
None | 圧縮を無効化。デバッグ用途や圧縮不要な場合に適します。 |
Parquet 最大行グループバイト数
このオプションは Parquet の行グループの最大サイズ(バイト単位)を指定します。行グループはデータ読み書きの基本単位で、バッファされたデータサイズがこの閾値を超えると EMQX は現在の行グループをフラッシュし新しい行グループを開始します。
- デフォルト値:
128 MB
ガイドライン:
- Athena や Spark などの分析クエリで読み取り性能を向上させたい場合は値を大きくします。
- 書き込み時のメモリ使用量を抑えたい場合や小規模データセットの場合は値を小さくします。
TIP
Parquet リーダーは行グループ単位でデータを読み込みます。大きな行グループはメタデータのオーバーヘッドを減らし分析クエリ性能を向上させます。
ルールのテスト
このセクションでは、Direct Upload メソッドで設定したルールのテスト方法を示します。
MQTTX を使ってトピック t/1 にメッセージをパブリッシュします。
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "Hello Azure" }'数件のメッセージを送信した後、Azure ポータルにアクセスして iot-data コンテナ内のアップロードされたオブジェクトを確認します。
Azure ポータルにログインし、ストレージアカウントに移動して iot-data コンテナを開くと、アップロードされたオブジェクトが表示されます。
詳細設定
このセクションでは、Azure Blob Storage Sink の詳細な設定オプションについて説明します。ダッシュボードの Sink 設定画面で 詳細設定 を展開し、用途に応じて以下のパラメータを調整できます。
| 項目名 | 説明 | デフォルト値 |
|---|---|---|
| バッファプールサイズ | EMQX と Azure Storage 間のデータフローを管理するバッファワーカープロセスの数を指定します。これらのワーカーはデータを一時的に保持・処理し、ターゲットサービスへの送信を最適化しスムーズなデータ伝送を保証します。 | 16 |
| リクエスト TTL | バッファに入ったリクエストが有効とみなされる最大時間(秒)を指定します。リクエストがこの TTL を超えてバッファ内に滞留するか、送信後に Azure Storage からの応答やアックがタイムリーに得られない場合、リクエストは期限切れとみなされます。 | 45 |
| ヘルスチェック間隔 | Sink が Azure Storage との接続状態を自動的にチェックする間隔(秒)を指定します。 | 15 |
| 最大バッファキューサイズ | Azure Blob Storage Sink の各バッファワーカーがバッファリングできる最大バイト数を指定します。バッファワーカーはデータを一時的に保持し、効率的なデータストリーム処理を行います。システム性能やデータ伝送要件に応じて調整してください。 | 256 |
| クエリモード | メッセージ伝送を最適化するために、同期(synchronous)または非同期(asynchronous)のリクエストモードを選択できます。非同期モードでは Azure Storage への書き込みが MQTT メッセージパブリッシュをブロックしませんが、クライアントがメッセージ到着前に受信する可能性があります。 | Asynchronous |
| バッチサイズ | EMQX から Azure Storage へ一度に転送するデータバッチの最大サイズを指定します。サイズを調整することでデータ転送の効率と性能を最適化できます。 「バッチサイズ」が「1」の場合、データレコードは個別に送信され、バッチ化されません。 | 1 |
| インフライトウィンドウ | 「インフライトキューリクエスト」とは、開始されたがまだ応答やアックを受け取っていないリクエストを指します。この設定は Sink と Azure Storage 間の通信で同時に存在可能なインフライトキューリクエストの最大数を制御します。 リクエストモードが asynchronous の場合、このパラメータは特に重要です。同一 MQTT クライアントからのメッセージを厳密に順序処理したい場合は、この値を 1 に設定してください。 | 100 |