对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
See how MongoDB 9.0 delivers up to 2x higher throughput.
MongoDB Branding Shape
Register now >
Docs 菜单

$iceberg 聚合阶段

$iceberg 阶段指定连接注册表中与 AWS S3 存储桶的连接,您可以在该存储桶中将数据写入Apache Iceberg 表。

$iceberg 必须是其所在的任何管道的最后阶段。每个管道只能使用一个 $iceberg 阶段。

$iceberg 管道阶段采用以下原型形式:

{
"$iceberg": {
"connectionName": "<registered-connection>",
"bucket": "<target-bucket>",
"databaseName": "<database>",
"tableName": "<string>" | <expression>,
"path": "<key-prefix>",
"region": "<target-region>",
"mode": "cdc" | "insert",
"idFieldName": "<field-name>",
"partitionedBy": {
"<column-name>": "<partition-transform>",
. . .
},
"catalog": {
"type": "hadoop" | "glue"
},
"schemaInference": {
"mode": "json" | "nested"
}
}
}

$iceberg 阶段采用包含以下字段的文档:

字段
类型
必要性
说明

connectionName

字符串

必需

用于读取和写入的 AWS S3 连接的名称。它必须与连接注册表中的连接名称匹配。

bucket

字符串

必需

包含目标Apache Iceberg数据库的S3 存储桶的名称。

databaseName

字符串

必需

包含目标表的 Apache Iceberg 数据库的名称。

tableName

字符串 |表达式

必需

目标Apache Iceberg 表的名称。必须是字符串或计算结果为字符串的表达式。使用表达式为每个文档进行动态路由。

path

字符串

必需

region

字符串

可选的

存储桶的 AWS地区。对于未在 AWS 上运行的流处理器是必需的。

mode

字符串

Optional

确定每个输入文档要执行的操作的策略。

  • "cdc" 导致 Atlas Stream Processing 通过读取 stream.source.operationType 元数据字段来确定操作类型。

  • "insert" 导致 Atlas Stream Processing 将每个文档作为新行添加到目标表中,忽略 stream.source.operationType 元数据字段中的操作类型声明。

默认值为 cdc。

idFieldName

字符串

Optional

在 cdc 模式下用作行键的字段和列名称。

默认值为 "_id"。

partitionedBy

文档

Optional

分区规格。如果您未设置此字段,$iceberg 将为 idFieldName 列设置默认分区转换。

必须是一份包含一个或多个键值对的文档。每个键必须是要执行分区转换的列的名称,每个值必须是要使用的分区转换。第一个分区转换必须针对 idFieldName。

要对嵌套在 struct 列内的字段进行分区,请指定该字段的点限定路径,示例outer.inner.key_field。该路径在任何级别都不能包含 list 列。

给定字段的分区转换值必须为以下值之一:

  • "identity"

  • "year"

  • "month"

  • "day"

  • "hour"

  • { truncate: int }

  • { bucket: int }

有关Apache Iceberg 的更多信息 partition transforms, see the Apache Iceberg documentation.

catalog

文档

Optional

定义要使用的 Iceberg 目录的文档。必须是包含值为type "hadoop"或 的"glue" 字段的文档。

schemaInference

文档

Optional

配置 $iceberg 如何推断表模式的文档。必须是包含具有以下值之一的 mode字段的文档:

  • "json" 导致Atlas Stream Processing将 object 和 array 字段写入JSON字符串列。

  • "nested" 导致Atlas Stream Processing将 object 字段写入struct 列,并将 array 字段写入 list 列。

默认为 "json"。要学习;了解更多信息,请参阅类型转换。

如果使用 $iceberg 阶段,则该阶段必须是流处理器中的最后阶段。

Atlas Stream Processing 仅支持 SP10、SP30 和 SP50 流处理器的 $iceberg 阶段。您的处理器层级决定动态路由支持的最大表数:

层级
最大表

SP 10

5

SP 30

10

SP 50

50

重要

