MQTTデータをTablestoreに取り込む
Tablestoreは、IoTシナリオに最適化されたスケーラブルでサーバレスなデータベースです。時系列データ、構造化データ、半構造化データを管理するためのワンストップソリューションであるIoTstoreを提供しています。IoT、車載ネットワーク、リスク管理、メッセージング、レコメンデーションシステムなどのシナリオに最適です。Tablestoreはコスト効率が高く高性能なデータストレージを提供し、ミリ秒単位のクエリや検索、柔軟なデータ分析機能を備えています。EMQXはTablestore Cloud、Tablestore OSS、Tablestore Enterpriseとシームレスに統合し、IoTユースケースにおける効率的なデータ管理を実現します。
動作概要
EMQXのTablestoreデータ統合は、EMQXのリアルタイムデータキャプチャと送信機能と、Tablestoreの高性能なデータストレージおよび分析機能をシームレスに組み合わせています。組み込みのルールエンジンを活用することで、EMQXからTablestoreへのデータ取り込みと保存のプロセスを簡素化し、複雑なコーディングを不要にします。EMQXはルールエンジンとSinkを通じてIoTデバイスのデータをTablestoreに転送し、効率的な保存と分析を可能にします。
データが保存されると、Tablestoreはレポートやチャート、その他の可視化を生成する強力なツールを提供し、これらはTablestoreの可視化機能を通じてユーザーに提示されます。
以下の図は、エネルギー蓄電シナリオにおけるEMQXとTablestore間の典型的なデータ統合アーキテクチャを示しています。

EMQXとTablestoreは、エネルギー消費データをリアルタイムに効率よく収集・分析するための拡張可能なIoTプラットフォームを提供します。このアーキテクチャでは、EMQXがデバイス接続、メッセージ送信、データルーティングを担当するIoTプラットフォームとして機能し、Tablestoreがデータ保存および分析プラットフォームとしてデータの保存と分析機能を担います。ワークフローは以下の通りです。
- メッセージのパブリッシュと受信:蓄電デバイスや産業用IoTデバイスはMQTTプロトコルを通じてEMQXに正常に接続し、電力消費量、入出力電力などの情報を含むエネルギー消費データを定期的にMQTTプロトコルでパブリッシュします。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
- メッセージデータの処理:組み込みのルールエンジンを使用して、特定のソースからのメッセージをトピックマッチングに基づいて処理できます。メッセージが到着するとルールエンジンを通過し、対応するルールとマッチングしてメッセージデータを処理します。例えば、データ形式の変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの拡充などです。
- Tablestoreへのデータ取り込み:ルールエンジンで定義されたルールがトリガーとなり、メッセージをTablestoreに書き込む操作が実行されます。Tablestore Sinkは設定可能なフィールドを提供し、書き込むデータ形式を柔軟に定義でき、メッセージの特定フィールドをTablestoreの対応するメジャメントやフィールドにマッピングします。
エネルギー消費データがTablestoreに書き込まれた後、以下のようにデータを分析できます。
- Grafanaなどの可視化ツールに接続し、データに基づくチャートを生成して蓄電データを表示する。
- ビジネスシステムに接続し、蓄電デバイスの状態監視やアラートを実施する。
特長と利点
Tablestoreデータ統合は以下の特長と利点を提供します。
- 効率的なデータ処理:EMQXは大量のIoTデバイス接続とメッセージスループットを処理可能であり、Tablestoreはデータ書き込み、保存、クエリに優れた性能を発揮します。IoTシナリオのデータ処理ニーズをシステムに過度な負荷をかけずに満たします。
- メッセージ変換:メッセージはEMQXのルールを通じて幅広い処理や変換が可能であり、その後Tablestoreに書き込まれます。
- スケーラビリティ:EMQXおよびTablestoreはクラスターのスケールアウトに対応しており、ビジネスの成長に応じて柔軟に水平拡張できます。
- 豊富なクエリ機能:Tablestoreは最適化された関数、演算子、インデックス技術を提供し、タイムスタンプ付きデータの効率的なクエリと分析を可能にし、IoT時系列データから価値ある洞察を正確に抽出します。
- 効率的なストレージ:Tablestoreは高圧縮率のエンコーディング方式を採用し、ストレージコストを大幅に削減します。また、異なるデータタイプごとに保存期間をカスタマイズでき、不必要なデータがストレージを占有するのを防ぎます。
はじめる前に
このセクションでは、Tablestoreデータ統合を作成する前に完了すべき準備について説明します。データベースインスタンスの作成、時系列テーブルの作成と管理が含まれます。
TIP
現在、Tablestoreとのデータ統合はTimeSeriesモデルのみをサポートしています。したがって、以下の手順はTimeSeriesモデルのデータ統合に焦点を当てています。
前提条件
以下を事前にご確認ください。
時系列テーブルの作成
- Tablestoreコンソールにログインします。
- 時系列モデルのインスタンスを作成します。インスタンス名には例えば
emqx-demoを指定します。インスタンス作成の詳細はTablestore公式ドキュメントを参照してください。 - インスタンス管理ページに移動します。
- インスタンス詳細タブで時系列テーブルを選択し、時系列テーブルの作成ボタンをクリックします。
- 時系列テーブル情報を設定し、テーブル名に
timeseries_demo_with_dataなどを入力して確認をクリックします。

