Skip to content

Apache PulsarへのMQTTデータストリーミング

Apache Pulsarは、アプリケーションやシステム間でリアルタイムデータストリームを効率的に送受信するために設計された、人気のあるオープンソースの分散イベントストリーミングプラットフォームです。Apache Pulsarは、より高いスケーラビリティ、より高速なスループット、そして低いレイテンシを提供します。IoTアプリケーションでは、デバイスから生成されるデータは通常、軽量なMQTTプロトコルを用いて送信されます。Apache PulsarとEMQXのデータ統合により、ユーザーはMQTTデータを簡単にApache Pulsarへストリーミングし、他のデータシステムと連携してIoTデバイスから生成されるデータのリアルタイム処理、保存、分析を行うことが可能になります。

本ページでは、EMQXとPulsar間のデータ統合の詳細な概要と、データ統合の作成および検証に関する実践的な手順を提供します。

動作の仕組み

Apache Pulsarデータ統合は、EMQXの標準機能として提供されており、EMQXのデバイス接続およびメッセージ送信機能と、Pulsarの強力なデータ処理機能を組み合わせています。組み込みのルールエンジンコンポーネントにより、両プラットフォーム間のデータストリーミングと処理のプロセスが簡素化されています。これにより、複雑なコーディングを必要とせずにMQTTデータをPulsarに送信し、Pulsarの強力なデータ処理機能を活用できるため、IoTデータの管理と活用がより効率的かつ便利になります。

EMQXのデータ統合 - Apache Pulsar

EMQXはルールエンジンと設定済みのSinkを通じてMQTTデータをApache Pulsarに転送し、その全体の流れは以下の通りです:

  1. メッセージのパブリッシュと受信:IoTデバイスはMQTTプロトコルを介して正常に接続を確立し、特定のトピックにテレメトリおよびステータスデータをパブリッシュします。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
  2. ルールエンジンによるメッセージ処理:組み込みのルールエンジンを使用して、特定のソースからのMQTTメッセージをトピックマッチングに基づいて処理します。ルールエンジンは対応するルールをマッチングし、データ形式の変換、特定情報のフィルタリング、コンテキスト情報の付加などの処理を行います。
  3. Apache Pulsarへのデータストリーミング:ルールがトリガーされると、メッセージをPulsarに転送するアクションが実行されます。データはPulsarのメッセージキーおよび値に簡単にマッピング可能です。MQTTトピックはPulsarトピックにマッピングすることもでき、データの整理や識別が容易になり、後続のデータ処理や分析を促進します。

MQTTメッセージデータがApache Pulsarに書き込まれた後は、以下のような柔軟なアプリケーション開発が可能です:

  • Pulsarのコンシューマーアプリケーションを作成し、これらのメッセージをサブスクライブして処理します。ビジネスニーズに応じて、MQTTデータを他のデータソースと関連付けたり集約したり変換したりして、リアルタイムのデータ同期と統合を実現できます。
  • 特定のMQTTメッセージを受信した際に、Pulsarのルールエンジンコンポーネントを用いて対応するアクションやイベントをトリガーし、システム間やアプリケーション間のイベント駆動機能を実現します。
  • Pulsar内でMQTTデータストリームをリアルタイムに分析し、異常や特定のイベントパターンを検出してアラート通知を行ったり、条件に応じた対応アクションを実行したりします。
  • 複数のMQTTトピックからのデータを統合し、Pulsarの計算能力を活用してリアルタイムの集約、計算、分析を行い、より包括的なデータインサイトを得ます。

特長と利点

Pulsarとのデータ統合により、以下の特長とメリットがビジネスにもたらされます:

  • 信頼性の高いIoTデータメッセージ配信:EMQXはMQTTメッセージを確実にバッチ送信でき、IoTデバイスとPulsarおよびアプリケーションシステムの統合を実現します。
  • MQTTメッセージ変換:ルールエンジンを用いて、EMQXはMQTTメッセージの抽出、フィルタリング、付加情報の追加、変換を行い、Pulsarへ送信します。
  • 柔軟なトピックマッピング:Pulsar SinkはMQTTトピックをPulsarトピックに柔軟にマッピングでき、Pulsarメッセージのキー(Key)や値(Value)の設定も容易です。
  • 柔軟なパーティション選択:Pulsar SinkはMQTTトピックやクライアントに基づき、異なる戦略でPulsarのパーティションを選択可能で、データの整理と識別に柔軟性を持たせます。
  • 高スループットシナリオでの処理能力:Pulsar Sinkは同期・非同期の両方の書き込みモードをサポートし、シナリオに応じてレイテンシとスループットのバランスを柔軟に調整できます。

