OpenTSDBへのMQTTデータ取り込み
OpenTSDBはスケーラブルで分散型の時系列データベースです。EMQXはOpenTSDBとの連携をサポートしており、MQTTメッセージをOpenTSDBに保存して後続の分析や取得に利用できます。
本ページでは、EMQXとOpenTSDB間のデータ統合について包括的に解説し、データ統合の作成および検証手順を実践的に説明します。
動作概要
OpenTSDBデータ統合はEMQXの標準機能であり、EMQXのリアルタイムデータキャプチャと転送機能をOpenTSDBのデータ保存・分析機能と組み合わせています。組み込みのルールエンジンコンポーネントにより、EMQXからOpenTSDBへのデータ取り込みを簡素化し、複雑なコーディングを不要にします。
以下の図はEMQXとOpenTSDB間の典型的なデータ統合アーキテクチャを示しています:

EMQXはルールエンジンとSinkを通じてデバイスデータをOpenTSDBに挿入します。OpenTSDBは豊富なクエリ機能を提供し、レポートやチャート、その他のデータ分析結果の生成をサポートします。産業用エネルギー管理シナリオを例にすると、ワークフローは以下の通りです:
- メッセージのパブリッシュと受信:産業用デバイスはMQTTプロトコルを通じてEMQXに正常に接続し、定期的にエネルギー消費データをパブリッシュします。このデータには生産ライン識別子やエネルギー消費値が含まれます。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
- ルールエンジンによるメッセージ処理:組み込みのルールエンジンはトピックマッチングに基づき特定のソースからのメッセージを処理します。メッセージが到着するとルールエンジンを通過し、対応するルールと照合されてメッセージデータを処理します。これにはデータ形式の変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などが含まれます。
- OpenTSDBへのデータ取り込み:ルールエンジンで定義されたルールがトリガーとなり、メッセージをOpenTSDBに書き込む操作を実行します。
データがOpenTSDBに書き込まれた後は、以下のように柔軟に利用できます:
- Grafanaなどの可視化ツールに接続し、エネルギー蓄積データを表示するチャートを生成する。
- 業務システムに接続してエネルギー蓄積デバイスの状態監視やアラートを行う。
特長と利点
OpenTSDBデータ統合は以下の特長と利点を提供します:
- 効率的なデータ処理:EMQXは大量のIoTデバイス接続とメッセージスループットを処理可能であり、OpenTSDBはデータ書き込み・保存・クエリに優れているため、IoTシナリオのデータ処理要件をシステムに過負荷をかけずに満たします。
- メッセージ変換:EMQXのルールを通じてメッセージは広範囲に処理・変換されてからOpenTSDBに書き込まれます。
- 大規模データ保存:EMQXとOpenTSDBを統合することで、大量のデバイスデータを直接OpenTSDBに保存可能です。OpenTSDBは大規模時系列データの保存とクエリに特化したデータベースであり、IoTデバイスが生成する膨大な時系列データを効率的に処理できます。
- 豊富なクエリ機能:OpenTSDBの最適化されたストレージ構造とインデックスにより、数十億のデータポイントの高速な書き込みとクエリが可能であり、IoTデバイスデータのリアルタイム監視・分析・可視化に非常に有用です。
- スケーラビリティ:EMQXとOpenTSDBはどちらもクラスター拡張に対応しており、ビジネスの成長に応じて柔軟に水平拡張が可能です。
はじめる前に
本節では、OpenTSDBデータ統合の作成に先立ち必要な準備事項を説明します。OpenTSDBサーバーのセットアップ方法も含みます。
前提条件
OpenTSDBのインストール
Dockerを用いてOpenTSDBをインストールし、Dockerイメージを起動します(現在はx86プラットフォームのみ対応)。
docker pull petergrace/opentsdb-docker
docker run -d --name opentsdb -p 4242:4242 petergrace/opentsdb-dockerコネクターの作成
本節では、SinkをOpenTSDBサーバーに接続するためのコネクター作成方法を示します。
以下の手順はEMQXとOpenTSDBを同一マシン上で実行していることを前提としています。リモート環境で実行している場合は設定を適宜調整してください。
- EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。
- ページ右上のCreateをクリックします。
- Create ConnectorページでOpenTSDBを選択し、Nextをクリックします。
- Configurationステップで以下を設定します:
- コネクター名を入力します。英数字の大文字・小文字を組み合わせた名前にしてください(例:
my_opentsdb)。 - Server Hostに
http://127.0.0.1:4242を入力します。OpenTSDBサーバーがリモートの場合は実際のURLを指定してください。 - その他のオプションはデフォルトのままにします。
- コネクター名を入力します。英数字の大文字・小文字を組み合わせた名前にしてください(例:
- 高度な設定(任意):詳細はSinkの特長を参照してください。
- Createをクリックする前に、Test ConnectivityをクリックしてコネクターがOpenTSDBサーバーに接続できるかテストできます。
- ページ下部のCreateボタンをクリックしてコネクター作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてルールとSinkの作成を続行できます。詳細はOpenTSDB Sinkを用いたルール作成を参照してください。
OpenTSDB Sinkを用いたルール作成
本節では、ダッシュボード上でMQTTトピックt/#からのメッセージを処理し、処理済みデータを設定済みSink経由でOpenTSDBに保存するルールの作成方法を示します。
EMQXダッシュボードにアクセスし、Integration -> Rulesをクリックします。
ページ右上のCreateをクリックします。
ルールIDに
my_ruleを入力し、SQL Editorに以下のステートメントを設定します。これはトピックt/#配下のMQTTメッセージをOpenTSDBに保存することを意味します。注意:独自のSQL構文を指定する場合は、Sinkが必要とする全てのフィールドを
SELECT部分に含めていることを確認してください。sqlSELECT payload.metric as metric, payload.tags as tags, payload.value as value FROM "t/#"注意:初心者の方はSQL ExamplesとEnable TestをクリックしてSQLルールの学習とテストを行うことができます。
+ Add Actionボタンをクリックし、ルールによってトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理したデータをOpenTSDBに送信します。
Type of Actionドロップダウンリストから
OpenTSDBを選択します。ActionドロップダウンはデフォルトのCreate Actionのままにします。既に作成済みのSinkがあれば選択可能ですが、本例では新規Sinkを作成します。Sinkの名前を入力します。名前は英数字の大文字・小文字を組み合わせてください。
Connectorドロップダウンから先ほど作成した
my_opentsdbを選択します。新規コネクターを作成する場合はドロップダウン横のボタンをクリックしてください。設定パラメーターはコネクターの作成を参照してください。Write DataフィールドでOpenTSDBに書き込むデータの形式を指定し、MQTTメッセージをOpenTSDBが要求する形式に正しく変換します。例えば、クライアントが以下のデータを報告するとします:
- トピック:
t/opents - ペイロード:
json{ "metric": "cpu", "tags": { "host": "serverA" }, "value": 12 }提供されたペイロードのデータ形式に基づき、以下の形式情報を設定します:
- Timestamp:OpenTSDBはデータポイントの時刻を記録するためにタイムスタンプを必要とします。MQTTメッセージにタイムスタンプが含まれていない場合は、EMQXのSink設定で現在時刻をタイムスタンプとして使用するか、クライアントの報告データ形式を修正してタイムスタンプフィールドを含める必要があります。
- Metric:この例では
"metric": "cpu"がメトリック名cpuを示します。 - Tags:タグはメトリックに関する追加情報を表します。ここでは
"tags": {"host": "serverA"}がこのメトリックデータがホストserverAからのものであることを示しています。 - Value:実際のメトリック値です。この例では
"value": 12でメトリック値が12であることを示します。
- トピック:
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はフォールバックアクションを参照してください。
高度な設定(任意):必要に応じてsyncまたはasyncクエリモードを選択します。詳細はSinkの特長の該当設定情報を参照してください。
Createをクリックする前に、Test ConnectivityをクリックしてSinkがOpenTSDBサーバーに接続できるかテスト可能です。
CreateボタンをクリックしてSink設定を完了します。新しいSinkがAction Outputsに追加されます。
Create Ruleページに戻り、設定内容を確認してCreateボタンをクリックしルールを生成します。
これでOpenTSDB Sinkを通じたデータ転送ルールが正常に作成されました。Integration -> Rulesページで新規作成したルールを確認できます。**Actions(Sink)**タブをクリックすると新しいOpenTSDB Sinkが表示されます。
また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/#配下のメッセージがルールmy_ruleによって解析されOpenTSDBに送信・保存されていることが確認できます。
ルールのテスト
MQTTXを使ってトピックt/opentsにメッセージをパブリッシュします。
mqttx pub -i emqx_c -t t/opents -m '{"metric":"cpu","tags":{"host":"serverA"},"value":12}'Sinkの稼働状況を確認すると、新規の受信メッセージと送信メッセージがそれぞれ1件ずつあるはずです。
OpenTSDBにデータが書き込まれているかを以下のコマンドで確認します:
curl -X POST -H "Accept: Application/json" -H "Content-Type: application/json" http://localhost:4242/api/query -d '{
"start": "1h-ago",
"queries": [
{
"aggregator": "last",
"metric": "cpu",
"tags": {
"host": "*"
}
}
],
"showTSUIDs": "true",
"showQuery": "true",
"delete": "false"
}'クエリ結果の整形済み出力例は以下の通りです:
[
{
"metric": "cpu",
"tags": {
"host": "serverA"
},
"aggregateTags": [],
"query": {
"aggregator": "last",
"metric": "cpu",
"tsuids": null,
"downsample": null,
"rate": false,
"filters": [
{
"tagk": "host",
"filter": "*",
"group_by": true,
"type": "wildcard"
}
],
"percentiles": null,
"index": 0,
"rateOptions": null,
"filterTagKs": [
"AAAB"
],
"explicitTags": false,
"useFuzzyFilter": true,
"preAggregate": false,
"rollupUsage": null,
"rollupTable": "raw",
"showHistogramBuckets": false,
"useMultiGets": true,
"tags": {
"host": "wildcard(*)"
},
"histogramQuery": false
},
"tsuids": [
"000001000001000001"
],
"dps": {
"1683532519": 12
}
}
]