定义
$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 阶段采用包含以下字段的文档:
字段 | 类型 | 必要性 | 说明 |
|---|---|---|---|
| 字符串 | 必需 | |
| 字符串 | 必需 | 包含目标Apache Iceberg数据库的S3 存储桶的名称。 |
| 字符串 | 必需 | 包含目标表的 Apache Iceberg 数据库的名称。 |
| 字符串 |表达式 | 必需 | 目标Apache Iceberg 表的名称。必须是字符串或计算结果为字符串的表达式。使用表达式为每个文档进行动态路由。 |
| 字符串 | 必需 | |
| 字符串 | 可选的 | 存储桶的 AWS地区。对于未在 AWS 上运行的流处理器是必需的。 |
| 字符串 | Optional | 确定每个输入文档要执行的操作的策略。
默认值为 |
| 字符串 | Optional | 在 默认值为 |
| 文档 | Optional | 分区规格。如果您未设置此字段, 必须是一份包含一个或多个键值对的文档。每个键必须是要执行分区转换的列的名称,每个值必须是要使用的分区转换。第一个分区转换必须针对 要对嵌套在 给定字段的分区转换值必须为以下值之一:
有关Apache Iceberg 的更多信息 partition transforms, see the Apache Iceberg documentation. |
| 文档 | Optional | 定义要使用的 Iceberg 目录的文档。必须是包含值为 |
| 文档 | Optional |
行为
如果使用 $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 | 详情 | |
|---|---|---|
|
| |
|
| |
|
| |
|
| |
|
| |
|
| 十六进制编码 |
|
| 字符化 UUID |
|
| 不适用于 UUID |
|
| 以微秒为单位的 UTC 时间 |
|
| UTC 时间,测量到微秒 |
|
| 默认序列化为基本JSON字符串。在 |
|
| 默认序列化为基本JSON字符串。在 |
不支持其他 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 表,从而创建该数据库操作历史记录的持久性存档。此聚合有两个阶段:
$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 存储桶的连接,通过读取每个文档的stream.source.operationType元数据字段来写入Apache Iceberg 表。 table namedmyTablein theiceberg-warehouse/path. Incdcmode, 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" } }
从多集合源到多表Apache Iceberg Target
在以下示例中,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" } }