Skip to content

GreptimeDBへのMQTTデータ取り込み ​

GreptimeDBは、スケーラビリティ、分析機能、効率性に特化したオープンソースの時系列データベースです。クラウド時代のインフラ上で動作するよう設計されており、その弾力性と汎用ストレージの利点をユーザーに提供します。EMQXは現在、主流のGreptimeDB、GreptimeCloud、またはGreptimeDB Enterpriseとの接続をサポートしています。

本ページでは、EMQXとGreptimeDB間のデータ統合について包括的に紹介し、データ統合の作成および検証に関する実践的な手順を提供します。

動作の仕組み ​

GreptimeDBデータ統合はEMQXに組み込まれた機能であり、EMQXのリアルタイムデータキャプチャおよび送信能力とGreptimeDBのデータ保存・分析能力を組み合わせています。組み込みのルールエンジンコンポーネントにより、EMQXからGreptimeDBへのデータ取り込みを簡素化し、複雑なコーディングを不要にします。ワークフローは以下の通りです。

下図はEMQXとGreptimeDB間のデータ統合の典型的なアーキテクチャを示しています。

EMQX Integration GreptimeDB

  1. メッセージのパブリッシュと受信:産業用デバイスはMQTTプロトコルを介してEMQXに正常に接続し、定期的にエネルギー消費データをパブリッシュします。このデータには生産ライン識別子やエネルギー消費値が含まれます。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
  2. ルールエンジンによるメッセージ処理:組み込みのルールエンジンはトピックマッチングに基づき特定のソースからのメッセージを処理します。メッセージが到着するとルールエンジンを通過し、対応するルールとマッチングしてメッセージデータを処理します。これにはデータ形式の変換、特定情報のフィルタリング、コンテキスト情報によるメッセージの付加などが含まれます。
  3. GreptimeDBへのデータ取り込み:ルールエンジンで定義されたルールは、メッセージをGreptimeDBに書き込む操作をトリガーします。GreptimeDB SinkはLine Protocolテンプレートを提供し、特定のメッセージフィールドをGreptimeDBの対応テーブルおよびカラムに柔軟に書き込むデータ形式を定義できます。

エネルギー消費データがGreptimeDBに書き込まれた後は、SQL文やPrometheusクエリ言語を使って柔軟にデータ分析が可能です。例えば:

  • Grafanaなどの可視化ツールに接続し、チャートを生成してエネルギー消費データを表示する。
  • ERPなどのアプリケーションシステムに接続し、生産分析や生産計画の調整を行う。
  • ビジネスシステムに接続し、リアルタイムのエネルギー使用状況分析を実施し、データ駆動型のエネルギー管理を促進する。

特長と利点 ​

GreptimeDBとのデータ統合は、以下の特長と利点をビジネスにもたらします。

  • 使いやすさ:EMQXとGreptimeDBはどちらも開発者にとって使いやすい設計です。EMQXは標準のMQTTプロトコルに加え、多様な認証、認可、クラスタリング機能を提供します。GreptimeDBは時系列テーブルやスキーマレスアーキテクチャなどユーザーフレンドリーな設計を持ちます。両者の統合により、ビジネス統合と開発のスピードを加速できます。
  • 効率的なデータ処理:EMQXは多数のIoTデバイス接続とメッセージスループットを効率的に処理可能です。GreptimeDBはデータの書き込み、保存、クエリに優れており、IoTシナリオのデータ処理要件をシステムに過負荷をかけずに満たします。
  • メッセージ変換:EMQXのルール内でメッセージに対して豊富な処理や変換を行い、GreptimeDBに書き込むことができます。
  • 効率的な保存とスケーラビリティ:EMQXとGreptimeDBはどちらもクラスターのスケールアウト機能を持ち、ビジネスの成長に応じて柔軟に水平スケーリングが可能です。
  • 高度なクエリ機能:GreptimeDBはタイムスタンプデータの効率的なクエリと分析のために最適化された関数、演算子、インデックス技術を提供し、IoT時系列データから正確なインサイトを抽出できます。

はじめる前に ​

このセクションでは、GreptimeDBデータ統合の作成を開始する前に必要な準備について説明します。GreptimeDBサーバーのインストール方法も含みます。

前提条件 ​

GreptimeDBサーバーのインストール ​

  1. Dockerを使ってGreptimeDBをインストールし、Dockerイメージを起動します。

    bash
    # GreptimeDBのDockerイメージを起動するコマンド
    docker run -p 127.0.0.1:4000-4003:4000-4003 \
      -v "$(pwd)/greptimedb_data:/greptimedb_data" \
      --name greptime --rm \
      greptime/greptimedb:latest standalone start \
      --http-addr 0.0.0.0:4000 \
      --rpc-bind-addr 0.0.0.0:4001 \
      --mysql-addr 0.0.0.0:4002 \
      --postgres-addr 0.0.0.0:4003 \
      --user-provider=static_user_provider:cmd:greptime_user=greptime_pwd
  2. user-providerパラメータはGreptimeDBの認証を設定します。ファイルによる設定も可能です。詳細はドキュメントを参照してください。

  3. GreptimeDBが起動したら、http://localhost:4000/dashboardにアクセスしてGreptimeDBダッシュボードを利用できます。ユーザー名とパスワードはそれぞれgreptime_userとgreptime_pwdです。

コネクターの作成 ​

このセクションでは、SinkをGreptimeDBサーバーに接続するためのコネクターの作成方法を説明します。

