# Azure Blob Storage に MQTT データを取り込む

[Azure Blob Storage](https://azure.microsoft.com/en-us/products/storage/blobs/) は、マイクロソフトのクラウドベースのオブジェクトストレージソリューションで、大量の非構造化データの取り扱いに特化しています。非構造化データとは、特定のデータモデルやフォーマットに従わないデータタイプ、例えばテキストファイルやバイナリデータを指します。EMQX は MQTT メッセージを効率的に Blob Storage コンテナに保存でき、IoT データの保存に柔軟なソリューションを提供します。

本ページでは、EMQX と Azure Blob Storage 間のデータ統合について詳しく紹介し、ルールおよび Sink の作成方法について実践的なガイダンスを提供します。

## 動作概要

EMQX における Azure Blob Storage データ統合はすぐに使える機能で、複雑なビジネス開発にも簡単に設定可能です。典型的な IoT アプリケーションでは、EMQX がデバイスの接続とメッセージの送受信を担う IoT プラットフォームとして機能し、Azure Blob Storage がメッセージデータの保存を担当するデータストレージプラットフォームとして機能します。

![azure-blob-storage-architecture](./assets/azure-blob-storage-architecture.png)

EMQX はルールエンジンと Sink を利用してデバイスのイベントやデータを Azure Blob Storage に転送します。アプリケーションは Azure Blob Storage からデータを読み取り、さらなるデータ活用が可能です。具体的なワークフローは以下の通りです。

1. **デバイスの EMQX への接続**：IoT デバイスは MQTT プロトコルで正常に接続するとオンラインイベントをトリガーします。このイベントにはデバイスID、送信元IPアドレスなどのプロパティ情報が含まれます。
2. **デバイスメッセージのパブリッシュと受信**：デバイスは特定のトピックを通じてテレメトリやステータスデータをパブリッシュします。EMQX はメッセージを受信し、ルールエンジン内で照合します。
3. **ルールエンジンによるメッセージ処理**：組み込みのルールエンジンはトピックマッチングに基づき特定のソースからのメッセージやイベントを処理します。対応するルールをマッチさせ、データフォーマット変換、特定情報のフィルタリング、コンテキスト情報の付加などの処理を行います。
4. **Azure Blob Storage への書き込み**：ルールがトリガーされると、メッセージをストレージコンテナに書き込むアクションが実行されます。Azure Blob Storage Sink を使うことで、処理結果からデータを抽出し Blob Storage に送信可能です。メッセージはテキストまたはバイナリ形式で保存でき、複数行の構造化データは CSV、JSON Lines、Parquet ファイルなどに集約可能で、メッセージ内容と Sink の設定に応じて選択されます。

イベントやメッセージデータがストレージコンテナに書き込まれた後は、Azure Blob Storage に接続してデータを読み取り、以下のような柔軟なアプリケーション開発に活用できます。

- データアーカイブ：デバイスメッセージを Azure Blob Storage のオブジェクトとして長期保存し、コンプライアンス要件やビジネスニーズに対応。
- データ分析：ストレージコンテナからデータを分析サービス（例：Snowflake）に取り込み、予知保全やデバイス効率評価などの分析に活用。

## 特長と利点

EMQX の 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 を作成する前に必要な準備を紹介します。

### 前提条件

- [ルール](./rules.md) の理解
- [データ統合](./data-bridges.md) の理解

### Azure Storage でコンテナを作成する

1. Azure Storage にアクセスするには Azure サブスクリプションが必要です。まだお持ちでない場合は、[無料アカウント](https://azure.microsoft.com/free/)を作成してください。

2. Azure Storage へのアクセスはすべてストレージアカウントを通じて行われます。このクイックスタートでは、[Azure ポータル](https://portal.azure.com/)、Azure PowerShell、または Azure CLI を使ってストレージアカウントを作成します。ストレージアカウントの作成方法は[ストレージアカウントの作成](https://learn.microsoft.com/en-us/azure/storage/common/storage-account-create)を参照してください。

3. Azure ポータルでコンテナを作成するには、新しく作成したストレージアカウントに移動し、左メニューの「データストレージ」セクションで「コンテナ」を選択します。+ **コンテナ** ボタンを押し、新しいコンテナ名に `iot-data` を入力して **作成** をクリックします。

   ![azure-storage-container-create](./assets/azure-storage-container-create.png)

4. ストレージアカウントの **セキュリティ+ネットワーク** -> **アクセスキー** に移動し、**キー** をコピーします。このキーは EMQX で Sink を設定する際に必要です。

   ![azure-storage-access-keys](./assets/azure-storage-access-keys.png)

## コネクターを作成する

Azure Blob Storage Sink を追加する前に、対応するコネクターを作成する必要があります。

1. ダッシュボードの **Integration** -> **Connector** ページに移動します。
2. 右上の **作成** ボタンをクリックします。
3. コネクタータイプで **Azure Blob Storage** を選択し、次へ進みます。
4. コネクター名を入力します。英数字の組み合わせで、ここでは `my-azure` と入力します。
5. 接続情報を入力します。
   - **Account Name**：ストレージアカウント名
   - **Account Key**：前ステップで取得したストレージアカウントのキー
6. **作成** をクリックする前に、**接続テスト** をクリックして Azure Storage への接続が成功するか確認できます。
7. 画面下部の **作成** ボタンをクリックしてコネクターの作成を完了します。

これでコネクターの作成が完了しました。次に、Azure Storage に書き込むデータを指定するルールと Sink を作成します。

## Azure Blob Storage Sink を使ったルールの作成

このセクションでは、EMQX でソース MQTT トピック `t/#` からのメッセージを処理し、処理結果を設定済みの Sink を通じて Azure Storage の `iot-data` コンテナに書き込むルールの作成方法を示します。

1. ダッシュボードの **Integration** -> **Rules** ページに移動します。

2. 右上の **作成** ボタンをクリックします。

3. ルールID に `my_rule` を入力し、SQL エディターに以下のルール SQL を入力します。

   ```sql
   SELECT
     *
   FROM
       "t/#"
   ```

   ::: tip

   SQL に不慣れな場合は、**SQL Examples** と **Enable Debug** をクリックしてルール SQL の学習や結果のテストが可能です。

   :::

4. アクションを追加し、**Action Type** ドロップダウンリストから `Azure Blob Storage` を選択します。アクションのドロップダウンはデフォルトの `create action` のままにするか、既存の Azure Blob Storage アクションを選択します。ここでは新しい Sink を作成してルールに追加します。

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

6. コネクターのドロップダウンから先ほど作成した `my-azure` を選択します。ドロップダウン横の作成ボタンをクリックするとポップアップで新規コネクターを素早く作成可能です。必要な設定パラメータは[コネクターの作成](#コネクターを作成する)を参照してください。

7. **Container** に `iot-data` を入力します。

8. **Upload Method** を選択します。2つの方法の違いは以下の通りです。

   - **Direct Upload**：ルールがトリガーされるたびに、設定済みのオブジェクトキーと内容に従ってデータを直接 Azure Blob Storage にアップロードします。バイナリや大きなテキストデータの保存に適していますが、多数のファイルが生成される可能性があります。
   - **Aggregated Upload**：複数のルールトリガーの結果を1つのファイル（例：CSVファイル）にまとめて Azure Blob Storage にアップロードします。構造化データの保存に適し、ファイル数を減らし書き込み効率を向上させます。

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

   :::: tabs type

   ::: tab Direct Upload

   Direct Upload で設定する項目は以下の通りです。

   - **Blob Name**：コンテナ内のアップロード先オブジェクトの場所を定義します。`${var}` 形式のプレースホルダーをサポートし、`/` でストレージディレクトリを指定可能です。管理や識別のためにオブジェクトのサフィックスも設定します。ここでは `msgs/${clientid}_${timestamp}.json` と入力します。`${clientid}` はクライアントID、`${timestamp}` はメッセージのタイムスタンプを表し、各デバイスのメッセージが異なるオブジェクトに書き込まれます。
   - **Object Content**：デフォルトはすべてのフィールドを含む JSON テキスト形式です。`${var}` 形式のプレースホルダーをサポートします。ここでは `${payload}` を入力し、メッセージ本文をオブジェクトの内容として使用します。オブジェクトの保存形式はメッセージ本文の形式に依存し、圧縮ファイル、画像、その他バイナリ形式もサポートします。

   :::

   ::: tab Aggregate Upload

   Aggregate Upload で設定するパラメータは以下の通りです。

   - **Blob Name**：オブジェクトの保存パスを指定します。以下の変数が使用可能です。

     - **`${action}`**：アクション名（必須）
     - **`${node}`**：アップロードを実行する EMQX ノード名（必須）
     - **`${datetime.{format}}`**：集約開始日時。`{format}` で指定するフォーマット（必須）
       - **`${datetime.rfc3339utc}`**：UTC フォーマットの RFC3339 日時
       - **`${datetime.rfc3339}`**：ローカルタイムゾーンの RFC3339 日時
       - **`${datetime.unix}`**：Unix タイムスタンプ
     - **`${datetime_until.{format}}`**：集約終了日時（フォーマットは上記と同様）
     - **`${sequence}`**：同一時間帯の集約アップロードのシーケンス番号（必須）

     必須マークのあるプレースホルダーがテンプレートに含まれていない場合、自動的にパスのサフィックスとして追加され、重複を防ぎます。その他のプレースホルダーは無効とみなされます。

   - **Aggregation Type**：Azure Storage に保存するバッチ化された MQTT メッセージのデータファイル形式を定義します。サポートされる値：

      - `CSV`：カンマ区切りの CSV 形式で書き込みます。
      - `JSON Lines`：[JSON Lines](https://jsonlines.org/) 形式で書き込みます。
      - `parquet`：[Apache Parquet](https://parquet.apache.org/) 形式で書き込みます。カラム型で大規模データの分析クエリに最適化されています。

        > スキーマ定義、圧縮、行グループ設定など詳細な構成は[Parquet フォーマットオプション](#parquet-format-options)を参照してください。

   - **Column Order**（`CSV` 用）：ルール結果のカラム順序をドロップダウンで調整可能。生成される CSV ファイルは選択したカラム順にソートされ、未選択カラムはアルファベット順に続きます。

   - **Max Records**：最大レコード数に達すると、1ファイルの集約が完了してアップロードされ、時間間隔がリセットされます。

   - **Time Interval**：時間間隔に達すると、最大レコード数に満たなくても集約が完了してアップロードされ、最大レコード数がリセットされます。

   :::

   ::::

10. **フォールバックアクション（オプション）**：メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義可能です。プライマリ Sink がメッセージ処理に失敗した場合にトリガーされます。詳細は[フォールバックアクション](./data-bridges.md#fallback-actions)を参照してください。

11. **詳細設定**を展開し、必要に応じて高度な設定オプションを構成します（任意）。詳細は[詳細設定](#advanced-settings)を参照してください。

12. 残りの設定はデフォルト値のままにし、**作成** ボタンをクリックして Sink の作成を完了します。作成成功後、ルール作成画面に戻り、新しい Sink がルールアクションに追加されます。

13. ルール作成画面に戻り、**作成** ボタンをクリックしてルール作成を完了します。

これでルールの作成が完了しました。**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 の [スキーマレジストリ](./schema-registry.md) で管理されている既存の [Avro スキーマ](./schema-registry-example-avro.md) を使用します。

  このオプションを選択した場合、シリアライズに使用するスキーマを識別するために **スキーマ名** の指定が必要です。

  ::: tip

  スキーマを中央管理し、複数システム間で一貫したスキーマ進化を行いたい場合に推奨します。

  :::

- **Avro スキーマを直接定義**：EMQX 内でスキーマ JSON 構造を直接入力して定義します。

  例：

  ```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` にメッセージをパブリッシュします。

```bash
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "Hello Azure" }'
```

数件のメッセージを送信した後、Azure ポータルにアクセスして `iot-data` コンテナにアップロードされたオブジェクトを確認します。

[Azure portal](https://portal.azure.com/) にログインし、ストレージアカウントに移動して `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 メッセージのパブリッシュをブロックしませんが、クライアントがメッセージを受信してから Azure Storage に到達するまでのタイムラグが生じる可能性があります。 | `Asynchronous` |
| **バッチサイズ**           | EMQX から Azure Storage へ一度の転送で送信するデータバッチの最大サイズを指定します。サイズを調整することでデータ転送の効率と性能を最適化できます。<br />「バッチサイズ」を「1」に設定すると、データレコードは個別に送信され、バッチ化されません。 | `1`          |
| **インフライトウィンドウ** | 「インフライトキューリクエスト」とは、送信済みで応答やアックをまだ受け取っていないリクエストを指します。この設定は Sink と Azure Storage 間の通信で同時に存在可能なインフライトキューリクエストの最大数を制御します。<br/>**Request Mode** が `asynchronous` の場合、このパラメータは特に重要です。同一 MQTT クライアントからのメッセージを厳密に順序処理する必要がある場合は、この値を `1` に設定してください。 | `100`        |
