定義
$meta 式は、ドキュメントのすべてのストリーミング メタデータを含むオブジェクトを返します。このデータは、ストリーム全体、または以下の Atlas Stream Processing の集計ステージのいずれかに対して公開できます。
$meta 式には次のプロトタイプ形式があります。
{ "$meta": <string> }
"source": { "type": "<source-type>", "ts": { "$date": "<datetime>" }, "topic": "<string>", "partition": <int>, "offset": <int>, "key": "<kafka-key>", "headers": [ { "k": "<header-key>", "v": "<header-value>" } ], "operationType": "<db-operation>", "ns": { "db": "<namespace-db>", "coll": "<namespace-coll>" }, "documentKey": { "_id": { "$oid": "<object-id>" } }, "initialSync": { "phase": "<sync-state>" } "kinesisStream": "<kinesis-name>", "shardId": "<kinesis-shard-id>", "sequenceNumber": "<doc-uuid>", "partitionKey": "<partition-id>", } "window": { "start": <ISODate>, "end": <ISODate>, "partition": "<session-partition>" }, "https": { "url": "<target-url>", "method": "<request-method>", "httpStatusCode": <http-code>, "responseTimeMs": <response-time-ms> }
構文
$meta 式は、メタデータのソースの完全修飾ドット構文パスに対応する単一の文字列入力を受け取ります。このパスのルートは "stream" である必要があります。次のパスをクエリできます。
パス | タイプ | 条件付き | 説明 |
|---|---|---|---|
| オブジェクト | 常に | |
| ドキュメント | 常に |
|
| string | 常に | ソースとして使用される接続のタイプ。 |
| ISODate | 常に | 取り込み時点でのレコードの日時。 |
| string | 条件付き | ストリームがレコードを取り込む Kafka トピック。Kafka ソースにのみ適用されます。 |
| integer | 条件付き | ストリームがレコードを取り込む Kafka トピックのパーティション。Kafka ソースにのみ適用されます。 |
| integer | 条件付き | Kafka ソース パーティション内のメッセージ順序とキューの位置をオフセット追跡します。Kafka ソースにのみ適用されます。 |
| string|int|long|double|object|binData | 条件付き | パーティショニングと負荷分散のために Kafka メッセージに割り当てられたキー。Kafka ソースにのみ適用されます。 |
| 配列 | 条件付き | Kafka メッセージのメタデータを記述するキーと値のペアのセット。Kafka ソースにのみ適用されます。 |
| string | 条件付き | Atlas Stream Processing が指定されたドキュメントに対して実行しようとしたデータベース操作の種類。Atlas 変更ストリーム ソースにのみ適用されます。 |
| ドキュメント | 条件付き | Atlas Stream Processing がドキュメントを取得する名前空間を含むドキュメント。Atlas 変更ストリーム ソースにのみ適用されます。 |
| string | 条件付き | Atlas Stream Processing が操作を試行するデータベース名。Atlas 変更ストリーム ソースにのみ適用されます。 この値は、コレクション変更ストリームまたはデータベース変更ストリームソースのすべてのドキュメントで同じです。クラスター変更ストリームソースの場合は異なります。 |
| string | 条件付き | Atlas Stream Processing が操作を試行するコレクションの名前。Atlas 変更ストリーム ソースにのみ適用されます。 この値は、コレクション変更ストリームソースのすべてのドキュメントで同じです。データベース変更ストリームまたはクラスター変更ストリームソースの場合は異なります。 |
| ドキュメント | 条件付き | ソースドキュメントのオブジェクト ID を含むドキュメント。Atlas 変更ストリーム ソースにのみ適用されます。 |
| string | 条件付き | 最初の同期作業の現在の状態。最初の同期中の Atlas 変更ストリームソースにのみ適用されます。 |
| string | 条件付き | Atlas Stream Processing がドキュメントを取得する Kinesis Data Stream の名前。AWS Kinesis ソースにのみ適用されます。 |
| string | 条件付き | Atlas Stream Processing がドキュメントを取得する Kinesis Data Stream 内のシャードの ID。AWS Kinesis ソースにのみ適用されます。 |
| string | 条件付き | Kinesis Data Stream から取得したドキュメントの一意の識別子。Amazon Web Services Kinesis ソースにのみ適用されます。 |
| string | 条件付き | ソース ドキュメントが属するパーティションの一意の識別子。AWS Kinesis ソースにのみ適用されます。 |
| ドキュメント | 条件付き | ウィンドウのメタデータを含むドキュメント。ドキュメントがウィンドウで処理された場合にのみ適用されます。 |
| ISODate | 条件付き | ウィンドウ開始時間。ドキュメントがウィンドウで処理された場合にのみ適用されます。 |
| ISODate | 条件付き | ウィンドウの閉じる時間。ドキュメントがウィンドウで処理された場合にのみ適用されます。 |
| string | 条件付き | ドキュメントが属すするセッション ウィンドウ パーティション。ドキュメントがセッション ウィンドウで処理された場合にのみ適用されます。 |
| ドキュメント | 条件付き | $https ステージのメタデータを含むドキュメント。プロセシングの失敗が |
| string | 条件付き |
|
| string | 条件付き |
|
| 整数 | 条件付き | リクエストの HTTP 応答ステータス コード。プロセシングの失敗が |
| 整数 | 条件付き | リクエストの応答時間(ミリ秒)。プロセシングの失敗が |
動作
Atlas Stream Processing$meta 式は、既存のMongoDB$meta 集計式のすべての機能を提供します。ただし、標準のMongoDB集計クエリでは、$meta の Atlas Stream Processing バージョンに固有の機能を使用することはできません。
例
次の例では、データが取り込まれた Kafka ソース トピックの配列を使用してストリームの出力を強化します。
{ $source: { connectionName: "kafka", topic: ["t1", "t2", "t3"] } }, { $emit: { connectionName: "kafka", topic: { $concat: [ { $meta: "stream.source.topic" }, "out" ] } } }
次の例では、各ウィンドウの開始時刻を報告するストリームにフィールドを追加します。
{ $source: { connectionName: "kafka", topic: "t1" } }, { $hoppingWindow: . . . }, { $addFields: { start: { $meta: "stream.window.start" } } }