定义
$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 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 必需 | |
| 字符串 | 必需 | 包含目标Apache Iceberg数据库的S3 存储桶的名称。 |
| 字符串 | 必需 | |
| 字符串 |表达式 | 必需 | 目标Apache Iceberg 表的名称。必须是字符串或计算结果为字符串的表达式。使用表达式为每个文档进行动态路由。 |
| 字符串 | 必需 | |
| 字符串 | 可选的 | 存储桶的 AWS地区。对于未在 AWS 上运行的流处理器是必需的。 |
| 字符串 | Optional | 确定每个输入文档要执行的操作的策略。
默认值为 |
| 字符串 | Optional | 在 默认值为 |
| 文档 | Optional | 分区规格。如果您未设置此字段, 必须是一份包含一个或多个键值对的文档。每个键必须是要执行分区转换的列的名称,每个值必须是要使用的分区转换。第一个分区转换必须针对 给定字段的分区转换值必须为以下值之一:
|
| 文档 | Optional | 定义要使用的 Iceberg 目录的文档。必须是包含值为 |
行为
如果使用 $iceberg 阶段,则该阶段必须是流处理器中的最后阶段。
Atlas Stream Processing 仅支持 SP10、SP30 和 SP50 流处理器的 $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 | 详情 | |
|---|---|---|
|
| |
|
| |
|
| |
|
| |
|
| |
|
| 十六进制编码 |
|
| 字符化 UUID |
|
| 不适用于 UUID |
|
| 以微秒为单位的 UTC 时间 |
|
| UTC 时间,测量到微秒 |
|
| 序列化为 Basic JSON string |
|
| 序列化为 Basic JSON string |
不支持其他 BSON types。Atlas Stream Processing 会将包含不支持的 BSON types 的文档发送到 DLQ。
示例
以下示例说明了 $iceberg 阶段的各种应用。
存档变更流
以下示例演示了如何以仅追加的方式写入Atlas数据库的初始内容和变更流写入Apache Iceberg 表,从而创建该数据库操作历史记录的持久性存档。此聚合有两个阶段:
$source阶段与 Atlas 数据库建立连接,具体针对db数据库中的orders集合。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。$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 表条目。它分为四个阶段:
$source阶段与 Atlas 数据库建立连接,具体针对db数据库中的orders集合。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。$match阶段对operationType进行筛选,以便仅处理具有有效操作类型声明的文档。$replaceRoot阶段根据操作类型更改文档根目录。对于删除操作,它会将文档根更改为文档的密钥。这会导致生成一条记录,表明文档已被删除,但会将其内容从进一步处理中排除。
对于所有其他操作,它会将文档根更改为
fullDocument,传递文档内容以供进一步处理,同时排除变更流元数据。
$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" } }
从多集合源到多表Apache Iceberg 目标
在以下示例中,Atlas Stream Processing 使用动态表达式将文档动态路由至各种输出目标。
$source阶段与 Atlas 数据库建立连接,具体针对位于db数据库中的集合a、b和c。它可以在处理器激活时启用初始同步以捕获数据库中的文档,并确保每个变更流事件都能捕获完整的文档。$match阶段会过滤operationType为"insert"、"update"、"delete"或"replace"之一的文档。$replaceRoot阶段根据操作类型更改文档根目录。对于删除操作,它会将文档根更改为文档的密钥。这会导致生成一条记录,表明文档已被删除,但会将其内容从进一步处理中排除。
对于所有其他操作,它会将文档根更改为
fullDocument,传递文档内容以供进一步处理,同时排除变更流元数据。
$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" } }