Skip to content

MQTT Streams ​

EMQX 6.1では、MQTTのリアルタイムなパブリッシュ/サブスクライブモデルを拡張し、永続的でリプレイ可能なメッセージストリームを実現するストリーミングおよびリプレイ機能であるMQTT Streamsを導入しました。これにより、MQTTのセマンティクスを維持しつつ、Kafkaのようなストリーミング機能を提供します。

本ページでは、EMQXにおけるMQTT Streams機能の設計動機、主要概念、内部アーキテクチャ、メッセージフロー、実際の適用シナリオまでを包括的に解説します。

Listener Mountpointとの非互換性について

mountpointが設定されたリスナー経由で接続されたクライアントに対しては、ストリームのリプレイは機能しません。EMQXは$stream/プレフィックスのマッチング前にmountpointを適用するため、$stream/<name>へのサブスクライブはマウントされたリテラルトピックへの通常のサブスクライブとして扱われ、クライアントにはエラーが通知されません。メッセージの取り込みも影響を受け、パブリッシュされたトピックにはmountpointプレフィックスが付与されるため、そのプレフィックスを含むトピックフィルターのストリームにのみマッチします。

MQTT Streamとは何か? ​

MQTT Streamは、名前付きの論理リソースであり、設定されたトピックフィルターにマッチするMQTTメッセージを継続的に収集します。メッセージはストリームの保持ポリシーに従って永続的に保存され、後からサブスクライブするクライアントによってリプレイ可能です。

ストリームは一意の名前で識別されます。トピックフィルターはストリームの設定の一部ですが、識別子ではありません。

各ストリームは以下を持ちます:

  • 一意の名前
  • 設定されたトピックフィルター
  • 保持ポリシー(時間ベースおよび/またはサイズベース)
  • 任意のLast-Valueセマンティクス
  • 明示的なライフサイクル(作成、更新、削除)

なぜMQTT Streamsを使うのか? ​

MQTTはリアルタイムメッセージングに最適化されていますが、以下のような制約があります:

  • メッセージは通常オンラインのサブスクライバーにのみ配信される。
  • 過去データのリプレイはネイティブにはサポートされていない。
  • 過去データの再処理には外部システムが必要。
  • 順序付けられたリプレイ可能なメッセージ履歴の維持が困難。

MQTT Streamsは、永続的なメッセージ保存とリプレイ機能をMQTTに拡張します。これにより、MQTTクライアントのパブリッシュやサブスクライブ方法を変更せずに、過去のメッセージやデバイスの最新状態を取得可能にします。

MQTT Streamsの主要概念 ​

  • MQTT Stream

    名前で識別・アドレス指定され、明示的なライフサイクルで管理される論理リソースです。アクティブな間は、設定された時間またはサイズの制限内でマッチするメッセージを継続的に保存します。保存されたメッセージは、パブリッシャー側の変更なしにサブスクライバーによってリプレイ可能です。

    ストリーム名に使用できる文字は以下のみです:

    • 英数字(A–Z, a–z, 0–9)
    • アンダースコア(_)
    • ハイフン(-)
    • ドット(.)

    サポートされるストリームタイプは2種類です:

    • 通常ストリーム:過去のデータを上書きせずにすべてのマッチするメッセージを保存します。stream-offsetサブスクリプションプロパティを使い、指定したタイムスタンプやオフセットからリプレイ可能です。
    • Last-Valueストリーム:Last-Valueセマンティクスを有効にします。同じストリームキーを持つメッセージは新しいものが古いものを上書きし、各キーに関連付けられた最新のメッセージのみを保持します。詳細はStream Key Expressionをご覧ください。
  • トピックフィルター

    sensors/+/dataのようなMQTTトピックフィルターで、どのパブリッシュメッセージをストリームに取り込むかを決定します。マッチするメッセージのみが取り込まれ、1つのメッセージが複数のストリームに属することもあります。

    TIP

    トピックフィルターはストリームの識別子ではなく、名前付きストリームの設定メタデータです。

  • ストリームサブスクリプション

    ストリームからメッセージを消費するための特別なMQTTサブスクリプションです。クライアントは以下のいずれかの形式でサブスクライブします:

    SUBSCRIBE $stream/<name>
    SUBSCRIBE $stream/<name>/<topic_filter>

    ここで:

    • <name> はストリーム名(必須)
    • <topic_filter> は既存ストリームにサブスクライブする際は省略可能
    • オートクリエーションが有効な場合、$stream/<name>/<topic_filter>で指定したトピックフィルターを使い、存在しない場合はEMQXがストリームを作成します。

    ストリームサブスクリプションは通常のMQTTサブスクリプションとは独立して動作し、External Subscription機構を通じて配信されます。

  • ストリームオフセット(リプレイ開始点)

    リプレイの開始点はトピックパスではなく、MQTT 5のユーザーサブスクリプションプロパティstream-offsetで指定します。

    stream-offsetプロパティはリプレイ開始位置を決定します。例:

    • タイムスタンプ
    • 論理オフセット
    • 最初や最後などの特別な位置(対応している場合)

    この設計により、トピック文字列からのオフセット解析が不要になり、MQTT 5プロパティによるリプレイ制御と整合します。

  • キー式

    各受信メッセージに対して評価されるユーザー定義式で、キーを抽出します。式はメッセージの内容やメタデータを参照可能です。抽出されたキーはストレージパーティション内のメッセージ順序保証に使われます。Last-Valueセマンティクスが有効な場合は、キーが上書き範囲を定義し、同じキーの新しいメッセージが古いものを置き換えます。

