Skip to content

Amazon S3へのMQTTデータ取り込み ​

Amazon S3 は、高い信頼性、安定性、セキュリティを誇るインターネットベースのストレージサービスであり、迅速なデプロイと使いやすさを特徴としています。EMQXはMQTTメッセージを効率的にAmazon S3バケットに保存でき、柔軟なIoTデータストレージ機能を実現します。

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

TIP

EMQXはAmazon S3以外にも、以下のようなS3プロトコル対応のストレージサービスに対応しています。

  • MinIO:MinIOは高性能な分散オブジェクトストレージシステムです。Amazon S3 API互換のオープンソースのオブジェクトストレージサーバーで、プライベートクラウド構築に適しています。
  • Google Cloud Storage:Google Cloud StorageはGoogle Cloudの統合オブジェクトストレージで、開発者や企業が大量のデータを保存できます。Amazon S3互換のインターフェースを提供しています。

ビジネスニーズや利用シナリオに応じて適切なストレージサービスを選択できます。

動作概要 ​

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

emqx-integration-s3

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

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

イベントやメッセージデータがAmazon S3に書き込まれた後、Amazon S3に接続してデータを読み取り、以下のような柔軟なアプリケーション開発が可能です。

  • データアーカイブ:デバイスメッセージをAmazon S3のオブジェクトとして長期保存し、コンプライアンス要件やビジネスニーズに対応。
  • データ分析:S3から分析サービス(例:Snowflake)にデータを取り込み、予知保全やデバイス効率評価などの分析を実施。

特長と利点 ​

EMQXのAmazon S3データ統合を利用することで、以下の特長と利点が得られます。

  • メッセージ変換:メッセージはEMQXルール内で高度に処理・変換されてからAmazon S3に書き込まれるため、後続の保存や活用が容易になります。
  • 柔軟なデータ操作:S3 Sinkを使うことで、特定のデータフィールドをAmazon S3バケットに簡単に書き込めます。バケット名やオブジェクトキーは動的に設定可能で柔軟なデータ保存をサポートします。
  • 統合されたビジネスプロセス:S3 Sinkにより、デバイスデータをAmazon S3の豊富なエコシステムアプリケーションと連携でき、データ分析やアーカイブなど多様なビジネスシナリオを実現します。
  • 低コストの長期保存:データベースと比較して、Amazon S3は高可用性・信頼性かつコスト効率の高いオブジェクトストレージサービスであり、長期保存に適しています。

これらの特長により、効率的で信頼性が高く、スケーラブルなIoTアプリケーションを構築し、ビジネスの意思決定や最適化に役立てることができます。

はじめる前に ​

ここでは、EMQXでAmazon S3 Sinkを作成する前の準備について説明します。

前提条件 ​

以下の内容に慣れていることを確認してください。

EMQXの概念: ​

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

AWSの概念: ​

AWS S3が初めての場合は以下を確認してください。

  • EC2:AWSの仮想マシンサービス(コンピュートインスタンス)。
  • IAM:AWS Identity and Access Management。インスタンスロールはそのインスタンス上で動作するプログラムに一時的な認証情報を発行可能。
  • IMDSv2:EC2のInstance Metadata Service v2。メタデータや一時認証情報を取得するためのトークンベースでより安全なサービス。

S3バケットの準備 ​

EMQXはAmazon S3およびその他のS3互換ストレージサービスをサポートしています。AWSクラウドサービスを利用するか、DockerでMinIOインスタンスを展開できます。

コネクターの作成 ​

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

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

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

  3. コネクタータイプでAmazon S3を選択し、Nextをクリックします。

  4. コネクター名を入力します。名前は英数字で始まり、英数字、ハイフン、アンダースコアを含めることができます。ここでは例としてmy-s3を入力します。

  5. 接続情報を入力します。

    • Amazon S3バケットを使用する場合は以下を入力します。
      • Host:リージョンによって異なり、s3.{region}.amazonaws.comの形式です。

      • Port:443を入力。

      • Access Key IDとSecret Access Key:

        • AWSで作成したアクセスキーを入力するか、
        • EC2上でIAMロールをアタッチしている場合は空欄のままにします。

        詳細はPrepare an S3 Bucketの「Amazon S3」タブを参照してください。

    • MinIOを使用する場合は以下を入力します。
      • Host:127.0.0.1を入力。リモートでMinIOを実行している場合は実際のホストアドレスを入力してください。
      • Port:9000を入力。
      • Access Key IDとSecret Access Key:MinIOで作成したアクセスキーを入力します。
  6. 残りの設定はデフォルト値を使用します。

  7. Createをクリックする前に、Test ConnectivityをクリックしてコネクターがS3サービスに接続できるかテスト可能です。

  8. 画面下部のCreateボタンをクリックしてコネクター作成を完了します。

