メッセージ変換
メッセージ変換は、メッセージがさらに処理される前やサブスクライバーに配信される前に、ユーザー定義のルールに基づいてメッセージを修正およびフォーマットする機能です。この機能は高度にカスタマイズ可能で、複数のエンコーディングや高度な変換をサポートしています。
ワークフロー
メッセージがパブリッシュされると、以下のワークフローを経ます。
スキーマ検証:メッセージがパブリッシュされ認可を通過すると、まずスキーマ検証が行われます。メッセージが検証に合格すると、次のステップに進みます。
メッセージ変換パイプライン:
- 変換マッチング:メッセージは、そのトピックに基づいてユーザー定義の変換リストと照合されます。異なるトピックやトピックフィルターに対して複数の変換を設定できます。
- 変換実行:マッチした変換は設定された順序で実行されます。パイプラインはJSON、Protobuf、Avroなどの各種エンコーダー・デコーダーをサポートし、Variform式を用いてメッセージの拡張や修正が可能です。
- 変換後処理:メッセージが変換パイプラインを正常に通過すると、ルールエンジンのトリガーやサブスクライバーへのメッセージ配信など、次の処理に進みます。
失敗時の処理:変換が失敗した場合、ユーザー設定のアクションが実行されます。
- メッセージ破棄:パブリッシュを終了しメッセージを破棄します。QoS 1およびQoS 2のメッセージにはPUBACKで特定の理由コード(131 - 実装固有のエラー)が返されます。
- 切断してメッセージ破棄:メッセージを破棄し、パブリッシュしたクライアントを切断します。
- 無視:追加のアクションは行いません。
変換失敗時には、設定されたアクションに関わらずログエントリが生成される場合があります。ログの出力レベルはユーザーが設定可能で、デフォルトは
warningです。さらに、変換失敗はルールエンジンのイベント($events/message_transformation/failed)をトリガーでき、これによりユーザーは誤ったメッセージを別トピックに再パブリッシュしたり、Kafkaへ送信して詳細分析を行うなどのカスタム処理を実装できます。
ユーザーガイド
このセクションでは、メッセージ変換機能の設定方法とテスト方法を説明します。
ダッシュボードでのメッセージ変換設定
ダッシュボードでメッセージ変換を作成および設定する手順は以下の通りです。
ダッシュボードにアクセスし、左のナビゲーションメニューから Smart Data Hub -> Message Transform をクリックします。
Message Transform ページ右上の Create をクリックします。
「Create Message Transform」ページで以下の情報を設定します:
Name:変換の名前を入力します。
Message Source Topic:変換対象のメッセージが属するトピックを設定します。複数のトピックやトピックフィルターを設定可能です。
Note(任意):メモを入力します。
Message Format Transformation:
Source Format:変換パイプラインに入るメッセージに適用するペイロードデコーダーを指定します。選択肢は以下の通りです:
None(デコードなし)JSONAvroProtobufCustom (External HTTP)
これらのデコーダーはバイナリの入力ペイロードを構造化マップに変換します。
Avro、Protobuf、Custom (External HTTP)を選択する場合は、あらかじめスキーマレジストリで作成されている必要があります。パイプラインに複数の変換がある場合、各ステップでデコードを行う必要はありません。例えば、変換
T1で既にペイロードがデコードされていれば、後続の変換T2はデコードをスキップし、正しい形式のペイロードを利用できます。Target Format:変換パイプラインの最後にメッセージペイロードをバイナリ値としてエンコードするためのペイロードエンコーダーを指定します。エンコーダーの選択肢は Source Format と同じです。
パイプラインの最後の変換のみがペイロードをバイナリ値にエンコードする必要があり、中間の変換はバイナリエンコードを処理する必要はありません。
Message Properties Transformation:
- Properties:式の結果として得られた変換値を書き込む先を指定します。有効な宛先は
payload、topic、qos、retain(対応するフラグを設定)、およびuser_property(MQTTのUser-Property)です。user_propertyを使用する場合は、このフィールドの下に正確に1つのキーを指定する必要があります(例:user_property.my_custom_prop)。payloadはそのまま使用してメッセージペイロード全体を上書きするか、ネストされたJSONオブジェクトとして特定のキーのパスを指定できます(例:payload.x.y)。 - Target Value:設定したプロパティに書き込む値を定義します。この値は
qos、retain、topic、payload、payload.x.yなどの他のフィールドからコピーするか、Variform式を使って生成できます。
- Properties:式の結果として得られた変換値を書き込む先を指定します。有効な宛先は
Transformation Failure Operation:
- Action After Failure:変換が失敗した場合に実行するアクションを選択します:
- Drop Message:パブリッシュ処理を終了しメッセージを破棄します。QoS 1およびQoS 2のメッセージにはPUBACKで特定の理由コードが返されます。
- Disconnect and Drop Message:メッセージを破棄し、パブリッシュしたクライアントを切断します。
- Ignore:追加のアクションは行いません。
- Action After Failure:変換が失敗した場合に実行するアクションを選択します:
Output Logs:変換失敗時にログを生成するか選択します。デフォルトでログは有効です。
Logs Level:ログの出力レベルを設定します。デフォルトは
warningです。
Create をクリックして設定を完了します。
作成前に Preview をクリックして変換をテストできます。新しいペインが開き、QoS、ペイロード、リテインフラグの有無、パブリッシャーのユーザー名やクライアントIDなど、受信メッセージのコンテキストを入力できます。必要な情報を入力後、Execute Transformation をクリックすると指定したコンテキストで変換を実行し、結果を確認できます。
変換が作成されると、メッセージ変換ページのリストにデフォルトで有効な状態で表示されます。必要に応じて無効化したり、Actions 列の Settings をクリックして設定を更新できます。削除や順序変更は More をクリックして行います。
設定ファイルでのメッセージ変換設定
Avro形式でエンコードされたメッセージを受信し、JSONにデコードしたいとします。デコード後、パブリッシュしたクライアントのクライアント属性から取得した tenant 属性をトピックの先頭に付加し、その後ルールエンジンで処理したい場合、以下の設定で実現できます。
message_transformation {
transformations = [
{
name = mytransformation
topics = ["t"]
failure_action = drop
payload_decoder = {type = avro, schema = myschema}
payload_encoder = {type = json}
operations = [
{key = "topic", value = "concat([client_attrs.tenant, '/', topic])"}
]
}
]
}この設定は、mytransformation という名前の変換を指定し、
- 指定したスキーマでAvro形式のメッセージペイロードをデコードし、
- ペイロードをJSON形式にエンコードし、
- クライアント属性の
tenantと元のトピックを連結してトピックを変更します。
詳細な設定方法については、設定マニュアルをご覧ください。
REST API
REST APIを通じたメッセージ変換の詳細な使用方法は、EMQX Enterprise APIをご参照ください。
デコード/エンコード用スキーマの作成
デコーダーおよびエンコーダースキーマの作成方法については、スキーマレジストリのセクションをご覧ください。
統計と指標
有効化すると、メッセージ変換はダッシュボード上で統計と指標を公開します。メッセージ変換ページで変換名をクリックすると、以下の情報が表示されます。
統計情報:
- 合計:システム起動以降のトリガー総数
- 成功:成功したデータ変換の数
- 失敗:失敗したデータ変換の数
レート指標:
- 現在の変換速度
- 過去5分間の速度
- 過去の最大速度
統計はリセット可能で、Prometheusの /prometheus/message_transformation からも取得可能です。