MQTTデータをRedisに取り込む
Redisは、オープンソースのインメモリデータストアであり、データベース、キャッシュ、ストリーミングエンジン、メッセージブローカーとして数百万の開発者に利用されています。EMQXはRedisとの統合をサポートしており、MQTTメッセージやクライアントイベントをRedisに保存できます。Redisデータ統合により、メッセージのキャッシュやクライアントイベントの統計にRedisを活用できます。
本ページでは、EMQXとRedis間のデータ統合の詳細な概要と、データ統合の作成および検証に関する実践的な手順を提供します。
動作の仕組み
Redisデータ統合はEMQXの標準機能であり、EMQXのリアルタイムデータキャプチャと送信機能を、Redisの豊富なデータ構造と強力なキー・バリューの読み書き性能と組み合わせています。組み込みのルールエンジンコンポーネントにより、EMQXからRedisへのデータ取り込みが簡素化され、複雑なコーディングが不要になります。
以下の図は、EMQXとRedis間のデータ統合の典型的なアーキテクチャを示しています。

MQTTデータをRedisに取り込む流れは以下の通りです。
- メッセージのパブリッシュと受信:産業用IoTデバイスはMQTTプロトコルを通じてEMQXに正常に接続し、機械、センサー、製品ラインの稼働状態、計測値、またはトリガーイベントに基づくリアルタイムMQTTデータをEMQXにパブリッシュします。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
- メッセージデータの処理:メッセージが到着するとルールエンジンを経由し、EMQXで定義されたルールに従って処理されます。ルールは事前定義された条件に基づき、どのメッセージをRedisにルーティングするかを決定します。ペイロード変換を指定するルールがあれば、データ形式の変換、特定情報のフィルタリング、追加コンテキストによるペイロードの拡充などが適用されます。
- Redisへのデータ取り込み:ルールエンジンがデータを処理した後、キャッシュやカウントなどの操作のために、あらかじめ設定されたRedisコマンドを実行するアクションがトリガーされます。
- データの保存と活用:Redisに保存されたデータを読み取ることで、企業はRedisの豊富なデータ操作機能を活用し、多様なユースケースを実現できます。例えば物流分野では、デバイスの最新状態の取得やGPS地理位置情報の分析、リアルタイムデータ分析やソートなどの操作が可能です。これにより、リアルタイム追跡やルート推奨などの機能を実現できます。
特長とメリット
Redisとのデータ統合は、効率的なデータ送信、処理、活用を実現するための多彩な特長とメリットを提供します。
- 高性能かつスケーラブル:EMQXの分散アーキテクチャとRedisのクラスター モードにより、データ量の増加に応じてアプリケーションをシームレスにスケール可能です。大規模データセットでも一貫した性能と応答性を確保します。
- リアルタイムデータストリーム:EMQXはリアルタイムデータストリーム処理に特化しており、デバイスからRedisへの効率的かつ信頼性の高いデータ送信を実現します。Redisは迅速なデータ操作を実行できるため、リアルタイムデータキャッシュに最適なデータストレージコンポーネントです。
- リアルタイムデータ分析:Redisはデバイス接続数、メッセージパブリッシュ数、特定のビジネス指標などのリアルタイムメトリクスを計算可能です。EMQXはリアルタイムメッセージの送受信と処理を担い、データ分析のためのリアルタイム入力を提供します。
- 地理位置情報分析:Redisは地理空間データ構造とコマンドを備え、地理位置情報の保存と検索が可能です。EMQXの強力なデバイス接続機能と組み合わせることで、物流、コネクテッドカー、スマートシティなど多様なIoTアプリケーションに広く応用できます。
はじめる前に
このセクションでは、Redisデータ統合を作成する前に必要な準備とRedisサーバーのセットアップ方法について説明します。
前提条件
Redisサーバーのインストール
Dockerを使ってRedisをインストールし、起動します。
# Redisコンテナを起動し、パスワードをpublicに設定
docker run --name redis -p 6379:6379 -d redis --requirepass "public"
# コンテナにアクセス
docker exec -it redis bash
# Redisサーバーにアクセスし、AUTHコマンドで認証
redis-cli
127.0.0.1:6379> AUTH public
OK
# インストールの確認
127.0.0.1:6379> set emqx "Hello World"
OK
127.0.0.1:6379> get emqx
"Hello World"これでRedisのインストールが成功し、SETおよびGETコマンドで動作確認ができました。Redisのコマンドの詳細はRedis Commandsをご参照ください。
コネクターの作成
このセクションでは、Redis SinkをRedisサーバーに接続するためのコネクター作成手順を示します。
以下の手順は、EMQXとRedisを同一ローカルマシン上で実行していることを前提としています。Redisが別の環境にある場合は設定を適宜調整してください。
- ダッシュボードに入り、Integration -> Connectorsをクリックします。
- ページ右上のCreateをクリックします。
- Create ConnectorページでRedisを選択し、Nextをクリックします。
- コネクターの名前を入力します。名前は英数字の組み合わせとしてください(例:
my_redis)。 - ビジネスニーズに応じてRedis Modeを設定します(例:
single)。 - 接続情報を入力します。
- Server Host:
127.0.0.1:6379を入力 - Username:
adminを入力 - Password:
publicを入力 - Database ID:
0を入力 - その他のオプションはビジネスニーズに応じて設定してください。
- 暗号化接続を確立したい場合は、Enable TLSのトグルスイッチをオンにします。TLS接続の詳細は外部リソースアクセスのTLSをご覧ください。
- Server Host:
- Createをクリックする前に、Test ConnectivityをクリックしてコネクターがRedisサーバーに接続できるかテストできます。
- ページ下部のCreateボタンをクリックしてコネクター作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてルールとSinkの作成を続けてください。詳細はルールとRedis Sinkの作成をご参照ください。
Redis Sinkを使ったルールの作成
このセクションでは、各クライアントの最新メッセージをキャッシュし、メッセージ破棄の統計を収集するルールの作成方法を示します。
メッセージキャッシュと統計機能用に2つの異なるRedis Sinkを作成する必要があります。作成するSinkの種類に応じて、以下のRedisコマンドテンプレートの設定手順に従ってください。
EMQXダッシュボードにアクセスし、Integration -> Rulesをクリックします。
ページ右上のCreateをクリックします。
ルールIDに
cache_to_redisを入力し、利用する機能に応じてSQL Editorにルールを設定します。メッセージキャッシュ用ルールを作成する場合、以下の文を入力します。これはトピック
t/#配下のMQTTメッセージをRedisに保存することを意味します。注意:独自のSQL構文を指定する場合は、Sinkで必要なすべてのフィールドが
SELECT句に含まれていることを確認してください。bashSELECT * FROM "t/#"メッセージ破棄統計用ルールを作成する場合、以下の文を入力します。
bashSELECT * FROM "$events/message_dropped", "$events/delivery_dropped"EMQXのルールは2種類のメッセージ破棄イベントを定義しており、これらをトリガーとしてRedisに記録できます。
イベント名 トピック パラメーター 転送中にメッセージが破棄される場合 $events/message_dropped $events/message_dropped 配信中にメッセージが破棄される場合 $events/delivery_dropped $events/delivery_dropped
TIP
初心者の方は、SQL ExamplesとEnable TestをクリックしてSQLルールを学習・テストしてください。
- Add Actionボタンをクリックして、ルールによりトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理されたデータをRedisに送信します。
Type of Actionドロップダウンリストから
Redisを選択します。ActionはデフォルトのCreate Actionのままにします。既に作成済みのSinkを選択することも可能ですが、このデモでは新しいSinkを作成します。Sinkの名前を入力します。名前は英数字の組み合わせとしてください。
Connectorドロップダウンから
my_redisを選択します。ドロップダウン横のボタンから新規コネクターを作成することも可能です。設定パラメーターの詳細はコネクターの作成をご参照ください。利用する機能に応じてRedis Command Templateを設定します。
メッセージキャッシュ用Sinkを作成する場合、RedisのHSETコマンドとハッシュデータ構造を使い、
clientidをキーとして、username、payload、timestampなどのフィールドを保存します。Redis内の他のキーと区別するために、メッセージにはemqx_messagesプレフィックスを付け、:で区切ります。bash# HSET key field value [field value...] HSET emqx_messages:${clientid} username ${username} payload ${payload} timestamp ${timestamp}メッセージ破棄統計用Sinkを作成する場合、以下のHINCRBYコマンドを使い、各トピックごとの破棄メッセージ数を集計します。
bash# HINCRBY key field increment HINCRBY emqx_message_dropped_count ${topic} 1このコマンドが実行されるたびに、対応するカウンターが1ずつ増加します。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はフォールバックアクションをご覧ください。
高度な設定(任意):必要に応じてsyncまたはasyncクエリモードを選択できます。詳細はSinkの機能をご参照ください。
Createをクリックする前に、Test ConnectivityをクリックしてSinkがRedisサーバーに接続できるかテストできます。
CreateボタンをクリックしてSink設定を完了します。新しいSinkがAction Outputsに追加されます。
Create Ruleページに戻り、設定内容を確認してCreateボタンをクリックしルールを生成します。
これでRedis Sink用のルールが正常に作成されました。Integration -> Rulesページで新規作成したルールを確認できます。**Actions(Sink)**タブをクリックすると新しいRedis Sinkが表示されます。
また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/#配下のメッセージがルールmy_ruleで解析されてRedisに送信・保存されていることが確認できます。
ルールのテスト
MQTTXを使ってトピックt/1にメッセージを送信し、メッセージキャッシュイベントをトリガーします。トピックt/1にサブスクライバーがいない場合、メッセージは破棄され、メッセージ破棄ルールがトリガーされます。
mqttx pub -i emqx_c -u emqx_u -t t/1 -m '{ "msg": "hello Redis" }'2つのSinkの実行状況を確認すると、1件の新しいMatchedと1件のSent Successfullyが表示されるはずです。
メッセージがキャッシュされているか確認します。
127.0.0.1:6379> HGETALL emqx_messages:emqx_c
1) "username"
2) "emqx_u"
3) "payload"
4) "{ \"msg\": \"hello Redis\" }"
5) "timestamp"
6) "1675263885119"テストを再実行すると、timestampフィールドが更新されているはずです。
破棄されたメッセージが集計されているか確認します。
127.0.0.1:6379> HGETALL emqx_message_dropped_count
1) "t/1"
2) "1"テストを繰り返すと、t/1に対応するカウンターの数値も増加します。