CockroachDBへのMQTTデータ取り込み
CockroachDBは、分散型でPostgreSQL互換のデータベースであり、フルマネージドクラウドサービス(CockroachDB Cloud)またはセルフホスト型のデプロイメントとして利用可能です。高いレジリエンス、水平スケーラビリティ、および完全なSQL互換性を必要とするグローバルアプリケーション向けに設計されています。EMQXはCockroachDBとスムーズに統合し、IoTデバイスからのMQTTデータをリアルタイムでキャプチャして保存します。これにより、グローバル展開における高速かつ信頼性の高いデータ取り込み、Raftベースのレプリケーションによる一貫性のあるデータ管理、および運用や分析のための低レイテンシ読み取りを実現します。
本ページでは、EMQXとCockroachDB間のデータ統合について包括的に紹介し、データ統合の作成および検証に関する実践的な手順を提供します。
動作概要
EMQXにおけるCockroachDBデータ統合は、MQTTベースのIoTデータストリームをCockroachDBの分散型PostgreSQL互換データベースに直接取り込む組み込み機能です。EMQXの組み込みルールエンジンを利用することで、複雑なカスタムコードを書くことなく、CockroachDBへ直接データを取り込み、グローバルに一貫したストレージとリアルタイムクエリを実現できます。
CockroachDBの共有なし(shared-nothing)分散アーキテクチャは、Raftベースのコンセンサスを用いて複数のノードやリージョンにデータを自動的にレプリケートし、障害時でも強い一貫性を維持します。これにより、IoTデータは常に安全かつ同期され、利用可能な状態が保たれます。
以下の図は、EMQXとCockroachDB間のデータ統合の典型的なアーキテクチャを示しています。

