对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
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"
}
}
}

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

字段
类型
必要性
说明

connectionName

字符串

必需

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

bucket

字符串

必需

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

databaseName

字符串

必需

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

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

  • "identity"

  • "year"

  • "month"

  • "day"

  • "hour"

  • { truncate: int }

  • { bucket: int }

有关Apache分区转换的更多信息,请参阅Apache文档。

catalog

文档

Optional

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

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

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

层级
最大表

SP 10

5

SP 30

10

SP 50

50

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

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

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

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

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

$iceberg 阶段写入表时, Atlas Stream Processing会执行从BSON到 Iceberg 基元类型的类型转换。

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

序列化为 Basic JSON string

array

string

序列化为 Basic JSON string

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

以下示例说明了 $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 存储桶的连接,写入 iceberg-warehouse/ 路径中名为 myTable 的Apache Iceberg 表。在 cdc模式下,此阶段通过读取每个文档的 stream.source.operationType元数据字段来确定对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 数据库中的集合 abc。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。

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