時系列テーブルの管理
作成した時系列テーブルを管理するには、テーブル名をクリックして時系列テーブル管理ページに入ります。ビジネス要件に応じて以下の操作を行えます。
データクエリタブをクリックします。
時系列の追加をクリックします。
TIP
このステップは任意です。時系列テーブルがまだ存在しない場合、データ書き込み時にTablestoreが自動的に作成します。そのため、この例では時系列の手動操作は示していません。

コネクターの作成
このセクションでは、SinkをTablestoreサーバーに接続するためのコネクターの作成方法を示します。
以下の手順はEMQXとTablestoreをローカルマシンで実行していることを前提としています。リモートで実行している場合は設定を適宜調整してください。
- EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。
- ページ右上のCreateをクリックします。
- Create ConnectorページでTablestoreを選択し、Nextをクリックします。
- Configurationステップで以下を設定します。
- コネクター名を入力します。英数字の組み合わせで、例:
my_tablestore。 - Tablestoreサーバー接続情報を入力します。
- Endpoint:TablestoreインスタンスのアクセスURLを入力します。Tablestoreコンソールのインスタンス詳細ページで確認可能です。公開ネットワークの場合は例として
https://emqx-demo.cn-hangzhou.ots.aliyuncs.comのように入力します。 - Instance Name:接続するTablestoreインスタンス名。ここでは作成済みの
emqx-demoを使用します。 - Access Key ID:Tablestore認証に使用するアクセスキーID。Alibaba Cloudが発行するキーです。
- Access Key Secret:アクセスキーIDに対応するシークレットキー。
- Storage Model Type:現在は
TimeSeriesのみサポート。
- Endpoint:TablestoreインスタンスのアクセスURLを入力します。Tablestoreコンソールのインスタンス詳細ページで確認可能です。公開ネットワークの場合は例として
- TLSパラメータを設定します。TablestoreはHTTPSエンドポイントを使用するためTLSはデフォルトで有効です。追加のTLS設定は不要です。TLS接続オプションの詳細は外部リソースアクセスのTLS有効化を参照してください。
- コネクター名を入力します。英数字の組み合わせで、例:
- Createをクリックする前に、Test ConnectivityをクリックしてコネクターがTablestoreサーバーに接続できるかテストできます。
- ページ下部のCreateをクリックしてコネクター作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてルールとSinkの作成に進めます。詳細はTablestore Sinkを使ったルール作成を参照してください。
Tablestore Sinkを使ったルール作成
このセクションでは、EMQXでソースMQTTトピックt/#からのメッセージを処理し、設定済みのSinkを通じてTablestoreに送信するルールの作成方法を示します。
EMQXダッシュボードで左メニューのIntegration -> Rulesをクリックします。
ページ右上のCreateをクリックします。
ルール作成ページでルールIDに
my_ruleを入力します。SQL Editorでルールを設定します。例えば、トピック
t/#のMQTTメッセージをTablestoreに保存したい場合、以下のSQL構文を使用します。TIP
独自のSQL構文を指定する場合は、後で設定するSinkのデータ形式に含まれるすべての変数が
SELECT部分に含まれていることを確認してください。sqlSELECT * FROM "t/#"注:初心者の方はSQL ExamplesとEnable TestをクリックしてSQLルールの学習とテストができます。
- Add Actionボタンをクリックしてルールがトリガーするアクションを定義します。このアクションにより、EMQXはルールで処理したデータをTablestoreに送信します。
Type of Actionドロップダウンリストから
Alibaba Tablestoreを選択します。ActionはデフォルトのCreate Actionのままにします。既に作成済みのSinkがあれば選択可能です。この例では新規Sinkを作成します。Sinkの名前を入力します。英数字の組み合わせで指定してください。
Connectorドロップダウンから先ほど作成した
my_tablestoreを選択します。新規コネクターを作成する場合はドロップダウン横のボタンをクリックします。設定パラメータはコネクター作成を参照してください。以下のフィールドを設定します。
Data Source:EMQXがメッセージを取得するデータソース。処理対象データの起点を表します。特定のトピックやデータストリームを指定できます。
Table Name:データを保存するTablestoreのテーブル名。事前に作成したテーブル名を入力します。
${table}などの変数を使って動的に指定することも可能です。Measurement:Tablestoreで使用するメジャメント名。通常はデータの論理的なグループやカテゴリを示します。例:
temperature_readingsやsensor_data。${measurement}などの変数も使用可能です。Storage Model Type:Tablestoreで使用するデータストレージモデルの種類。現在は時系列データに最適化された
timeseriesのみサポートしています。Tags:Tablestoreの各データエントリに関連付けるキーと値のペア。メタデータやラベルとして利用し、クエリやフィルタリングを容易にします。Addをクリックして複数のタグを定義できます。例:
Key Value locationoffice1devicesensor1Fields:Tablestoreに送信するデータのフィールドリスト。各フィールドはTablestoreテーブルのカラムにマッピングされます。Addをクリックして以下を追加します。
- Column:Tablestoreのカラム名。
${column_name}などの変数を使って定義可能で、後述のメッセージペイロードのフィールドと一致させます。 - Message value:カラムに割り当てる値。
${value}のような動的参照、trueのようなブール値、1.3のような数値、バイナリデータなどが指定可能です。 - Is Int:カラムが数値型の場合、デフォルトでは浮動小数点型としてTablestoreに挿入されます。整数値として挿入したい場合はこのフラグを
trueに設定します。設定ファイル経由の場合、${isint}のような変数で動的に指定可能です。 - Is Binary:カラムがバイナリ型の場合、デフォルトでは文字列型として挿入されます。バイナリデータとして挿入するにはこのフラグを
trueに設定します。設定ファイル経由の場合、${isbinary}のような変数で動的に指定可能です。
- Column:Tablestoreのカラム名。
Timestamp:Tablestoreに記録するタイムスタンプ。マイクロ秒単位の整数値で指定します。固定値を指定するか、文字列"NOW"を使ってEMQXがメッセージ処理時に現在時刻を動的に埋め込むこともできます。
${microsecond_timestamp}のような変数も使用可能です。Meta Update Model:Tablestoreのメタデータ更新戦略を定義します。
MUM_IGNORE:メタデータの更新を無視し、競合があっても変更しません。MUM_NORMAL:通常のメタデータ更新を行います。メタデータが存在しない場合は動的に作成してからデータを書き込みます。既存メタデータと競合する場合は上書きされる可能性があります。
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。プライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はフォールバックアクションを参照してください。
詳細設定(任意):詳細設定を参照してください。
Createをクリックする前に、Test ConnectivityをクリックしてSinkがTablestoreサーバーに接続できるかテストできます。
CreateをクリックしてSink作成を完了します。ルール作成ページのAction Outputsタブに新しいSinkが表示されます。
ルール作成ページで設定内容を確認し、Createをクリックしてルールを生成します。
これでルールが正常に作成され、Ruleページに新しいルールが表示されます。**Actions(Sink)**タブをクリックすると、新しいTablestore Sinkが確認できます。
また、Integration -> Flow Designerをクリックしてトポロジーを確認できます。トピックt/#のメッセージがルールmy_ruleで解析され、Tablestoreに送信・保存されていることが分かります。
ルールのテスト
MQTTXを使ってトピック
t/1にメッセージを送信し、オンライン/オフラインイベントをトリガーします。bashmqttx pub -i emqx_c -t t/1 -m '{ "table": "timeseries_demo_with_data", "measurement": "foo", "microsecond_timestamp": 1734924039271024, "column_name": "cc", "value": 1}'Sinkの稼働状況を確認し、新しい受信メッセージと送信メッセージが1件ずつあることを確認します。
Tablestoreコンソールにアクセスし、データがTablestoreに書き込まれているか確認します。
- Metric Nameにメジャメント名(このデモでは
foo)を入力します。 - Tagに
location=office1とdevice=sensor1をクエリ条件として入力し、Searchをクリックします。

