Definición
La etapa $iceberg especifica una conexión en el Registro de conexiones a un bucket S3 de AWS donde puede escribir datos en una tabla de Apache Iceberg.
Ubicación
$iceberg debe ser la última etapa de cualquier pipeline en la que aparezca. Solo puede utilizar una etapa $iceberg por pipeline.
Sintaxis
La etapa de pipeline $iceberg tiene la siguiente forma de prototipo:
{ "$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" } } }
La etapa $iceberg procesa un documento con los siguientes campos:
Campo | Tipo | Necesidad | Descripción |
|---|---|---|---|
| string | Requerido | Nombre de la conexión AWS S3 que se utilizará para lecturas y escrituras. Este nombre debe coincidir con el nombre de la conexión en el Registro de conexiones. |
| string | Requerido | Nombre del bucket S3 que contiene la base de datos Apache Iceberg de destino. |
| string | Requerido | Nombre de la base de datos de Apache Iceberg que contiene la tabla de destino. |
| string | expresión | Requerido | Nombre de la tabla de destino de Apache Iceberg. Debe ser una cadena de texto o una expresión que se evalúe como una cadena de texto. Utilice expresiones para el enrutamiento dinámico por documento. |
| string | Requerido | Clave de prefijo de la ruta a su base de datos Apache Iceberg. |
| string | Condicional | Región deAWS del bucket. Requerido para los procesadores de flujo que no se ejecutan en AWS. |
| string | Opcional | Estrategia para determinar qué operación realizar por documento de entrada.
Se establece por defecto en |
| string | Opcional | Campo y nombre de columna utilizados como clave de fila en el modo Se establece por defecto en |
| Documento | Opcional | Especificación de particionamiento. Si no establece este campo, Debe ser un documento que contenga uno o más pares clave-valor. Cada clave debe ser el nombre de una columna en la que se realizará una transformación de partición, y cada valor debe ser la transformación de partición que se utilizará. La primera transformación de partición debe realizarse en Para particionar por un campo anidado dentro de una columna El valor de transformación de partición de un campo determinado debe ser uno de los siguientes:
Para obtener más información sobre las transformaciones de particiones de Apache Iceberg, consulte la documentación de Apache Iceberg. |
| Documento | Opcional | Documento que define el catálogo Iceberg a utilizar. Debe ser un documento que contenga el |
| Documento | Opcional | Documento que configura cómo
Por defecto es |
Comportamiento
Si utiliza la etapa $iceberg, debe ser la última etapa de su procesador de transmisión.
Atlas Stream Processing es compatible con la etapa $iceberg solo para los procesadores de transmisión SP10, SP30 y SP50. Su nivel de procesador determina el número máximo de tablas compatibles con el enrutamiento dinámico:
Nivel | Tablas máximas |
|---|---|
SP10 | 5 |
SP30 | 10 |
SP50 | 50 |
Importante
Si un procesador de flujo con enrutamiento dinámico supera el número máximo de tablas admitidas en su nivel, el procesador entra en el estado FAILED. Para obtener información sobre las causas de los fallos y la recuperación, consulte la sección sobre el manejo de errores y las políticas de reintento.
La etapa $iceberg infiere el esquema de la tabla resultante de Apache Iceberg a partir del esquema de los datos de salida del procesador de flujo. A medida que Atlas Stream Processing observa nuevos campos en el flujo, incluidos los campos dentro de las columnas struct, el esquema de la tabla evoluciona en consecuencia.
Si se especifica una tabla que no existe aún, Apache Iceberg crea la tabla cuando recibe el primer mensaje que la dirige.
Atlas Stream Processing garantiza el procesamiento al menos una vez para la salida a las tablas de Apache Iceberg.
Puede utilizar una expresión dinámica como valor del campo tableName. Al utilizar una expresión dinámica para capturar valores específicos del documento, puede enrutar los documentos de entrada a diferentes tablas según estos valores. La expresión debe evaluarse como una string. Para obtener un ejemplo, consulte Enrutamiento dinámico. Para obtener más información, consulte operadores de expresión.
Si especifica un tema con una expresión dinámica, pero Atlas Stream Processing no puede evaluar la expresión para un mensaje determinado, Atlas Stream Processing envía ese mensaje a la fila de letra muerta, si está configurada, y procesa los mensajes subsiguientes. Si no hay ninguna fila de letra muerta configurada, Atlas Stream Processing omite el mensaje por completo y procesa los mensajes subsiguientes.
Conversión de tipo
Atlas Stream Processing realiza una conversión de tipos de BSON a tipos primitivos Iceberg al escribir en tablas en la $iceberg etapa.
BSON | Apache Iceberg Primitivo | Detalles |
|---|---|---|
|
| |
|
| |
|
| |
|
| |
|
| |
|
| Codificado en hexadecimal |
|
| UUID en cadena |
|
| No se aplica a UUID |
|
| Hora UTC, medida en microsegundos |
|
| Hora UTC, medida en microsegundos |
|
| Se serializa como una cadena JSON básica por defecto. En el modo |
|
| Se serializa como una cadena JSON básica por defecto. En el modo |
No se admiten otros BSON types. Atlas Stream Processing envía documentos con BSON types no admitidos a la DLQ.
Campos de objetos y matrices
El campo schemaInference.mode determina cómo Atlas Stream Processing escribe los campos object y array en las tablas de Apache Iceberg:
En el modo
json, Atlas Stream Processing serializa los camposobjectyarraycomo cadenas JSON básicas y los escribe en las columnasstring. El modo predeterminado esjson.En el modo
nested, Atlas Stream Processing infiere los tipos anidados de Iceberg a partir de los camposobjectyarray:Cada campo
objectse convierte en una columnastruct. Atlas Stream Processing infiere el tipo de cada campo en la estructura a partir del campo correspondiente en el objeto. Esto se aplica en todos los niveles de anidamiento.Cada campo
arrayse convierte en una columnalist. Atlas Stream Processing infiere el tipo de elemento de la lista a partir de los elementos no nulos del array. Todos los elementos primitivos deben ser del mismo tipo. Si los elementos son objetos, el tipo de elemento es unstructque contiene los campos de todos los objetos.
Este modo se aplica únicamente a los campos que no tienen una columna coincidente en la tabla de destino. Si la columna ya existe, Atlas Stream Processing ignora el modo y escribe el campo según el tipo de columna.
Si la columna es
string, Atlas Stream Processing serializa el campo como una cadena JSON básica.Si la columna es de tipo
structolist, Atlas Stream Processing escribe el campo como ese tipo anidado.
Si la entrada contiene nuevos campos object o array, Atlas Stream Processing utiliza el modo para inferir los tipos de las nuevas columnas.
Ejemplos
Los siguientes ejemplos ilustran varias aplicaciones de la etapa $iceberg.
Archivar un Changestream
El siguiente ejemplo demuestra cómo escribir el contenido inicial y el flujo de cambios de una base de datos Atlas en una tabla de Apache Iceberg de forma de solo escritura, creando un archivo duradero del historial operativo de esa base de datos. Esta agregación tiene dos etapas:
La etapa
$sourceestablece una conexión con la base de datos de Atlas, específicamente dirigida a la colecciónordersen la base de datosdb. Permite que la sincronización inicial capture documentos en la base de datos en el momento de la activación del procesador y garantiza que cada evento de flujo de cambios capture el documento completo.La etapa
$icebergestablece una conexión con el bucket S3 de AWS, escribiendo en una tabla llamadamyTableen la rutaiceberg-warehouse/. Al especificar solo operacionesinsert, garantiza un flujo de escritura de solo adición, estilo registro.
{ "$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" } }
Reflejar una colección
El siguiente ejemplo muestra cómo replicar una colección de Atlas en su totalidad en una tabla de Apache Iceberg.
Antes de definir la agregación, establezca la siguiente variable:
const isDeleteExpr = {$eq: [{$meta: "stream.source.operationType"}, "delete"]};
La siguiente agregación agrega, actualiza y elimina entradas de la tabla Apache Iceberg en sincronía con los cambios en la colección de origen de Atlas. Consta de cuatro etapas:
La etapa
$sourceestablece una conexión con la base de datos de Atlas, específicamente dirigida a la colecciónordersen la base de datosdb. Permite que la sincronización inicial capture documentos en la base de datos en el momento de la activación del procesador y garantiza que cada evento de flujo de cambios capture el documento completo.La etapa
$matchfiltra enoperationTypepara que solo se procesen los documentos con declaraciones de tipo de operación válidas.La etapa
$replaceRootcambia la raíz del documento según el tipo de operación.Para las operaciones de borrado, cambia la raíz del documento a la clave del documento. Esto da como resultado un registro de que el documento se borró, pero excluye su contenido de un procesamiento posterior.
Para todas las demás operaciones, cambia la raíz del documento a
fullDocument, pasando el contenido del documento para su posterior procesamiento mientras excluye los metadatos de changestream.
La etapa
$icebergestablece una conexión con el bucket S3 de AWS, escribiendo en una tabla de Apache Iceberg llamadamyTableen la rutaiceberg-warehouse/. En el modocdc, esta etapa determina la operación a realizar sobre la tabla de Apache Iceberg leyendo el campo de metadatosstream.source.operationTypede cada documento.
{ "$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" } }
Desde una fuente de múltiples colecciones hasta un objetivo Apache Iceberg de múltiples tablas
En el siguiente ejemplo, Atlas Stream Processing utiliza expresiones dinámicas para enrutar dinámicamente documentos a varios destinos de salida.
La etapa
$sourceestablece una conexión con la base de datos de Atlas, específicamente dirigida a las coleccionesa,bycen la base de datosdb. Permite que la sincronización inicial capture documentos en la base de datos en el momento de la activación del procesador y garantiza que cada evento de flujo de cambios capture el documento completo.La etapa
$matchfiltra los documentos cuyooperationTypees uno de"insert","update","delete"o"replace".La etapa
$replaceRootcambia la raíz del documento según el tipo de operación.Para las operaciones de borrado, cambia la raíz del documento a la clave del documento. Esto da como resultado un registro de que el documento se borró, pero excluye su contenido de un procesamiento posterior.
Para todas las demás operaciones, cambia la raíz del documento a
fullDocument, pasando el contenido del documento para su posterior procesamiento mientras excluye los metadatos de changestream.
La etapa
$icebergestablece una conexión con el bucket S3 de AWS denominadomyData, escribiendo en una tabla de Apache Iceberg en la rutaiceberg-warehouse/. Determina el nombre de la tabla según el nombre de la fuentecollection, obtenido de los metadatos del documento. Asimismo, determina la operación a realizar según los metadatos del documento.
{ "$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" } }