如果具有动态路由的流处理器超过其层级支持的最大表数,则该处理器将进入 FAILED 状态。要学习;了解故障原因和恢复,请参阅错误处理和重试策略。

$iceberg 阶段从流处理器输出数据的模式推断生成的Apache Iceberg 表的模式。当Atlas Stream Processing观察到流中的新字段(包括 struct 列中的字段)时,表模式也会相应地发展。

如果您指定的表尚不存在,则Apache Iceberg会在收到第一条指向该表的消息时创建该表。

Atlas Stream ProcessingACID 一致性保证对Apache Iceberg 表的输出至少进行一次处理。

您可以使用 动态表达式 作为 tableName 字段的值。通过使用动态表达式捕获文档特定值,您可以根据这些值将输入文档路由到不同的表中。表达式的结果必须为字符串。有关示例,请参阅 动态路由。要了解详情,请参阅 表达式操作符。

如果您使用动态表达式指定主题,但 Atlas Stream Processing 无法评估给定消息的表达式,则 Atlas Stream Processing 会将该消息发送到 死信队列(DLQ)(如果已配置)并处理后续消息。如果未配置 死信队列(DLQ),则 Atlas Stream Processing 会完全跳过该消息并处理后续消息。

在 $iceberg 阶段写入表时, Atlas Stream Processing会执行从BSON开始的类型转换。 to Iceberg primitive types

BSON
详情

string

string

int

int

long

long

double

double

bool

boolean

ObjectId

string

十六进制编码

UUID

string

字符化 UUID

BinData

binary

不适用于 UUID

date

timestamptz

以微秒为单位的 UTC 时间

timestamp

timestamptz

UTC 时间,测量到微秒

object

string or struct

默认序列化为基本JSON字符串。在 nested模式下,成为 struct 列。要学习;了解更多信息,请参阅对象和数组字段。

array

string or list

默认序列化为基本JSON字符串。在 nested模式下,成为 list 列。要学习;了解更多信息,请参阅对象和数组字段。

不支持其他 BSON types。Atlas Stream Processing 会将包含不支持的 BSON types 的文档发送到 DLQ。

schemaInference.mode字段确定Atlas Stream Processing如何将 object 和 array 字段写入Apache Iceberg 表:

  • 在 json模式下, Atlas Stream Processing将 object 和 array 字段序列化为基本JSON字符串,并将它们写入 string 列。 json 是默认模式。

  • 在 nested模式下, Atlas Stream Processing会从 object 和 array 字段推断 Iceberg 嵌套类型:

    • 每个 object字段都成为 struct 列。 Atlas Stream Processing会根据对象中的相应字段推断结构体中每个字段的类型。这适用于每个嵌套级别。

    • 每个 array字段都成为 list 列。 Atlas Stream Processing从大量的非空元素推断列表元素类型。原始元素必须都是同一类型。如果元素是对象,则元素类型是包含所有对象字段的 struct。

该模式仅适用于目标表中没有匹配列的字段。如果列已经存在, Atlas Stream Processing会忽略模式并根据列类型写入字段:

  • 如果列是 string, Atlas Stream Processing会将该字段序列化为 Basic JSON字符串。

  • 如果列是 struct 或 list, Atlas Stream Processing会将该字段写入该嵌套类型。

如果输入包含新的 object 或 array 字段, Atlas Stream Processing会使用该模式来推断新列的类型。

以下示例说明了 $iceberg 阶段的各种应用。

以下示例演示了如何以仅追加的方式写入Atlas数据库的初始内容和变更流写入Apache Iceberg 表,从而创建该数据库操作历史记录的持久性存档。此聚合有两个阶段:

  1. $source 阶段与 Atlas 数据库建立连接,具体针对 db 数据库中的 orders 集合。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。

  2. $iceberg 阶段建立与 AWS S3 存储桶的连接,写入 iceberg-warehouse/ 路径中名为 myTable 的表。通过仅指定 insert 操作,可确保仅追加的日志式写入流。