- Metric Nameにメジャメント名(このデモでは
詳細設定
このセクションでは、TablestoreコネクターおよびSinkの詳細な設定オプションについて説明します。ダッシュボードでコネクターやSinkを設定する際、Advanced Settingsに移動して以下のパラメータをニーズに合わせて調整できます。
| フィールド | 説明 | 推奨値 |
|---|---|---|
| Buffer Pool Size | EMQXとTablestore間のegressタイプのブリッジでデータフロー管理に割り当てるバッファワーカープロセス数を指定します。これらのワーカーはデータ送信前に一時的にデータを保持・処理します。egress(送信)シナリオのパフォーマンス最適化に関連します。Ingress(受信)データのみ扱うSinkの場合は"0"に設定可能です。 | 16 |
| Request TTL | バッファに入ったリクエストが有効とみなされる最大時間(秒)を指定します。リクエストがTTLを超えてバッファに残るか、Tablestoreからの応答・アックがタイムリーに得られない場合、リクエストは期限切れと見なされます。 | 45 |
| Health Check Interval | SinkがTablestoreへの接続状態を自動的にヘルスチェックする間隔(秒)を指定します。 | 15 |
| Max Buffer Queue Size | Tablestore Sinkの各バッファワーカーがバッファリング可能な最大バイト数を指定します。バッファワーカーはデータ送信前に一時的にデータを保持し、データフローを効率化します。システムの性能やデータ転送要件に応じて調整してください。 | 256 |
| Batch Size | EMQXからTablestoreへ一度に転送可能なデータバッチのサイズを指定します。サイズ調整によりデータ転送の効率と性能を最適化できます。 | 1 |
| Query Mode | メッセージ送信を最適化するためにasynchronous(非同期)またはsynchronous(同期)クエリモードを選択できます。非同期モードではTablestoreへの書き込みがMQTTメッセージのパブリッシュ処理をブロックしませんが、クライアントがTablestore到着前にメッセージを受信する可能性があります。 | Async |
| Inflight Window | "in-flight query"は開始されたがまだ応答やアックを受け取っていないクエリを指します。SinkがTablestoreと通信する際に同時に存在可能なin-flightクエリの最大数を制御します。 Query Modeが asyncの場合、このパラメータは特に重要です。同一MQTTクライアントからのメッセージを厳密に順序処理したい場合は1に設定してください。 | 100 |