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はルールエンジンと設定されたSinkを通じてMQTTデータをApache Pulsarに転送し、その全体の流れは以下の通りです:
- メッセージのパブリッシュと受信:IoTデバイスはMQTTプロトコルを介して正常に接続を確立し、その後特定のトピックにテレメトリおよびステータスデータをパブリッシュします。EMQXはこれらのメッセージを受信すると、ルールエンジン内でマッチング処理を開始します。
- ルールエンジンによるメッセージ処理:組み込みのルールエンジンを使用して、特定のソースからのMQTTメッセージをトピックマッチングに基づき処理できます。ルールエンジンは対応するルールをマッチングし、データフォーマットの変換、特定情報のフィルタリング、コンテキスト情報の付加などのメッセージ処理を行います。
- Apache Pulsarへのデータストリーミング:ルールがトリガーされると、メッセージをPulsarに転送するアクションが実行されます。データはPulsarのメッセージキーおよび値に簡単に設定可能です。MQTTトピックはPulsarトピックにマッピングすることもでき、データの整理や識別が容易になり、後続のデータ処理や分析が促進されます。
MQTTメッセージデータがApache Pulsarに書き込まれた後は、以下のような柔軟なアプリケーション開発が可能です:
- Pulsarのコンシューマーアプリケーションを作成し、これらのメッセージをサブスクライブして処理します。ビジネスニーズに応じて、MQTTデータを他のデータソースと関連付けたり集約したり変換したりして、リアルタイムのデータ同期と統合を実現できます。
- 特定のMQTTメッセージを受信した際に、Pulsarのルールエンジンコンポーネントを使って対応するアクションやイベントをトリガーし、システム間やアプリケーション間のイベント駆動型機能を実装できます。
- Pulsar内でMQTTデータストリームをリアルタイムに分析し、異常や特定のイベントパターンを検出して、アラート通知や対応アクションを実行できます。
- 複数のMQTTトピックからのデータを統合し、Pulsarの計算能力を活用してリアルタイムの集約、計算、分析を行い、より包括的なデータインサイトを得ることができます。
特長とメリット
Pulsarとのデータ統合により、以下の特長と利点がビジネスにもたらされます:
- 信頼性の高いIoTデータメッセージ配信:EMQXはMQTTメッセージをバッチ処理で確実にPulsarに送信でき、IoTデバイスとPulsarおよびアプリケーションシステムの統合を実現します。
- MQTTメッセージ変換:ルールエンジンを使用して、EMQXはMQTTメッセージのフィルタリングや変換が可能です。メッセージはPulsarに送信される前にデータ抽出、フィルタリング、付加、変換を受けられます。
- 柔軟なトピックマッピング:Pulsar SinkはMQTTトピックをPulsarトピックに柔軟にマッピングでき、Pulsarメッセージのキー(Key)および値(Value)を簡単に設定可能です。
- 柔軟なパーティション選択:Pulsar SinkはMQTTトピックやクライアントに基づいて異なる戦略でPulsarのパーティションを選択でき、データの整理や識別に柔軟性を提供します。
- 高スループットシナリオでの処理能力:Pulsar Sinkは同期および非同期の書き込みモードをサポートし、シナリオに応じてレイテンシとスループットのバランスを柔軟に調整できます。
はじめる前に
このセクションでは、EMQXダッシュボードでPulsarデータ統合を作成する前に完了すべき準備について説明します。
前提条件
Pulsarのインストール
DockerでPulsarを起動します。
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というトピックを作成します。
docker exec -it pulsar bin/pulsar-admin topics create-partitioned-topic persistent://public/default/my-topic -p 1コネクターの作成
このセクションでは、SinkをPulsarサーバーに接続するためのコネクターの作成方法を説明します。
以下の手順は、EMQXとPulsarの両方をローカルマシンで実行していることを前提としています。リモートで実行している場合は設定を適宜調整してください。
- EMQXダッシュボードに入り、Integration -> Connectorsをクリックします。
- ページ右上のCreateをクリックします。
- Create ConnectorページでPulsarを選択し、Nextをクリックします。
- Configurationステップで以下の情報を設定します:
- コネクター名を入力します。大文字・小文字の英数字の組み合わせとしてください。例:
my_pulsar - Bridge Roleはデフォルトで
Producerが選択されています。 - Pulsarサーバーへの接続およびメッセージ書き込み情報を設定します:
- Serversに
pulsar://localhost:6650を入力します。リモート環境の場合は適宜調整してください。 - Authenticationで認証方式を選択します:
none、Basic auth、またはtoken。Basic authの場合、EMQXはUsernameとPasswordを:で連結して認証文字列を作成します。 - Enable TLS:暗号化接続を確立したい場合はトグルスイッチをオンにします。TLS接続の詳細は外部リソースアクセスのTLSを参照してください。
- Serversに
- コネクター名を入力します。大文字・小文字の英数字の組み合わせとしてください。例:
- 高度な設定(任意):高度な設定を参照してください。
- Createをクリックする前に、Test ConnectivityをクリックしてコネクターがPulsarサーバーに接続できるかテストできます。
- ページ下部のCreateボタンをクリックしてコネクターの作成を完了します。ポップアップダイアログでBack to Connector Listをクリックするか、Create RuleをクリックしてルールとSinkの作成を続行できます。詳細はCreate a Rule with Pulsar Sinkを参照してください。
Pulsar Sinkを用いたルールの作成
このセクションでは、DashboardでソースMQTTトピックt/#からのメッセージを処理し、処理済みデータを設定済みのSinkを介してPulsarトピックmy-topicに保存するルールの作成方法を示します。
EMQXダッシュボードで、Integration -> Rulesをクリックします。
ページ右上のCreateをクリックします。
ルールIDを入力します。例:
my_ruleSQL Editorに以下のステートメントを入力します。これはトピック
t/#のMQTTメッセージをPulsarに保存するためのものです。注意:独自のSQL構文を指定する場合は、Sinkで必要なすべてのフィールドが
SELECT部分に含まれていることを確認してください。sqlSELECT * FROM "t/#"注意:初心者の方はSQL Examplesをクリックし、Enable TestでSQLルールを学習・テストできます。
+ Add Actionボタンをクリックして、ルールによってトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理されたデータをPulsarに送信します。
Action Typeドロップダウンリストから
Pulsarを選択します。Actionドロップダウンはデフォルトの
Create Actionのままにします。既に作成済みのSinkを選択することも可能ですが、この例では新しいSinkを作成します。Sinkの名前を入力します。名前は大文字・小文字の英数字の組み合わせとしてください。
Connectorドロップダウンから先ほど作成した
my_pulsarを選択します。ドロップダウン横のボタンから新しいコネクターを作成することも可能です。設定パラメーターはCreate a Connectorを参照してください。Sinkの以下のオプションを設定します:
- Pulsar Topic Name:先に作成した
persistent://public/default/my-topicを入力します。注意:ここで変数はサポートされていません。 - Partition Strategy:プロデューサーがメッセージをPulsarのパーティションに振り分ける方法を選択します:
random、roundrobin、またはkey_dispatch。 - Compression:圧縮アルゴリズムの使用有無と、Pulsarメッセージ内のレコードの圧縮・解凍に使用するアルゴリズムを指定します。選択肢は
no_compression、snappy、zlibです。 - Retention Period:Pulsarトピックにパブリッシュされたメッセージが保持される期間を定義します。この設定により、サブスクライバーがメッセージを消費可能な期間を制御できます。デフォルトは
infinityで、メッセージの自動期限切れはありません。秒数で数値を指定すると、その時間を超えたメッセージは自動的に期限切れとなりトピックから削除されます。 - Message Key:Pulsarメッセージのキーを入力します。プレーンな文字列またはプレースホルダー(${var})を含む文字列が使用可能です。
- Message Value:Pulsarメッセージの値を入力します。プレーンな文字列またはプレースホルダー(${var})を含む文字列が使用可能です。
- Pulsar Topic Name:先に作成した
フォールバックアクション(任意):メッセージ配信失敗時の信頼性向上のために、1つ以上のフォールバックアクションを定義できます。これらはプライマリSinkがメッセージ処理に失敗した場合にトリガーされます。詳細はフォールバックアクションを参照してください。
高度な設定(任意):高度な設定を参照してください。
Createをクリックする前に、Test ConnectivityをクリックしてコネクターがPulsarサーバーに接続できるかテストできます。
CreateボタンをクリックしてSinkの設定を完了します。新しいSinkがAction Outputsに追加されます。
Create Ruleページに戻り、設定内容を確認後、Createボタンをクリックしてルールを生成します。
これでルールの作成が完了しました。Integration -> Rulesページで新規作成したルールを確認できます。**Actions(Sink)**タブをクリックすると、新しいPulsar Sinkが表示されます。
また、Integration -> Flow Designerをクリックするとトポロジーが表示され、トピックt/#のメッセージがPulsarに送信・保存されていることが確認できます。
ルールのテスト
MQTTXを使ってトピックt/1にメッセージを送信します:
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "Hello Pulsar" }'Sinkの稼働状況を確認すると、新規の受信メッセージと送信メッセージがそれぞれ1件あるはずです。
以下のPulsarコマンドで、メッセージがトピックpersistent://public/default/my-topicに書き込まれているか確認します:
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 Size | Pulsarメッセージ内にバッチングされる個別リクエストの最大数。 | 100 |
| Max Batch Bytes | Pulsarバッチ内で収集可能なメッセージの最大サイズ(バイト)。通常、Pulsarブローカーのデフォルトは1MBですが、EMQXはメッセージエンコードのオーバーヘッドを考慮し、デフォルト値を1MB未満に設定しています。単一メッセージがこの制限を超える場合は別バッチで送信されます。 | 900 KB |
| Query Mode | メッセージ送信を最適化するために、asynchronousまたはsynchronousのクエリモードを選択可能。非同期モードではPulsarへの書き込みがMQTTメッセージパブリッシュ処理をブロックしませんが、クライアントがPulsar到着前にメッセージを受信する可能性があります。 | Async |
| Buffer Mode | メッセージ送信前にバッファリングするかどうかを定義。メモリバッファリングは送信速度を向上させます。memory: メッセージはメモリにバッファされ、EMQXノード再起動時に失われます。disk: メッセージはディスクにバッファされ、EMQXノード再起動後も保持されます。hybrid: 初期はメモリにバッファし、一定量(segment_bytes設定参照)を超えると徐々にディスクにオフロードされます。メモリモード同様、ノード再起動時にメッセージは失われます。 | 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 Interval | Sinkの稼働状況をチェックする間隔。 | 1 秒 |