以下の手順は、EMQXとGreptimeDBの両方をローカルマシンで実行していることを前提としています。リモートで実行している場合は設定を適宜調整してください。

  1. EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。
  2. ページ右上のCreateをクリックします。
  3. Create ConnectorページでGreptimeDBを選択し、Nextをクリックします。
  4. Configurationステップで以下の情報を設定します:
    • コネクター名を入力します。英数字の大文字・小文字の組み合わせで、例:my_greptimedb。
    • Server Host:127.0.0.1:4001を入力します。GreptimeCloudに接続する場合はポートを443にして{url}:443と入力してください。
    • Database:publicを入力します。GreptimeCloudの場合はサービス名を入力してください。
    • UsernameとPassword:greptime_userとgreptime_pwdを入力します(GreptimeDBサーバーのインストールで設定したもの)。GreptimeCloud接続時はサービスのユーザー名とパスワードを入力してください。
  5. Advanced Settingsを展開し、必要に応じて詳細設定を行います(任意)。詳細は高度な設定を参照してください。
  6. Createをクリックする前に、Test ConnectivityをクリックしてコネクターがGreptimeDBサーバーに接続できるかテストできます。
  7. ページ下部のCreateボタンをクリックしてコネクターの作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてGreptimeDB Sinkを使ったルール作成に進めます。詳細はGreptimeDB Sinkを使ったルール作成を参照してください。

GreptimeDB Sinkを使ったルール作成 ​

このセクションでは、EMQXでルールを作成し、ソースMQTTトピックt/#からのメッセージを処理して、設定したSink経由でGreptimeDBに送信する方法を説明します。

  1. EMQXダッシュボードにアクセスし、Integration -> Rulesをクリックします。

  2. ページ右上のCreateをクリックします。

  3. ルールIDにmy_ruleを入力し、SQL Editorでルールを設定します。ここではトピックt/#のMQTTメッセージをGreptimeDBに保存するため、以下のSQL文を使用します。

    注意:独自のSQL文を指定する場合は、Sinkが必要とするすべてのフィールドをSELECT句に含めてください。

    sql
    SELECT
      *
    FROM
      "t/#"

    TIP

    初心者の方はSQL Examplesをクリックし、Enable Testを有効にしてSQLルールの学習とテストが可能です。

    • Add Actionボタンをクリックし、ルールによってトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理したデータをGreptimeDBに送信します。
  4. Type of ActionドロップダウンからGreptimeDBを選択します。ActionドロップダウンはデフォルトのCreate Actionのままにします。既に作成したSinkを選択することも可能ですが、ここでは新規Sinkを作成します。

  5. Sinkの名前を入力します。名前は英数字の大文字・小文字の組み合わせで指定してください。

  6. Connectorドロップダウンから先ほど作成したmy_greptimedbを選択します。ドロップダウン横のボタンから新規コネクター作成も可能です。設定パラメータはコネクターの作成を参照してください。

  7. Write Syntaxを設定します。これは、データポイントの計測名、タグ、フィールド、タイムスタンプを指定するテキストベースのフォーマットで、InfluxDB line protocolの構文に準拠したプレースホルダーをサポートしています。GreptimeDBはInfluxDB互換のデータフォーマットをサポートしています。

    TIP

    • GreptimeDBに符号付き整数型の値を書き込む場合は、プレースホルダーの後にiを付けます。例:${payload.int}i
    • 符号なし整数型の場合はuを付けます。例:${payload.int}u
  8. Time Precisionを指定します。デフォルトはmillisecondです。

  9. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はフォールバックアクションを参照してください。

  10. 高度な設定(任意):同期(sync)または非同期(async)クエリモードの選択、キューやバッチの有効化設定を行えます。詳細はSinkの機能を参照してください。

  11. Createをクリックする前に、Test ConnectivityをクリックしてSinkがGreptimeDBサーバーに接続できるかテスト可能です。

  12. CreateボタンをクリックしてSinkの設定を完了します。新しいSinkがAction Outputsに追加されます。

  13. Create Ruleページに戻り、設定内容を確認してCreateをクリックしルールを生成します。

これでGreptimeDB Sinkを経由してデータを転送するルールの作成が完了しました。Integration -> Rulesページで新規作成したルールを確認できます。**Actions(Sink)**タブをクリックすると、新しいGreptimeDB Sinkが表示されます。

また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/#のメッセージがルールmy_ruleによって解析されGreptimeDBに送信・保存されている様子を確認できます。

ルールのテスト ​

MQTTクライアントのMQTTXを使ってトピックt/1にメッセージを送信し、オンライン/オフラインイベントをトリガーします。

bash
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello GreptimeDB" }'

Sinkの稼働状況を確認すると、新規の受信メッセージと送信メッセージがそれぞれ1件ずつあるはずです。

GreptimeDBダッシュボードでSQLを使い、メッセージがGreptimeDBに書き込まれていることを確認できます。

高度な設定 ​

このセクションでは、コネクターのパフォーマンスを最適化し、特定のシナリオに合わせて動作をカスタマイズするための高度な設定オプションを説明します。コネクター作成時にAdvanced Settingsを展開し、ビジネスニーズに応じて以下の設定を行えます。

フィールド名説明デフォルト値
Time-To-Live (TTL)GreptimeDBで自動作成されるテーブルの有効期限設定。-
Custom Timestamp Column Name定義すると、クエリ時に表示されるカスタムタイムスタンプカラム名を指定。-
Start Timeoutコネクターが自動起動したリソースの正常状態到達を待つ最大時間(秒)。リソース作成要求に応答する前に、接続先リソースが完全に稼働しデータ処理可能であることを確認するための設定。5秒
Health Check Intervalコネクターの稼働状態をチェックする間隔時間。15秒
Health Check TimeoutGreptimeDBサーバーとの接続に対する自動ヘルスチェックのタイムアウト時間。60秒