はじめる前に

このセクションでは、EMQXダッシュボードでPulsarデータ統合を作成する前に必要な準備について説明します。

前提条件

Pulsarのインストール

DockerでPulsarを起動します。

bash
docker run --rm -it -p 6650:6650 --name pulsar apachepulsar/pulsar:2.11.0 bin/pulsar standalone -nfw -nss

詳細な操作手順は、Pulsarドキュメントのクイックスタートセクションを参照してください。

Pulsarトピックの作成

EMQXでデータ統合を作成する前に、対象となるPulsarトピックを作成しておく必要があります。以下のコマンドで、publicテナントのdefaultネームスペースに、パーティション数1のmy-topicトピックを作成します。

bash
docker exec -it pulsar bin/pulsar-admin topics create-partitioned-topic persistent://public/default/my-topic -p 1

コネクターの作成

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

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

  1. EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。
  2. ページ右上のCreateをクリックします。
  3. Create Connectorページで、Pulsarを選択し、Nextをクリックします。
  4. Configurationステップで以下の情報を設定します:
    • コネクター名を入力します。英数字の大文字・小文字の組み合わせで、例:my_pulsar
    • Bridge RoleはデフォルトでProducerが選択されています。
    • Pulsarサーバーへの接続およびメッセージ書き込み情報を設定します:
      • Serversにはpulsar://localhost:6650を入力します。リモート環境の場合は適宜変更してください。
      • Authenticationは認証方式を選択します:noneBasic auth、またはtokenBasic authの場合、EMQXはUsernamePassword:で連結して認証文字列を作成します。
      • Enable TLS:暗号化接続を確立したい場合はトグルをオンにします。TLS接続の詳細は外部リソースアクセスのTLSを参照してください。
  5. 詳細設定(任意):高度な設定を参照してください。
  6. Createをクリックする前に、Test ConnectivityでコネクターがPulsarサーバーに接続できるかテストできます。
  7. ページ下部のCreateボタンをクリックしてコネクターを作成します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてルールとSinkの作成を続行できます。詳細はCreate a Rule with Pulsar Sinkを参照してください。

Pulsar Sinkを使ったルールの作成

このセクションでは、ソースMQTTトピックt/#からのメッセージを処理し、処理済みデータを設定済みのSinkを介してPulsarトピックmy-topicに保存するルールの作成方法をダッシュボードで説明します。

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

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

  3. ルールIDを入力します。例:my_rule

  4. SQL Editorに以下のステートメントを入力します。これはトピックt/#のMQTTメッセージをPulsarに保存する例です。

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

    sql
    SELECT
      *
    FROM
      "t/#"

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

  5. + Add Actionボタンをクリックし、ルールでトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理したデータをPulsarに送信します。

  6. Action TypeのドロップダウンリストからPulsarを選択します。

  7. ActionドロップダウンはデフォルトのCreate Actionのままにします。既に作成済みのSinkを選択することも可能です。この例では新しいSinkを作成します。

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

  9. Connectorドロップダウンから先ほど作成したmy_pulsarを選択します。新しいコネクターを作成する場合は、ドロップダウン横のボタンをクリックしてください。設定パラメータの詳細はCreate a Connectorを参照してください。

  10. Sinkの以下のオプションを設定します:

    • Pulsar Topic Name:先に作成したpersistent://public/default/my-topicを入力します。注意:ここでは変数はサポートされていません。
    • Partition Strategy:プロデューサーがPulsarのパーティションにメッセージを振り分ける方法を選択します:randomroundrobin、またはkey_dispatch
    • Compression:圧縮アルゴリズムの使用有無と、Pulsarメッセージ内のレコード圧縮・解凍に使うアルゴリズムを指定します。選択肢はno_compressionsnappyzlibです。
    • Retention Period:Pulsarトピックにパブリッシュされたメッセージの保持期間を設定します。この設定により、メッセージがサブスクライバーに利用可能な期間を制御できます。デフォルトはinfinityで、メッセージの自動期限切れはありません。秒数で数値を指定すると、その時間を超えたメッセージは自動的に期限切れとなりトピックから削除されます。
    • Message Key:Pulsarメッセージのキー。プレーン文字列またはプレースホルダー(${var})を含む文字列を入力します。
    • Message Value:Pulsarメッセージの値。プレーン文字列またはプレースホルダー(${var})を含む文字列を入力します。
  11. フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した際にトリガーされます。詳細はフォールバックアクションを参照してください。

  12. 詳細設定(任意)高度な設定を参照してください。

  13. Createをクリックする前に、Test ConnectivityでコネクターがPulsarサーバーに接続できるかテストできます。

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

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

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

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

