定义
$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"。您可以查询以下路径:
路径 | 类型 | 条件性 | 说明 |
|---|---|---|---|
| 对象 | 始终 | |
| 文档 | 始终 | 包含 |
| 字符串 | 始终 | 用作数据源的连接类型。 |
| ISODate | 始终 | 记录在数据引入点的日期和时间。 |
| 字符串 | 可选的 | 流从中摄取记录的 Kafka 主题。仅适用于 Kafka 源。 |
| 整型 | 可选的 | 流从中提取记录的 Kafka 主题分区。仅适用于 Kafka 源。 |
| 整型 | 可选的 | 在 Kafka 源分区中追踪消息顺序和队列位置的偏移量。仅适用于 Kafka 源。 |
| string|int|long|double|object|binData | 可选的 | 分配给 Kafka 消息的密钥,用于分区和负载分配。仅适用于 Kafka 源。 |
| 阵列 | 可选的 | 描述 Kafka 消息元数据的键值对集。仅适用于 Kafka 源。 |
| 字符串 | 可选的 | Atlas Stream Processing 试图对给定文档执行的数据库操作类型。仅适用于 Atlas 变更流源。 |
| 文档 | 可选的 | 包含 Atlas Stream Processing 源文档所在命名空间的文档。仅适用于 Atlas 变更流源。 |
| 字符串 | 可选的 | |
| 字符串 | 可选的 | |
| 文档 | 可选的 | 包含源文档的对象 ID 的文档。仅适用于 Atlas 变更流源。 |
| 字符串 | 可选的 | 初始同步操作的当前状态。仅在初始同步期间适用于 Atlas 变更流源。 |
| 字符串 | 可选的 | Atlas Stream Processing 源文档的 Kinesis 数据流的名称。仅适用于 AWS Kinesis 源。 |
| 字符串 | 可选的 | Kinesis Data Stream 中的分片 ID,Atlas Stream Processing 从中获取文档。仅适用于 Amazon Web Services Kinesis 源。 |
| 字符串 | 可选的 | 来自 Kinesis Data Stream 的源文档的唯一标识符。仅适用于 Amazon Web Services Kinesis 源。 |
| 字符串 | 可选的 | 源文档所属分区的唯一标识符。仅适用于 Amazon Web Services Kinesis 源。 |
| 文档 | 可选的 | 包含窗口元数据的文档。仅当文档在窗口中处理时才应用。 |
| ISODate | 可选的 | 窗口打开时间。仅当文档在窗口中处理时应用。 |
| ISODate | 可选的 | 窗口关闭时间。仅当文档在窗口中处理时应用。 |
| 字符串 | 可选的 | 文档所属的会话窗口分区。仅当文档在会话窗口中处理时才应用。 |
| 文档 | 可选的 | 包含 $https 阶段元数据的文档。仅当处理失败发生在 |
| 字符串 | 可选的 |
|
| 字符串 | 可选的 |
|
| int | 可选的 | 请求的 HTTP 响应状态代码。仅当处理失败发生在 |
| int | 可选的 | 请求的响应时间(毫秒)。仅当处理失败发生在 |
行为
Atlas Stream Processing $meta表达式提供现有MongoDB $meta聚合表达式的所有功能。但是,您无法在标准MongoDB聚合查询中使用Atlas Stream Processing版本的 $meta 的特定功能。
示例
以下示例通过包含数据摄取来源的 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" } } }