对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
Docs 菜单

$meta 聚合阶段(流处理)

$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"。您可以查询以下路径:

路径
类型
条件性
说明

stream

对象

始终

所有 $source 阶段、任何窗口阶段或 $https 阶段在管道中配置的元数据。

stream.source

文档

始终

包含 $source 阶段元数据的文档。

stream.source.type

字符串

始终

用作数据源的连接类型。

stream.source.ts

ISODate

始终

记录在数据引入点的日期和时间。

stream.source.topic

字符串

可选的

流从中摄取记录的 Kafka 主题。仅适用于 Kafka 源。

stream.source.partition

整型

可选的

流从中提取记录的 Kafka 主题分区。仅适用于 Kafka 源。

stream.source.offset

整型

可选的

在 Kafka 源分区中追踪消息顺序和队列位置的偏移量。仅适用于 Kafka 源。

stream.source.key

string|int|long|double|object|binData

可选的

分配给 Kafka 消息的密钥,用于分区和负载分配。仅适用于 Kafka 源。

stream.source.headers

阵列

可选的

描述 Kafka 消息元数据的键值对集。仅适用于 Kafka 源。

stream.source.operationType

字符串

可选的

Atlas Stream Processing 试图对给定文档执行的数据库操作类型。仅适用于 Atlas 变更流源。

stream.source.ns

文档

可选的

包含 Atlas Stream Processing 源文档所在命名空间的文档。仅适用于 Atlas 变更流源。

stream.source.ns.db

字符串

可选的

Atlas Stream Processing 尝试操作的数据库名称。仅适用于 Atlas 变更流源。

对于 集合变更流数据库变更流 源,此值对所有文档都相同。对于 集群变更流 源,此值因源而异。

stream.source.ns.coll

字符串

可选的

Atlas Stream Processing 尝试操作的集合名称。仅适用于 Atlas 变更流源。

对于集合变更流源,所有文档的该值都相同。对于数据库变更流集群变更流源,该值可变。

stream.source.documentKey._id

文档

可选的

包含源文档的对象 ID 的文档。仅适用于 Atlas 变更流源。

stream.source.initialSync.phase

字符串

可选的

初始同步操作的当前状态。仅在初始同步期间适用于 Atlas 变更流源。

stream.source.kinesisStream

字符串

可选的

Atlas Stream Processing 源文档的 Kinesis 数据流的名称。仅适用于 AWS Kinesis 源。

stream.source.shardId

字符串

可选的

Kinesis Data Stream 中的分片 ID,Atlas Stream Processing 从中获取文档。仅适用于 Amazon Web Services Kinesis 源。

stream.source.sequenceNumber

字符串

可选的

来自 Kinesis Data Stream 的源文档的唯一标识符。仅适用于 Amazon Web Services Kinesis 源。

stream.source.partitionKey

字符串

可选的

源文档所属分区的唯一标识符。仅适用于 Amazon Web Services Kinesis 源。

stream.window

文档

可选的

包含窗口元数据的文档。仅当文档在窗口中处理时才应用。

stream.window.start

ISODate

可选的

窗口打开时间。仅当文档在窗口中处理时应用。

stream.window.end

ISODate

可选的

窗口关闭时间。仅当文档在窗口中处理时应用。

stream.window.partition

字符串

可选的

文档所属的会话窗口分区。仅当文档在会话窗口中处理时才应用。

stream.https

文档

可选的

包含 $https 阶段元数据的文档。仅当处理失败发生在 $https 阶段时才应用。

stream.https.url

字符串

可选的

$https 阶段的目标 URL。仅在处理失败发生在 $https 阶段时应用。

stream.https.method

字符串

可选的

$https 阶段使用的 HTTP 请求方法。仅在处理失败发生在 $https 阶段时应用。

stream.https.httpStatusCode

int

可选的

请求的 HTTP 响应状态代码。仅当处理失败发生在 $https 阶段时才应用。

stream.https.responseTimeMs

int

可选的

请求的响应时间(毫秒)。仅当处理失败发生在 $https 阶段时才应用。

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" }
}
}