{
"$source": {
"connectionName": "atlas1",
"db": "db",
"coll": "orders",
"initialSync": {
"enable": true
},
"config": {
"fullDocument": "required"
}
},
"$iceberg": {
"connectionName": "myS3Connection",
"bucket": "myData",
"path": "iceberg-warehouse/",
"tableName": "myTable",
"mode": "insert"
}
}

以下示例演示了如何将整个Atlas集合镜像到Apache Iceberg 表。

在定义聚合之前,设置以下变量:

const isDeleteExpr = {$eq: [{$meta: "stream.source.operationType"}, "delete"]};

以下聚合与Atlas源集合的更改同步添加、更新和删除Apache Iceberg 表条目。它分为四个阶段:

  1. $source 阶段与 Atlas 数据库建立连接,具体针对 db 数据库中的 orders 集合。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。

  2. $match 阶段对 operationType 进行**过滤器**,以便仅处理具有有效操作类型声明的文档。

  3. $replaceRoot 阶段会根据操作类型更改文档根。

    • 对于删除操作,它会将文档根更改为文档的密钥。这会导致生成一条记录,表明文档已被删除,但会将其内容从进一步处理中排除。

    • 对于所有其他操作,它会将文档根更改为 fullDocument,传递文档内容以供进一步处理,同时排除变更流元数据。

  4. $iceberg 阶段建立与 AWS S3 存储桶的连接,通过读取每个文档的 stream.source.operationType元数据字段来写入Apache Iceberg 表。 table named myTable in the iceberg-warehouse/ path. In cdc mode, this stage determines the operation to perform against the Apache Iceberg

{
"$source": {
"connectionName": "atlas1",
"db": "db",
"coll": "orders",
"initialSync": {
"enable": true
},
"config": {
"fullDocument": "required"
}
},
"$match": {
"operationType": {
"$in": ["insert", "update", "delete", "replace"]
}
},
"$replaceRoot": {
"newRoot": {
"$cond": {
"if": isDeleteExpr,
"then": "$documentKey",
"else": "$fullDocument"
}
}
}
"$iceberg": {
"connectionName": "myS3Connection",
"bucket": "myData",
"path": "iceberg-warehouse/",
"tableName": "myTable",
"mode": "cdc"
}
}

在以下示例中,Atlas Stream Processing 使用动态表达式将文档动态路由至各种输出目标。

  1. $source 阶段与 Atlas 数据库建立连接,具体针对位于 db 数据库中的集合 a、b 和 c。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。

  2. $match 阶段会过滤 operationType 为 "insert"、"update"、"delete" 或 "replace" 之一的文档。

  3. $replaceRoot 阶段会根据操作类型更改文档根。

    • 对于删除操作,它会将文档根更改为文档的密钥。这会导致生成一条记录,表明文档已被删除,但会将其内容从进一步处理中排除。

    • 对于所有其他操作,它会将文档根更改为 fullDocument,传递文档内容以供进一步处理,同时排除变更流元数据。

  4. $iceberg 阶段建立与名为 myData 的 AWS S3 存储桶的连接,写入 iceberg-warehouse/ 路径中的Apache Iceberg 表。它根据从文档元数据中检索的源 collection 的名称来确定表的名称。它还根据文档元数据确定要执行的操作。

{
"$source": {
"connectionName": "atlas1",
"db": "db",
"coll": ["a", "b", "c"],
"initialSync": {
"enable": true
},
"config": {
"fullDocument": "required"
}
}
},
{
"$match": {
"operationType": {
"$in": ["insert", "update", "delete", "replace"]
}
}
},
{
"$replaceRoot": {
"newRoot": {
"$cond": {
"if": {
"$eq": [{ "$meta": "stream.source.operationType" }, "delete"]
},
"then": "$documentKey",
"else": "$fullDocument"
}
}
}
},
{
"$iceberg": {
"connectionName": "myS3Connection",
"databaseName": "iceberg-db",
"bucket": "myData",
"path": "iceberg-warehouse/",
"tableName": { "$meta": "stream.source.ns.coll" },
"mode": "cdc"
}
}