Para agentes de IA: hay un índice de documentación disponible en https://www.mongodb.com/es/docs/llms.txt — versiones en markdown de todas las páginas están disponibles agregando .md a cualquier ruta URL.
Docs Menu

$iceberg etapa de agregació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.

$iceberg debe ser la última etapa de cualquier pipeline en la que aparezca. Solo puede utilizar una etapa $iceberg por pipeline.

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

connectionName

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.

bucket

string

Requerido

Nombre del bucket S3 que contiene la base de datos Apache Iceberg de destino.

databaseName

string

Requerido

Nombre de la base de datos de Apache Iceberg que contiene la tabla de destino.

tableName

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.

path

string

Requerido

Clave de prefijo de la ruta a su base de datos Apache Iceberg.

region

string

Condicional

Región deAWS del bucket. Requerido para los procesadores de flujo que no se ejecutan en AWS.

mode

string

Opcional

Estrategia para determinar qué operación realizar por documento de entrada.

  • "cdc" hace que Atlas Stream Processing determine el tipo de operación leyendo el campo de metadatos stream.source.operationType.

  • "insert" hace que Atlas Stream Processing agregue cada documento a la tabla de destino como una nueva fila, ignorando las declaraciones de tipo de operación en el campo de metadatos stream.source.operationType.

Se establece por defecto en cdc.

idFieldName

string

Opcional

Campo y nombre de columna utilizados como clave de fila en el modo cdc.

Se establece por defecto en "_id".

partitionedBy

Documento

Opcional

Especificación de particionamiento. Si no establece este campo, $iceberg establece una transformación de partición por defecto para la columna idFieldName.

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 idFieldName.

Para particionar por un campo anidado dentro de una columna struct, especifique la ruta al campo separada por puntos, por ejemplo, outer.inner.key_field. La ruta no puede incluir una columna list en ningún nivel.

El valor de transformación de partición de un campo determinado debe ser uno de los siguientes:

  • "identity"

  • "year"

  • "month"

  • "day"

  • "hour"

  • { truncate: int }

  • { bucket: int }

Para obtener más información sobre las transformaciones de particiones de Apache Iceberg, consulte la documentación de Apache Iceberg.

catalog

Documento

Opcional

Documento que define el catálogo Iceberg a utilizar. Debe ser un documento que contenga el type campo con un valor de "hadoop" "glue"o.

schemaInference

Documento

Opcional

Documento que configura cómo $iceberg infiere el esquema de la tabla. Debe ser un documento que contenga el campo mode con uno de los siguientes valores:

  • "json" hace que Atlas Stream Processing escriba los campos object y array como columnas de cadena JSON.

  • "nested" hace que Atlas Stream Processing escriba los campos object como columnas struct y los campos array como columnas list.

Por defecto es "json". Para obtener más información, consulte Conversión de tipos.

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.

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

string

string

int

int

long

long

double

double

bool

boolean

ObjectId

string

Codificado en hexadecimal

UUID

string

UUID en cadena

BinData

binary

No se aplica a UUID

date

timestamptz

Hora UTC, medida en microsegundos

timestamp

timestamptz

Hora UTC, medida en microsegundos

object

string or struct

Se serializa como una cadena JSON básica por defecto. En el modo nested, se convierte en una columna struct. Para obtener más información, consulte Campos de objeto y matriz.

array

string or list

Se serializa como una cadena JSON básica por defecto. En el modo nested, se convierte en una columna list. Para obtener más información, consulte Campos de objeto y matriz.

No se admiten otros BSON types. Atlas Stream Processing envía documentos con BSON types no admitidos a la DLQ.

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 campos object y array como cadenas JSON básicas y los escribe en las columnas string. El modo predeterminado es json.

  • En el modo nested, Atlas Stream Processing infiere los tipos anidados de Iceberg a partir de los campos object y array:

    • Cada campo object se convierte en una columna struct. 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 array se convierte en una columna list. 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 un struct que 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 struct o list, 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.

Los siguientes ejemplos ilustran varias aplicaciones de la etapa $iceberg.

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:

  1. La etapa $source establece una conexión con la base de datos de Atlas, específicamente dirigida a la colección orders en la base de datos db. 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.

  2. La etapa $iceberg establece una conexión con el bucket S3 de AWS, escribiendo en una tabla llamada myTable en la ruta iceberg-warehouse/. Al especificar solo operaciones insert, 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"
}
}

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:

  1. La etapa $source establece una conexión con la base de datos de Atlas, específicamente dirigida a la colección orders en la base de datos db. 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.

  2. La etapa $match filtra en operationType para que solo se procesen los documentos con declaraciones de tipo de operación válidas.

  3. La etapa $replaceRoot cambia 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.

  4. La etapa $iceberg establece una conexión con el bucket S3 de AWS, escribiendo en una tabla de Apache Iceberg llamada myTable en la ruta iceberg-warehouse/. En el modo cdc, esta etapa determina la operación a realizar sobre la tabla de Apache Iceberg leyendo el campo de metadatos stream.source.operationType de 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"
}
}

En el siguiente ejemplo, Atlas Stream Processing utiliza expresiones dinámicas para enrutar dinámicamente documentos a varios destinos de salida.

  1. La etapa $source establece una conexión con la base de datos de Atlas, específicamente dirigida a las colecciones a, b y c en la base de datos db. 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.

  2. La etapa $match filtra los documentos cuyo operationType es uno de "insert", "update", "delete" o "replace".

  3. La etapa $replaceRoot cambia 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.

  4. La etapa $iceberg establece una conexión con el bucket S3 de AWS denominado myData, escribiendo en una tabla de Apache Iceberg en la ruta iceberg-warehouse/. Determina el nombre de la tabla según el nombre de la fuente collection, 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"
}
}