CockroachDBへのMQTTデータ取り込みは以下のように動作します。
- IoTデバイスがEMQXに接続:IoTデバイスがMQTTプロトコルを通じて正常に接続されると、オンラインイベントがトリガーされます。イベントにはデバイスID、送信元IPアドレスなどの情報が含まれます。
- メッセージのパブリッシュと受信:デバイスは特定のトピックにテレメトリおよびステータスデータをパブリッシュします。EMQXはこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
- ルールエンジンによるメッセージ処理:EMQXのルールエンジンは、トピックやメッセージ内容に基づいて定義されたルールにマッチさせ、イベントやメッセージを処理します。処理にはデータ変換(例:JSONからSQL用フォーマットへの変換)、フィルタリング、コンテキスト情報によるデータ強化などが含まれ、データベースへの挿入準備を行います。
- CockroachDBへの書き込み:マッチしたルールがトリガーされ、CockroachDBに対してSQLが実行されます。SQLテンプレートを用いて、処理済みデータのフィールドをCockroachDBのテーブルやカラムにマッピングします。CockroachDBの分散SQL実行およびベクトル化クエリエンジンにより、高スループットの書き込みと低レイテンシの分析クエリが可能です。マルチリージョン展開では、ジオパーティショニングによる最適化も可能です。
イベントおよびメッセージデータがCockroachDBに書き込まれた後は、以下のような活用が可能です。
- CockroachDBをGrafanaなどのツールに接続し、ライブのIoTメトリクスを表示するダッシュボードやチャートを作成。
- デバイス管理プラットフォームやAI/MLモデルと連携し、ヘルスチェック、異常検知、アラートトリガーを実現。
- CockroachDBの分散クエリエンジンを用いて、ライブのIoTデータに対する集計、結合、時系列分析を実行しつつ、新しいテレメトリの処理を並行して継続。
特長とメリット
CockroachDBとのデータ統合により、以下の特長と利点が得られます。
- 柔軟なイベント処理:EMQXのルールエンジンを利用して、CockroachDBにデバイスのライフサイクルイベント(接続、切断、ステータス変更)を低レイテンシで保存・処理可能です。CockroachDBの分散実行と自動リバランシングにより、イベントデータは高可用性を保ち、リアルタイム分析で障害や異常、トレンド検出に活用できます。
- メッセージ変換:メッセージはEMQXルールを通じて高度に処理・変換されてからCockroachDBに書き込まれるため、保存データは最初から分析に適した形となります。これによりクエリの複雑さが軽減され、下流処理が最適化されます。
- SQLテンプレートによる柔軟なデータ操作:EMQXのSQLテンプレートマッピングを使い、構造化されたIoTデータをCockroachDBのテーブルやカラムに挿入・更新可能です。PostgreSQL互換のため、標準SQL、JSONBストレージ、インデックス作成が利用でき、ベクトル化実行エンジンによる高速分析やフォロワーリードによる低レイテンシなリージョンローカルアクセスが可能です。
- 業務プロセスとの統合:CockroachDBのPostgreSQL互換性により、ERP、CRM、GISなどの業務システムと統合可能です。EMQXと組み合わせることで、複雑なETLパイプラインを構築せずにイベント駆動型の自動化やクロスシステムオーケストレーションを実現できます。
- 高度な地理空間機能:PostGISなどのPostgreSQL拡張を通じて、CockroachDBは地理空間データの保存、インデックス作成、クエリをサポートします。これにより、ジオフェンシング、位置ベースのアラート、ルート追跡、リアルタイム資産監視がEMQXの信頼性の高いIoTデータ取り込みと連携して可能となります。
- 組み込みのメトリクスと監視:EMQXは各CockroachDBシンクのランタイムメトリクス(メッセージ数、成功/失敗率、スループット)を提供し、CockroachDBは組み込みの可観測性ツールを備え、PrometheusやGrafanaと連携して詳細なパフォーマンスおよびヘルス監視を実現します。
はじめる前に
このセクションでは、CockroachDB統合の作成を開始する前に必要な準備について説明します。CockroachDBのデプロイメントやデータベースおよびテーブルの作成方法を含みます。
前提条件
CockroachDBでのデータベースおよびテーブル作成
EMQXでCockroachDBコネクターを作成する前に、CockroachDBクラスターが稼働していること、およびIoTデータを保存するためのデータベースとテーブルが準備されていることを確認してください。
CockroachDBクラスターを作成します。
- CockroachDB Cloudの場合は、CockroachDB Cloudドキュメントに従ってクラスターをプロビジョニングしてください。
- セルフホスト型の場合は、インストールガイドに従ってください。
EMQX用の専用SQLユーザーを作成します。詳細はCockroachDBユーザー管理ガイドを参照してください。ここでは例として
emqx_userという名前のSQLユーザーを使用します。このユーザーには以下の権限が必要です。- 対象データベースへの接続権限
- テーブル作成権限
- EMQXデータテーブルへの読み書き権限
データベースの作成に従い、データベースを作成します。例としてデータベース名は
emqx_dataとします。emqx_dataデータベースに接続し、MQTTメッセージとクライアントイベントデータを格納する2つのテーブルを作成します。テーブル作成の手順に従ってください。以下のSQL文で、クライアントID、トピック、QoS、ペイロード、到着時間などのメタデータを含むMQTTメッセージを格納する
t_mqtt_msgテーブルを作成します。sqlCREATE TABLE t_mqtt_msg ( id SERIAL primary key, msgid character varying(64), sender character varying(64), topic character varying(255), qos integer, retain integer, payload text, arrived timestamp without time zone );以下のSQL文で、クライアントのオンライン/オフラインイベントをタイムスタンプ付きで格納する
emqx_client_eventsテーブルを作成します。sqlCREATE TABLE emqx_client_events ( id SERIAL primary key, clientid VARCHAR(255), event VARCHAR(255), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
CockroachDBコネクターの作成
CockroachDBシンクを追加する前に、EMQXでCockroachDBコネクターを作成する必要があります。コネクターは、EMQXがセルフホスト型またはCockroachDB Cloudのクラスターに接続する方法を定義します。
EMQXダッシュボードで、Integration -> Connector に移動します。
ページ右上の Create をクリックします。
Create Connector ページで CockroachDB を選択し、Next をクリックします。
コネクター名を入力します。名前は英数字で始まり、英数字、ハイフン、アンダースコアを含めることができます。例:
my_cockroachdb接続情報を入力します。
- Server Host:CockroachDBクラスターのホスト名またはIPアドレス
- CockroachDB Cloud:CockroachDB Cloudコンソールで提供される接続文字列のホスト値を使用(例:
free-tier.gcp-us-central1.cockroachlabs.cloud) - セルフホスト型:CockroachDBが稼働しているアドレスを使用(例:ローカルは
127.0.0.1、またはサーバーのパブリック/プライベートIP)
- CockroachDB Cloud:CockroachDB Cloudコンソールで提供される接続文字列のホスト値を使用(例:
- Database Name:EMQXがデータを保存する対象データベース名。例:
emqx_data - Username:認証および識別に使用するCockroachDBのSQLユーザー名。例:
emqx_user - Password:
emqx_userのパスワード - Enable TLS:暗号化接続を確立する場合はトグルスイッチをオンにします。TLS接続の詳細は外部リソースアクセスのTLSを参照してください。
- Server Host:CockroachDBクラスターのホスト名またはIPアドレス
詳細設定(任意):接続プールサイズ、アイドルタイムアウト、リクエストタイムアウトなどの追加接続プロパティを設定できます。詳細はシンクの機能を参照してください。
Test Connectivity をクリックして、EMQXが指定した設定でCockroachDBクラスターに正常に接続できるかを確認します。
Create をクリックしてコネクターを保存します。
作成後は以下のいずれかを選択できます。
- Back to Connector List をクリックして全コネクターを表示
- Create Rule をクリックして、このコネクターを使ったデータ転送ルールを即座に作成
詳細な例は以下を参照してください。
メッセージ保存用CockroachDBシンクのルール作成
このセクションでは、ダッシュボードでソースMQTTトピックt/#からのメッセージを処理し、処理済みデータを設定済みのCockroachDBシンクを介してt_mqtt_msgテーブルに保存するルールの作成方法を示します。
ダッシュボードの Integration -> Rules ページに移動します。
ページ右上の Create をクリックします。
ルールIDに
my_ruleを入力し、SQLエディターにルールを記述します。ここでは、トピックt/#のMQTTメッセージをCockroachDBに保存するため、ルールのSELECT句でSQLテンプレートで使用するすべての変数を含むフィールドを選択してください。ルールSQLは以下の通りです。sqlSELECT * FROM "t/#"TIP
初心者の方は、SQL Examplesをクリックし、Enable Testを有効にしてSQLルールの学習とテストを行うことができます。
- Add Action ボタンをクリックして、ルールによりトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理したデータをCockroachDBに送信します。
Type of Action ドロップダウンからCockroachDBを選択し、Action ドロップダウンはデフォルトの
Create Actionのままにするか、既存のCockroachDBアクションを選択します。この例では新規シンクを作成しルールに追加します。シンクの名前と説明を入力します。
Connector ドロップダウンから先ほど作成した
my_cockroachdbを選択します。新しいコネクターを作成する場合は、ドロップダウン横のボタンをクリックしてください。設定パラメーターはCockroachDBコネクターの作成を参照してください。SQLテンプレートを設定します。以下のSQL文を使ってデータを挿入します。
※これはプリペアドステートメントのため、フィールドは引用符で囲まず、文末にセミコロンを付けないでください。
sqlINSERT INTO t_mqtt_msg(msgid, sender, topic, qos, payload, arrived) VALUES( ${id}, ${clientid}, ${topic}, ${qos}, ${payload}, TO_TIMESTAMP((${timestamp} :: bigint)/1000) )フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。詳細はフォールバックアクションを参照してください。
詳細設定(任意):詳細はシンクの機能を参照してください。
Createをクリックする前に、Test ConnectivityをクリックしてシンクがCockroachDBクラスターに接続できるかをテストできます。
Createボタンをクリックしてシンク設定を完了します。新しいシンクがAction Outputsに追加されます。
Create Ruleページで設定内容を確認し、Saveをクリックしてルールを生成します。
ルール作成後、Integration -> Rules ページで新規ルールを確認でき、Action (Sink) タブで新規CockroachDBシンクも確認できます。
また、Integration -> Flow Designer でトポロジーを確認でき、トピックt/#のメッセージがルールmy_ruleで解析されてCockroachDBに書き込まれている様子を可視化できます。
イベント記録用CockroachDBシンクのルール作成
このセクションでは、クライアントのオンライン/オフライン状態を記録し、イベントデータを設定済みのCockroachDBシンクを介してemqx_client_eventsテーブルに保存するルールの作成方法を示します。
手順はメッセージ保存用CockroachDBシンクのルール作成と同様ですが、SQLテンプレートとSQLルールが異なります。
オンライン/オフライン状態記録用のSQLルール文は以下の通りです。
SELECT
*
FROM
"$events/client_connected", "$events/client_disconnected"イベント記録用のSQLテンプレートは以下の通りです。
※これはプリペアドステートメントのため、フィールドは引用符で囲まず、文末にセミコロンを付けないでください。
INSERT INTO emqx_client_events(clientid, event, created_at) VALUES (
${clientid},
${event},
TO_TIMESTAMP((${timestamp} :: bigint)/1000)
)ルールのテスト
MQTTXを使ってトピックt/1にメッセージを送信し、オンライン/オフラインイベントをトリガーします。
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello CockroachDB" }'2つのシンクの稼働状況を確認します。メッセージ保存用シンクでは新規の受信メッセージ1件と送信メッセージ1件があるはずです。イベント記録用シンクでは2件のイベントレコードが存在します。
t_mqtt_msgデータテーブルにデータが書き込まれているか確認します。
emqx_data=# select * from t_mqtt_msg;
id | msgid | sender | topic | qos | retain | payload
| arrived
----+----------------------------------+--------+-------+-----+--------+-------------------------------+---------------------
1 | 0005F298A0F0AEE2F443000012DC0002 | emqx_c | t/1 | 0 | | { "msg": "hello CockroachDB" } | 2023-01-19 07:10:32
(1 row)emqx_client_eventsテーブルにデータが書き込まれているか確認します。
emqx_data=# select * from emqx_client_events;
id | clientid | event | created_at
----+----------+---------------------+---------------------
3 | emqx_c | client.connected | 2023-01-19 07:10:32
4 | emqx_c | client.disconnected | 2023-01-19 07:10:32
(2 rows)