Definición
La etapa $iceberg especifica una conexión en el Registro de conexiones a un bucket de AWS S3 donde puede guardar 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" } } }
La etapa $iceberg procesa un documento con los siguientes campos:
Campo | Tipo | Necesidad | Descripción |
|---|---|---|---|
| string | Requerido | Nombre de la conexión de Amazon Web Services S3 que se utilizará para lecturas y guardar. Esto 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 de destino Apache Iceberg. |
| string | Requerido | Nombre de la base de datos Apache Iceberg que contiene la tabla de destino. |
| string | expresión | Requerido | Nombre de la tabla de destino Apache Iceberg. Debe ser una string o una expresión que se evalúe como una string. Utiliza expresiones para el enrutamiento dinámico por documento. |
| string | Requerido | Prefijo de la clave de la ruta a su base de datos Apache Iceberg. |
| string | Condicional | Región de Amazon Web Services del bucket. Obligatorio para los procesadores de streaming que no se ejecutan en Amazon Web Services. |
| 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 contra 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 ser contra 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 de Iceberg que se va a utilizar. Debe ser un documento que contenga el campo |
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 |
La etapa $iceberg infiere el esquema de la tabla Apache Iceberg resultante a partir del esquema de los datos de salida del procesador de transmisión. A medida que Atlas Stream Processing observa nuevos campos en la transmisión, el esquema de la tabla evoluciona en consecuencia.
Si especifica una tabla que aún no existe, 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 la conversión de tipos de BSON a tipos primitivos de Iceberg al escribir en tablas en la etapa $iceberg.
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 |
|
| Serializado como una string JSON básica |
|
| Serializado como una string JSON básica |
No se admiten otros BSON types. Atlas Stream Processing envía documentos con BSON types no admitidos a la DLQ.
Ejemplos
Los siguientes ejemplos ilustran varias aplicaciones de la etapa $iceberg.
Archivar un Changestream
El siguiente ejemplo muestra cómo guardar el contenido inicial y el changestream de una base de datos de Atlas en una tabla de Apache Iceberg de solo anexar, creando un fichero 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 AWS S3, escribiendo en una tabla denominadamyTableen la rutaiceberg-warehouse/. Al especificar solo operacionesinsert, garantiza un flujo de escritura de solo anexar y 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 demuestra cómo reflejar 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 tabla de Apache Iceberg en sincronización con los cambios en la colección de origen de Atlas. Tiene 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 AWS S3, escribiendo en una tabla Apache Iceberg denominadamyTableen la rutaiceberg-warehouse/. En el modocdc, esta etapa determina la operación que se realizará en la tabla Apache Iceberg leyendo del 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" } }
De origen de varias colecciones a destino de varias tablas Apache Iceberg
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 AWS S3 denominadomyData, escribiendo en una tabla Apache Iceberg en la rutaiceberg-warehouse/. Determina el nombre de la tabla según el nombre del origencollectiontal como se recupera de los metadatos del documento. También 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" } }