ルールのテスト

MQTTXを使ってトピックt/1にメッセージを送信します:

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

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

以下のPulsarコマンドで、メッセージがトピックpersistent://public/default/my-topicに書き込まれているか確認します:

bash
docker exec -it pulsar bin/pulsar-client consume -n 0 -s mysubscriptionid -p Earliest persistent://public/default/my-topic

高度な設定

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

項目説明推奨値
Max Inflightプロデューサーが各パーティションに送信できるメッセージバッチの最大数。
この数を増やすとスループットが向上します。
10
Sync Publish Timeout同期パブリッシュ操作で、メッセージが正常に配信されたことを確認するまでの最大待機時間(秒)。
配信問題やネットワーク障害時に無限待機を防ぎ、データ信頼性を確保します。
3
Socket Send Buffer Sizeネットワーク送信性能を最適化するためのソケットバッファサイズ。1 MB
Batch SizePulsarメッセージ内にバッチングされる個別リクエストの最大数。100
Max Batch BytesPulsarバッチ内で収集可能なメッセージの最大サイズ(バイト)。
通常、Pulsarブローカーのデフォルトは1MBですが、EMQXはPulsarのメッセージエンコードオーバーヘッドを考慮し、1MBよりやや小さい値を設定しています。個別メッセージがこの制限を超える場合は単独バッチとして送信されます。
900 KB
Query Modeメッセージ送信を最適化するため、asynchronousまたはsynchronousのクエリモードを選択可能。非同期モードではPulsarへの書き込みがMQTTメッセージのパブリッシュ処理をブロックしませんが、クライアントがPulsar到着前にメッセージを受信する可能性があります。Async
Buffer Modeメッセージ送信前のバッファリング方法を定義。メモリバッファリングで送信速度が向上します。
memory:メッセージをメモリにバッファリング。EMQXノード再起動時にメッセージは失われます。
disk:メッセージをディスクにバッファリング。EMQXノード再起動後もメッセージは保持されます。
hybrid:初めはメモリにバッファリングし、一定容量(segment_bytes設定参照)に達すると徐々にディスクにオフロードします。メモリモード同様、EMQXノード再起動時にメッセージは失われます。
memory
Pulsar Per-partition Buffer Limit各Pulsarパーティションの最大バッファサイズ(バイト)。この制限に達すると、古いメッセージを破棄してバッファスペースを確保します。
メモリ使用量とパフォーマンスのバランス調整に役立ちます。
2 GB
Segment File Bytesバッファモードがdiskまたはhybridの場合に適用。メッセージ保存用の分割ファイルサイズを制御し、ディスクストレージの最適化に影響します。100 MB
Memory Overload Protectionバッファモードがmemoryの場合に適用。EMQXはメモリ圧迫時に古いバッファメッセージを自動破棄し、システムの安定性を確保します。
注意:Linuxシステムのみ有効です。
disabled
Start Timeoutコネクターが自動起動したリソースの正常状態を待機する最大時間(秒)。この設定により、Polarなどの接続リソースが完全に稼働し、データ処理準備が整うまで操作を進めません。5
Health Check IntervalSinkの稼働状態をチェックする間隔。1