MQTT Streamsのアーキテクチャ ​

MQTT Streamsは、ブローカーコアと疎結合の独立したEMQXアプリケーションとして実装され、既存のインフラを再利用しています。EMQXとの統合は内部フックとExternal Subscriptionフレームワークを通じて行われます。

TIP

External Subscriptionは、ライブMQTTパブリッシュ経路外から発生したメッセージをMQTTクライアントセッションに接続し、クライアントの動作を変えずに標準MQTTサブスクリプション経由でメッセージを配信するEMQXの仕組みです。

主なコンポーネント ​

  • Streams Registry:MQTTストリームのライフサイクルを管理し、ストリーム名、トピックフィルター、保持ポリシー、キー式などのメタデータを保持します。Mnesiaテーブルを使い効率的にストリームを検索します。
  • Streams Message Database:ストリームメッセージの永続ストレージを提供し、EMQXのパーシステンス上に構築されています。メッセージを永続化し、保持制限を適用し、Last-Valueセマンティクスを有効にし、保持ポリシーに従って効率的にメッセージを取得可能です。
  • Streams ExtSub Handler:メッセージストリームをMQTTクライアントセッションと統合します。パーシステントストレージからメッセージを取得し、External Subscriptionフレームワークを通じてサブスクライバーに配信します。

MQTT Streamsのデータフローダイアグラム ​

以下の図はMQTT Streamsコンポーネント間のデータフローを示しています:

streams_data_flow

パブリッシュフロー ​

  1. クライアントがMQTTトピックにメッセージをパブリッシュします。
  2. MQTT Streamsのフックがトリガーされ、パブリッシュ処理を開始します。
  3. フックはStreams Registryに問い合わせ、メッセージトピックにマッチするストリームを特定します。
  4. マッチした各ストリームに対してメッセージを書き込み、パーシステントストレージに保存します。

サブスクライブおよび消費フロー ​

  1. クライアントがストリームトピック($stream/<name>または$stream/<name>/<topic_filter>)にサブスクライブし、任意でstream-offsetサブスクリプションプロパティを含めます。

    廃止予定

    旧形式の$s/<offset>/<topic_filter>は後方互換のためサポートされていますが廃止予定です。詳細は互換性の注意点をご覧ください。

  2. External Subscriptionフレームワークがサブスクリプションを処理し、ストリームトピック用のStreams ExtSubハンドラーを初期化します。

  3. ハンドラーは指定されたstream-offsetと保持ルールに従い、パーシステントストレージからメッセージを取得します。

  4. 取得したメッセージはExternal Subscriptionフレームワークに渡されます。

  5. ExtSubアプリケーションが標準MQTT配信を通じてクライアントにメッセージを届けます。

MQTT Streamsのコア機能 ​

MQTT Streamsは、メッセージの保存、順序付け、保持、リプレイ消費のための配信方法を定義する一連のコア機能を提供します。

  • オフセットベースのリプレイ

    リプレイ開始位置はトピックパスではなく、stream-offsetサブスクリプションプロパティで指定します。指定オフセットより前のメッセージはスキップされます。

  • 保持

    ストリームの保持ポリシーはメッセージリプレイの範囲を制限します。メッセージは時間またはサイズで制限された期間保持され、期限切れのメッセージは消費済みか否かに関わらず自動的に削除されます。

  • キーごとの順序保証

    MQTT Streamsは単一のグローバルな配信順序を保証しません。同じキーを持つメッセージは常にパブリッシュされた順序で配信されます。異なるキーのメッセージは任意の順序で配信される可能性があります。

  • Last-Valueセマンティクス

    ストリームはLast-Valueセマンティクスを有効にできます。同じキーを持つメッセージは新しいものが古いものを上書きし、最新の値のみが保持されます。キーが解決できないメッセージは通常通り保存されます。

  • MQTTネイティブ配信

    ストリームメッセージは標準のMQTTメカニズムを使って配信されます。パブリッシャーは動作を変更する必要がありません。サブスクライバーへのメッセージ配信はExternal Subscriptionを通じて統合されています。

セキュリティ上の考慮事項 ​