これでコネクター作成が完了し、次にS3サービスに書き込むデータを指定するルールとSinkの作成に進みます。

Amazon S3 Sinkを使ったルールの作成 ​

ここでは、EMQXでソースMQTTトピックt/#からのメッセージを処理し、処理結果を設定済みのSink経由でS3のiot-dataバケットに書き込むルールの作成方法を示します。

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

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

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

    sql
    SELECT
      *
    FROM
        "t/#"

    TIP

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

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

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

  6. コネクタードロップダウンから先に作成したmy-s3コネクターを選択します。ドロップダウン横の作成ボタンをクリックしてポップアップで新規コネクターを素早く作成することも可能です。必要な設定パラメータはCreate a Connectorを参照してください。

  7. Bucketにiot-dataを入力します。このフィールドは${var}形式のプレースホルダーもサポートしますが、対応するバケットが事前にS3に作成されている必要があります。

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

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

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

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

  10. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はFallback Actionsを参照してください。

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

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

  13. ルール作成画面でCreateボタンをクリックし、ルール作成を完了します。

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

また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/#のメッセージがルールmy_ruleで解析されS3に書き込まれる流れを視覚的に確認できます。

Parquetフォーマットオプション ​

Aggregation Typeがparquetに設定されている場合、EMQXは集約されたルール結果をApache Parquet形式で保存します。Parquetは列指向かつ圧縮されたファイル形式で、分析処理に最適化されています。

ここではParquet出力フォーマットの全設定オプションを説明します。

Parquetスキーマ(Avro) ​

このオプションはMQTTメッセージのフィールドをParquetファイルの列にマッピングする方法を定義します。EMQXはApache Avroスキーマ仕様を用いてParquetデータの構造を記述します。

以下のいずれかを選択可能です。

  • スキーマレジストリに存在するAvroスキーマ:EMQXのスキーマレジストリで管理されている既存のAvroスキーマを使用します。

    この場合、シリアライズに使用するスキーマを識別するためにSchema Nameを指定する必要があります。

    TIP

    スキーマを中央管理し、複数システム間で一貫したスキーマ進化を行いたい場合にこのオプションを使用してください。

  • Avroスキーマを直接定義:EMQX内でスキーマJSON構造をSchema Definitionフィールドに直接入力し定義します。

    例:

    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 S3" }'

数件メッセージを送信した後、MinIOコンソールまたはAmazon S3コンソールにアクセスして結果を確認します。

詳細設定 ​

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

項目名説明デフォルト値
Buffer Pool SizeEMQXとS3間のデータフローを管理するバッファワーカープロセス数を指定します。これらのワーカーはデータを一時的に保存・処理し、ターゲットサービスに送信する前のパフォーマンス最適化とスムーズなデータ転送に重要です。16
Request TTLバッファに入ったリクエストが有効とみなされる最大時間(秒)を指定します。リクエストがこのTTLを超えてバッファに滞留するか、送信後にS3から適時の応答やアックが得られない場合、リクエストは期限切れと判断されます。45
Health Check IntervalSinkがS3との接続状態を自動的にヘルスチェックする間隔(秒)を指定します。15秒
Health Check Interval Jitter複数ノードが同時にヘルスチェックを開始する確率を減らすため、基本間隔に加える一様ランダム遅延です。複数のアクションやソースが同一コネクターを共有する場合、ジッターを有効にするとヘルスチェックがずれて実行されます。0ミリ秒
Health Check TimeoutコネクターがS3テーブルとの接続ヘルスチェックを行う際のタイムアウト時間を指定します。60秒
Max Buffer Queue SizeS3 Sinkの各バッファワーカーがバッファリングできる最大バイト数を指定します。バッファワーカーはデータを一時保存し、S3への送信を効率化します。システム性能やデータ転送要件に応じて調整してください。256 MB
Query Modeメッセージ送信を最適化するため、synchronousまたはasynchronousのリクエストモードを選択可能です。非同期モードではS3への書き込みがMQTTメッセージのパブリッシュ処理をブロックしませんが、クライアントがS3到達前にメッセージを受信する可能性があります。Asynchronous
In-flight Window「インフライトキューリクエスト」とは、開始されたがまだ応答やアックを受け取っていないリクエストを指します。この設定はSinkとS3間の通信中に同時に存在できるインフライトキューリクエストの最大数を制御します。
Request Modeがasynchronousの場合、同一MQTTクライアントからのメッセージを厳密に順序処理する必要がある場合はこの値を1に設定してください。
100
Min Part Size集約完了後のパートアップロードの最小チャンクサイズです。アップロードするデータはこのサイズに達するまでメモリに蓄積されます。5MB
Max Part Sizeパートアップロードの最大チャンクサイズです。S3 Sinkはこのサイズを超えるパートのアップロードを試みません。5GB