Amazon Kinesis への MQTT データストリーミング
AWS Kinesis は、AWS 上で提供されるフルマネージドのリアルタイムストリーミングデータ処理サービスで、ストリーミングデータの収集、処理、分析を容易に行えます。あらゆる規模のストリーミングデータをリアルタイムかつ経済的かつ効率的に処理でき、高い柔軟性を持ち、数十万のソースからの大量のストリーミングデータを低レイテンシで処理可能です。
EMQX は Amazon Kinesis Data Streams とシームレスに統合でき、大量の IoT デバイスを接続してリアルタイムにメッセージを収集・送信できます。このデータ統合により、Amazon Kinesis Data Streams と連携してリアルタイムデータ分析や複雑なストリーム処理を実現します。
本ページでは、EMQX と Amazon Kinesis 間のデータ統合について包括的に紹介し、データ統合の作成と検証方法を実践的に解説します。
動作の仕組み
Amazon Kinesis とのデータ統合は、EMQX の標準機能として提供されており、MQTT データストリームを Amazon Kinesis とシームレスに連携させ、IoT アプリケーション開発における豊富なサービスと機能を活用できるよう設計されています。

EMQX はルールエンジンと Sink を介して MQTT データを Amazon Kinesis に転送します。全体の流れは以下の通りです。
- IoT デバイスがメッセージをパブリッシュ:デバイスは特定のトピックを通じてテレメトリや状態データをパブリッシュし、ルールエンジンをトリガーします。
- ルールエンジンがメッセージを処理:組み込みのルールエンジンは、特定のソースからの MQTT メッセージをトピックマッチングに基づいて処理します。ルールエンジンは対応するルールをマッチさせ、データ形式の変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などを行います。
- Amazon Kinesis へのブリッジング:ルールによってトリガーされたアクションでメッセージを Amazon Kinesis に転送します。パーティションキー、書き込み先のデータストリーム、メッセージフォーマットなどをカスタム設定でき、柔軟なデータ統合が可能です。
MQTT メッセージデータが Amazon Kinesis に書き込まれた後は、以下のような柔軟なアプリケーション開発が可能です。
- リアルタイムデータ処理と分析:強力な Amazon Kinesis のデータ処理・分析ツールとストリーミング機能を活用し、メッセージデータのリアルタイム処理・分析を行い、有益なインサイトや意思決定支援を得られます。
- イベント駆動型機能:Amazon のイベント処理をトリガーし、動的かつ柔軟な関数の起動・処理を実現します。
- データの保存と共有:メッセージデータを Amazon Kinesis のストレージサービスに送信し、大量データの安全な保存・管理を行います。これにより他の Amazon サービスと連携してデータを共有・分析し、多様なビジネスニーズに対応可能です。
特長とメリット
EMQX と AWS Kinesis Data Streams のデータ統合は、以下の機能と利点をビジネスにもたらします。
- 信頼性の高いデータ伝送と順序保証:EMQX と AWS Kinesis Data Streams はともに信頼性の高いデータ伝送機構を備えています。EMQX は MQTT プロトコルを通じてメッセージの確実な送信を保証し、AWS Kinesis Data Streams はパーティションとシーケンス番号でメッセージの順序を保証します。これにより、デバイスから送信されたメッセージが正確に目的地に届き、正しい順序で処理されることを確実にします。
- リアルタイムデータ処理:デバイスからの高頻度データは、EMQX のルール SQL によるリアルタイムの一次処理を経て、MQTT メッセージのフィルタリング、抽出、付加、変換を容易に行えます。AWS Kinesis Data Streams への送信後は、AWS Lambda や AWS 管理の Apache Flink と連携してさらなるリアルタイム分析が可能です。
- 弾力的なスケーラビリティ対応:EMQX は数百万の IoT デバイス接続を容易に実現し、弾力的なスケーラビリティを提供します。一方、AWS Kinesis Data Streams はオンデマンドの自動リソース割り当てと拡張を採用しています。両者を組み合わせたアプリケーションは接続数やデータ量に応じてスケールし、ビジネスの成長に継続的に対応します。
- パーシステンスなデータ保存:AWS Kinesis Data Streams はパーシステンスなデータ保存機能を備え、毎秒数百万件のデバイスデータストリームを確実に保存します。必要に応じて過去データの取得が可能で、オフライン分析や処理を支援します。
AWS Kinesis Data Streams を活用したストリーミングデータパイプラインの構築は、EMQX と AWS プラットフォームの統合の難易度を大幅に軽減し、ユーザーにより豊富で柔軟なデータ処理ソリューションを提供します。これにより、EMQX ユーザーは AWS 上で機能的に充実した高性能なデータ駆動型アプリケーションを構築できます。
はじめる前に
このセクションでは、Amazon Kinesis データ統合を作成する前に必要な準備について説明します。Kinesis サービスのセットアップやデータストリームサービスのエミュレーション方法も含みます。
前提条件
Amazon Kinesis Data Streams でストリームを作成する
以下の手順に従い、AWS マネジメントコンソールからストリームを作成します(詳細はこちらのチュートリアルを参照してください)。
AWS マネジメントコンソールにサインインし、Kinesis コンソールを開きます。
ナビゲーションバーでリージョンセレクターを展開し、リージョンを選択します。
Create data stream を選択します。
Create Kinesis stream ページでデータストリーム名を入力し、On-demand キャパシティモードを選択します。
Amazon Kinesis Data Streams をローカルでエミュレートする
開発やテストを容易にするため、LocalStack を使って Amazon Kinesis Data Streams サービスをローカルでエミュレートできます。LocalStack を利用すると、リモートのクラウドプロバイダーに接続せずにローカルマシン上で AWS アプリケーションを完全に実行可能です。
Docker イメージを使ってインストールおよび起動します。
bash# LocalStack の Docker イメージをローカルで起動 docker run --name localstack -p '4566:4566' -e 'KINESIS_LATENCY=0' -d localstack/localstack:2.1 # コンテナにアクセス docker exec -it localstack bashシャード数 1 のストリーム
my_streamを作成します。bashawslocal kinesis create-stream --stream-name "my_stream" --shard-count 1
コネクターを作成する
このセクションでは、Sink を Amazon Kinesis Data Streams サービスに接続するためのコネクター作成方法を説明します。
- EMQX ダッシュボードに入り、Integration -> Connectors をクリックします。
- ページ右上の Create をクリックします。
- Create Connector ページで Amazon Kinesis を選択し、Next をクリックします。
- Configuration ステップで以下の情報を設定します。
- コネクター名を入力します。英数字の大文字・小文字の組み合わせとしてください。例:
my_kinesis - Amazon Kinesis Endpoint:Kinesis サービスのエンドポイントを入力します。LocalStackを使う場合は
http://localhost:4566を入力します。 - AWS Access Key ID:アクセスキーIDを入力します。LocalStackの場合は任意の値を入力してください。
- AWS Secret Access Key:シークレットアクセスキーを入力します。LocalStackの場合は任意の値を入力してください。
- コネクター名を入力します。英数字の大文字・小文字の組み合わせとしてください。例:
- Create をクリックする前に、Test Connectivity をクリックしてコネクターが Amazon Kinesis Data Streams サービスに接続できるかテストできます。
- ページ下部の Create ボタンをクリックしてコネクターの作成を完了します。ポップアップダイアログで Back to Connector List をクリックするか、Create Rule をクリックしてルールと Sink の作成を続行し、Amazon Kinesis に転送するデータを指定できます。詳細は Amazon Kinesis Sink を使ったルール作成 を参照してください。
Amazon Kinesis Sink を使ったルール作成
このセクションでは、ソース MQTT トピック t/# からのメッセージを処理し、処理結果を Amazon データストリーム my_stream にストリーミングするルールの作成方法を示します。
EMQX ダッシュボードで Integration -> Rules をクリックします。
ページ右上の Create をクリックします。
ルール ID に
my_ruleを入力します。SQL Editor でルールを設定します。トピック
t/#の MQTT メッセージを Amazon Kinesis Data Streams に保存したい場合、以下の SQL 文を使用できます。注意:独自の SQL 文を指定する場合は、Sink のペイロードテンプレートで必要なすべてのフィールドが
SELECT部分に含まれていることを確認してください。sqlSELECT * FROM "t/#"TIP
初心者の方は SQL Examples と Enable Test をクリックして、SQL ルールの学習とテストを行うことをおすすめします。
- Add Action ボタンをクリックして、ルールでトリガーされるアクションを定義します。このアクションにより、EMQX はルールで処理したデータを Kinesis に送信します。
Type of Action ドロップダウンリストから
Amazon Kinesisを選択します。Action ドロップダウンはデフォルトのCreate Actionのままにします。既に作成済みの Sink があれば選択可能ですが、この例では新規に Sink を作成します。Sink の名前と説明を入力します。名前は英数字の大文字・小文字の組み合わせにしてください。
Connector ドロップダウンから、先に作成した
my_kinesisを選択します。新しいコネクターを作成する場合は、ドロップダウン横のボタンをクリックしてください。設定パラメータは コネクター作成 を参照してください。以下の情報を入力します。
- Amazon Kinesis Stream: Amazon Kinesis Data Streams でストリームを作成する で作成したストリーム名を入力します。
- Partition Key:このストリームに送信されるレコードに関連付けるパーティションキーを入力します。
${variable_name}の形式のプレースホルダーも使用可能です(次のステップの例を参照)。
Payload Template フィールドは空欄のままにするか、テンプレートを定義します。
- 空欄の場合、クライアントID、トピック、ペイロードなど MQTT メッセージの可視入力すべてを JSON 形式でエンコードします。
- テンプレートを定義した場合、
${variable_name}の形式のプレースホルダーは MQTT コンテキストの対応する値で置き換えられます。例:${topic}は MQTT メッセージのトピックがmy/topicならmy/topicに置換されます。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリ Sink がメッセージ処理に失敗した場合にトリガーされます。詳細は フォールバックアクション を参照してください。
詳細設定(任意):必要に応じて詳細設定オプションを構成します。詳細は 詳細設定 を参照してください。
Create をクリックする前に、Test Connectivity をクリックして Sink が Amazon Kinesis Data Streams サービスに接続できるかテストできます。
Create ボタンをクリックして Sink の設定を完了します。新しい Sink が Action Outputs に追加されます。
Create Rule ページに戻り、設定内容を確認します。Create ボタンをクリックしてルールを生成します。
これで、Amazon Kinesis Sink を通じてデータを転送するルールが正常に作成されました。Integration -> Rules ページで新規作成したルールを確認できます。Actions(Sink) タブをクリックすると、新しい Amazon Kinesis Sink が表示されます。
また、Integration -> Flow Designer をクリックしてトポロジーを表示すると、トピック t/# のメッセージがルール my_rule によって解析され、Amazon Kinesis Data Streams に送信・保存されていることが確認できます。
ルールのテスト
MQTTX を使ってトピック
t/my_topicにメッセージを送信します。bashmqttx pub -i emqx_c -t t/my_topic -m '{ "msg": "hello Amazon Kinesis" }'Sink の稼働状況を確認すると、新規の受信メッセージと送信メッセージがそれぞれ1件ずつあるはずです。
Amazon Kinesis Data Viewer にアクセスし、レコードを取得するとメッセージが確認できます。
LocalStack を使った確認
LocalStack を使用する場合は、以下の手順で受信データを確認します。
メッセージを EMQX に送信する前に、ShardIterator を取得します。
bashawslocal kinesis get-shard-iterator --stream-name my_stream --shard-id shardId-000000000000 --shard-iterator-type LATEST { "ShardIterator": "AAAAAAAAAAG3YjBK9sp0uSIFGTPIYBI17bJ1RsqX4uJmRllBAZmFRnjq1kPLrgcyn7RVigmH+WsGciWpImxjXYLJhmqI2QO/DrlLfp6d1IyJFixg1s+MhtKoM6IOH0Tb2CPW9NwPYoT809x03n1zL8HbkXg7hpZjWXPmsEvkXjn4UCBf5dBerq7NLKS3RtAmOiXVN6skPpk=" }MQTTX を使ってトピック
t/my_topicにメッセージを送信します。bashmqttx pub -i emqx_c -t t/my_topic -m '{ "msg": "hello Amazon Kinesis" }'レコードを読み取り、受信データをデコードします。
bashawslocal kinesis get-records --shard-iterator="AAAAAAAAAAG3YjBK9sp0uSIFGTPIYBI17bJ1RsqX4uJmRllBAZmFRnjq1kPLrgcyn7RVigmH+WsGciWpImxjXYLJhmqI2QO/DrlLfp6d1IyJFixg1s+MhtKoM6IOH0Tb2CPW9NwPYoT809x03n1zL8HbkXg7hpZjWXPmsEvkXjn4UCBf5dBerq7NLKS3RtAmOiXVN6skPpk=" { "Records": [ { "SequenceNumber": "49642650476690467334495639799144299020426020544120356866", "ApproximateArrivalTimestamp": 1689389148.261, "Data": "eyAibXNnIjogImhlbGxvIEFtYXpvbiBLaW5lc2lzIiB9", "PartitionKey": "key", "EncryptionType": "NONE" } ], "NextShardIterator": "AAAAAAAAAAFj5M3+6XUECflJAlkoSNHV/LBciTYY9If2z1iP+egC/PtdVI2t1HCf3L0S6efAxb01UtvI+3ZSh6BO02+L0BxP5ssB6ONBPfFgqvUIjbfu0GOmzUaPiHTqS8nNjoBtqk0fkYFDOiATdCCnMSqZDVqvARng5oiObgigmxq8InciH+xry2vce1dF9+RRFkKLBc0=", "MillisBehindLatest": 0 } echo 'eyAibXNnIjogImhlbGxvIEFtYXpvbiBLaW5lc2lzIiB9' | base64 -d { "msg": "hello Amazon Kinesis" }
詳細設定
このセクションでは、Amazon Kinesis Sink の詳細設定オプションについて説明します。ダッシュボードの Sink 設定画面で Advanced Settings を展開し、用途に応じて以下のパラメータを調整できます。
| フィールド名 | 説明 | デフォルト値 |
|---|---|---|
| Buffer Pool Size | EMQX と Kinesis 間のデータフローを管理するバッファワーカーの数を指定します。これらのワーカーはデータを一時的に保存・処理し、ターゲットサービスへの送信を最適化しスムーズなデータ伝送を支えます。 | 16 |
| Request TTL | バッファに入ったリクエストが有効とみなされる最大時間(秒)を指定します。リクエストがこの TTL を超えてバッファに滞留するか、送信後に Kinesis からの応答やアックがタイムリーに得られない場合、リクエストは期限切れと判断されます。 | 45 秒 |
| Health Check Interval | Sink が Kinesis との接続の自動ヘルスチェックを行う間隔(秒)を指定します。 | 15 秒 |
| Health Check Interval Jitter | 複数ノードが同時にヘルスチェックを開始する確率を減らすため、基本間隔に加える一様ランダム遅延です。複数のアクションやソースが同じコネクターを共有する場合、ジッターを有効にするとヘルスチェックがわずかにずれて実行されます。 | 15 秒 |
| Health Check Timeout | コネクターが Kinesis との接続ヘルスチェックを行う際のタイムアウト時間(秒)を指定します。 | 60 秒 |
| Max Buffer Queue Size | Kinesis Sink の各バッファワーカーがバッファリングできる最大バイト数を指定します。バッファワーカーはデータを一時保存し、より効率的にデータストリームを処理します。システム性能やデータ伝送要件に応じて調整してください。 | 256 |
| Query Mode | メッセージ送信を最適化するために、synchronous(同期)または asynchronous(非同期)リクエストモードを選択できます。非同期モードでは Kinesis への書き込みが MQTT メッセージのパブリッシュ処理をブロックしませんが、クライアントがメッセージを受信するタイミングが Kinesis 到達前になる可能性があります。 | Async |
| Batch Size | EMQX から Kinesis へ一度に転送するデータバッチの最大サイズを指定します。サイズを調整することで、EMQX と Kinesis 間のデータ転送効率と性能を微調整できます。 「Batch Size」を「1」に設定すると、データレコードはバッチ化されず個別に送信されます。 | 1 |
| Inflight Window | 「インフライトキューリクエスト」とは、送信済みで応答やアックをまだ受け取っていないリクエストを指します。この設定は、Sink と Kinesis 間の通信中に同時に存在できるインフライトキューリクエストの最大数を制御します。 Request Mode が asynchronous の場合、このパラメータは特に重要です。同一 MQTT クライアントからのメッセージを厳密に順序通り処理したい場合は、この値を 1 に設定してください。 | 100 |