EMQXは、$stream/プレフィックスとストリーム名を含む完全なサブスクリプショントピックフィルターに対してストリームサブスクリプションの認可を行います。ストリームサブスクリプションの認可により、クライアントはサブスクリプション作成前にストリームが保存したメッセージをリプレイ可能になる場合があります。

ストリームサブスクリプションには専用のルールが必要 ​

通常のトピック空間用に書かれたルールは対応するストリームサブスクリプションには適用されません:

  • t/#用のルールは$stream/events/t/#には適用されません。EMQXはこれらを異なるトピックとして扱います。
  • #や+で始まる認可トピックフィルターは$で始まるサブスクリプショントピックフィルターにマッチしません。デフォルトのacl.confにある{eq, "#"}ルールを含む#拒否ルールは$stream/events/#を拒否しません。

そのため、#が拒否されたクライアントでも$stream/events/#にサブスクライブし、stream-offsetをearliestに設定してストリームにまだ保存されているすべてのメッセージをリプレイ可能です。$stream/用の明示的なルールを追加し、ストリームのトピックフィルターにマッチするトピックのルールと同等以上に厳しくしてください:

erlang
%% 自動作成を許可せず、事前作成済みの"events"ストリームのリプレイを1つのコンシューマに許可
{allow, {username, "event_reader"}, subscribe, ["$stream/events"]}.

%% 廃止予定のプレフィックスを含むその他すべてのストリームサブスクリプションを拒否
{deny, all, subscribe, ["$stream/#", "$s/#"]}.

クライアントがまだ使用している間は廃止予定の$s/プレフィックスもルールに残してください。詳細は廃止予定プレフィックスを参照してください。

この挙動は共有サブスクリプションとは異なります。$share/<group>/t/#の場合、EMQXはプレフィックスを除去し t/#を認可しますが、$stream/<name>/t/#の場合は完全なサブスクリプショントピックフィルターを認可します。

オートクリエーションはクライアント指定のトピックフィルターを許可 ​

ストリームのオートクリエーションが有効な場合、サブスクライブするクライアントが新規ストリームのトピックフィルターを決定します。EMQXは$stream/<name>/<topic_filter>の<topic_filter>を使ってストリームを作成し、クライアントが直接<topic_filter>のサブスクライブ権限を持つかは別途チェックしません。

例えば、クライアントが$stream/+/#へのサブスクライブを許可されている場合、$stream/events/#にサブスクライブ可能です。eventsストリームが存在しなければ、EMQXはトピックフィルター#で作成します。これにより、$で始まらないすべてのトピックのメッセージを保存します。クライアントは直接サブスクライブ権限のないトピックのメッセージをリプレイできる可能性があります。

オートクリエーションはデフォルトでLast-Valueストリームに対して有効で、通常ストリームでは無効ですが、MQTT Streamsが有効な場合のみ機能します。MQTT Streamsはデフォルトで無効(streams.enable = false)です。信頼できないクライアントを受け入れる環境では、$stream/eventsなどの特定の事前作成済みストリームのみアクセスを許可し、$stream/#や$s/#にマッチする他のサブスクリプションはすべて拒否してください。あるいは、オートクリエーションを無効にして、ダッシュボードやREST APIからストリームを作成してください。詳細はダッシュボード経由でのストリーム自動作成を参照してください。

互換性の注意点 ​

既存のデプロイメントに関する互換性の考慮点を説明します。

名前付きストリーム ​

  • すべてのストリームは明示的に名前付きリソースとなりました。
  • ストリーム名は許可された文字セットのルールに従う必要があります。

旧ストリーム ​

以前に作成された名前なしストリームには、トピックフィルターから派生した名前が自動的に割り当てられます。

派生名は/<topic_filter>となります。

廃止予定プレフィックス ​

ストリームサブスクリプションの$sプレフィックスは後方互換のため引き続きサポートされていますが、廃止予定です。

新規デプロイメントでは$stream/<name>を使用してください。

典型的なユースケース ​

  • 過去データのリプレイ:過去のMQTTイベントを再処理し、デバッグや新しいビジネスロジックに活用。
  • 時系列分析:センサーデータを保存・リプレイし、分析や予知保全に利用。
  • イベントソーシング:すべての状態変化を不変のイベントログとして永続化。
  • IoTデジタルツイン:物理デバイスの最新状態をデジタル形式で維持。
  • 設定同期:デバイスが常に最新の設定を受け取ることを保証。

次のステップ ​

MQTT Streamsの基本を理解したら、実際の利用方法を学びましょう:

  • ストリームの作成と設定:ダッシュボードやREST APIでストリームを宣言し、Last-Valueセマンティクスや保持ポリシーを設定する方法を解説。
  • クイックスタートチュートリアル:MQTTXを使って実際のパブリッシャーとサブスクライバーのシナリオをシミュレートするステップバイステップガイド。