Para agentes de IA: um índice de documentação está disponível em https://www.mongodb.com/pt-br/docs/llms.txt — as versões de markdown de todas as páginas estão disponíveis anexando .md a qualquer caminho de URL.
See how MongoDB 9.0 delivers up to 2x higher throughput.
MongoDB Branding Shape
Register now >
Menu Docs

$iceberg Estágio de agregação

O estágio $iceberg especifica uma conexão no Registro de Conexão com um bucket do AWS S3 onde você pode escrever dados em uma tabela do Apache icebergue.

$iceberg deve ser o último estágio de qualquer pipeline em que apareça. Você pode usar apenas um estágio $iceberg por pipeline.

O estágio de pipeline do $iceberg tem a seguinte forma de protótipo:

{
"$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"
}
}
}

O estágio $iceberg recebe um documento com os seguintes campos:

Campo
Tipo
necessidade
Descrição

connectionName

string

Obrigatório

Nome da conexão AWS S3 para usar para leituras e gravações. Isso deve corresponder ao nome da conexão no Registro de Conexão.

bucket

string

Obrigatório

Nome do bucket S3 que contém o banco de dados Apache icebergue de destino.

databaseName

string

Obrigatório

Nome do banco de dados Apache Iceberg que contém a tabela de destino.

tableName

corda | expressão

Obrigatório

Nome da tabela de destino Apache Iceberg. Deve ser uma string ou uma expressão que avalie para uma string. Use expressões para roteamento dinâmico por documento.

path

string

Obrigatório

Chave de prefixo do caminho para seu banco de dados Apache Iceberg.

region

string

Condicional

RegiãoAWS do bucket. Necessário para processadores de stream que não executam no AWS.

mode

string

Opcional

Estratégia para determinar qual operação executar por documento de entrada.

  • "cdc" faz com que o Atlas Stream Processing determine o tipo de operação lendo o campo de metadados stream.source.operationType.

  • "insert" faz com que o Atlas Stream Processing anexe cada documento à tabela de destino como uma nova linha, ignorando as declarações de tipo de operação no campo de metadados stream.source.operationType.

Padrão é cdc.

idFieldName

string

Opcional

Nome do campo e da coluna usado como chave de linha no modo cdc.

Padrão é "_id".

partitionedBy

documento

Opcional

Especificação de particionamento. Se você não definir esse campo, $iceberg definirá uma transformação de partição padrão para a coluna idFieldName.

Deve ser um documento contendo um ou mais pares de valores-chave. Cada chave deve ser o nome de uma coluna na qual executar uma transformação de partição, e cada valor deve ser a transformação de partição a ser usada. A primeira transformação de partição deve ser em idFieldName.

Para fazer a partição em um campo aninhado em uma coluna struct, especifique o caminho qualificado por ponto para o campo, por exemplo outer.inner.key_field. O caminho não pode incluir uma coluna list em qualquer nível.

O valor de transformação de partição de um determinado campo deve ser um dos seguintes:

  • "identity"

  • "year"

  • "month"

  • "day"

  • "hour"

  • { truncate: int }

  • { bucket: int }

Para obter mais informações sobre as transformações de partição do Apache Iceberg, consulte a documentação do Apache Iceberg.

catalog

documento

Opcional

Documento que define o catálogo Iceberg a ser usado. Deve ser um documento que contenha o campo type com um valor de "hadoop" ou "glue".

schemaInference

documento

Opcional

Documento que configura como $iceberg infere o esquema de tabela. Deve ser um documento contendo o campo mode com um dos seguintes valores:

  • "json" faz com que o Atlas Stream Processing grave campos object e array como colunas de string JSON.

  • "nested" faz com que o Atlas Stream Processing grave campos object como struct colunas e array campos como list colunas.

O padrão é "json". Para saber mais, consulte Conversão de tipo.

Se você usar o estágio $iceberg, ele deverá ser o último estágio do seu processador de fluxo.

O Atlas Stream Processing suporta o estágio $iceberg apenas para processadores de stream SP10, SP30 e SP50. Seu tier de processador determina o número máximo de tabelas suportadas para roteamento dinâmico:

Nível
Tabelas máximas

SP10

5

SP30

10

SP50

50

Importante

Se um processador de stream com roteamento dinâmico exceder o número máximo de tabelas suportadas em seu nível, o processador entrará em um estado FAILED. Para saber mais sobre as causas de falhas e recuperação, consulte Políticas de tratamento de erros e novas tentativas.

O estágio $iceberg infere o esquema da tabela Apache icebergue resultante a partir do esquema dos dados de saída do processador de fluxo. À medida que o Atlas Stream Processing observa novos campos no stream, incluindo campos dentro de colunas struct, o esquema da tabela desenvolve adequadamente.

Se você especificar uma tabela que ainda não existe, o Apache Iceberg criará a tabela quando receber a primeira mensagem destinada a ela.

O Atlas Stream Processing garante o processamento de pelo menos uma vez para saída para tabelas Apache Iceberg.

Você pode usar uma expressão dinâmica como o valor do campo tableName. Ao usar uma expressão dinâmica para capturar valores específicos do documento, você pode rotear documentos de entrada para diferentes tabelas de acordo com esses valores. A expressão deve ser avaliada como uma string. Para obter um exemplo, consulte Roteamento dinâmico. Para saber mais, consulte operadores de expressão.

Se você especificar um tópico com uma expressão dinâmica, mas o Atlas Stream Processing não puder avaliar a expressão de uma determinada mensagem, ele enviará essa mensagem para a fila de mensagens não entregues (DLQ), se ela estiver configurada, e processará as mensagens subsequentes. Se não houver uma fila de mensagens não entregues (DLQ) configurada, o Atlas Stream Processing ignorará completamente a mensagem e processará as mensagens subsequentes.

O Atlas Stream Processing executa a conversão de tipo de BSON para tipos primitivos Iceberg ao gravar em tabelas no estágio $iceberg.

BSON
Apache Iceberg Primitivo
Detalhes

string

string

int

int

long

long

double

double

bool

boolean

ObjectId

string

Codificado em hexadecimal

UUID

string

UUID stringificado

BinData

binary

Não se aplica ao UUID

date

timestamptz

Horário UTC, medido em microssegundos

timestamp

timestamptz

Horário UTC, medida em microssegundos

object

string ou struct

Serializado como uma string JSON básica por padrão. No modonested , torna-se uma coluna struct . Para saber mais, consulte Campos de objeto e array.

array

string ou list

Serializado como uma string JSON básica por padrão. No modonested , torna-se uma coluna list . Para saber mais, consulte Campos de objeto e array.

Outros BSON types não são suportados. O Atlas Stream Processing envia documentos com BSON types não suportados para o DLQ.

O campo schemaInference.mode determina como o Atlas Stream Processing grava os campos object e array nas tabelas do Apache icebergue:

  • No modojson, o Atlas Stream Processing serializa os campos object e array como strings JSON básicas e os grava em colunas string. json é o modo padrão.

  • No modonested, o Atlas Stream Processing infere tipos aninhados de icebergue a partir de campos object e array:

    • Cada campo object se torna uma coluna struct . O Atlas Stream Processing infere o tipo de cada campo na estrutura a partir do campo correspondente no objeto. Isso se aplica a todos os níveis de aninhamento.

    • Cada campo array se torna uma coluna list . O Atlas Stream Processing infere o tipo de elemento da lista a partir dos elementos não nulos da array. Os elementos primitivos devem ser todos do mesmo tipo. Se os elementos forem objetos, o tipo de elemento é um struct que contém os campos de todos os objetos.

O modo se aplica apenas a campos que não têm uma coluna correspondente na tabela de destino. Se a coluna já existir, o Atlas Stream Processing ignorará o modo e gravará o campo baseado no tipo de coluna:

  • Se a coluna for um string, o Atlas Stream Processing serializa o campo como uma string JSON básica.

  • Se a coluna for struct ou list, o Atlas Stream Processing gravará o campo como esse tipo aninhado.

Se a entrada contiver novos campos object ou array, o Atlas Stream Processing usará o modo para inferir os tipos das novas colunas.

Os exemplos a seguir ilustram vários aplicativos do estágio $iceberg.

O exemplo a seguir demonstra como gravar o conteúdo inicial e o changestream de um banco de dados Atlas em uma tabela Apache Iceberg de forma somente anexa, criando um arquivo durável do histórico operacional desse banco de dados. Esta agregação tem duas etapas:

  1. O estágio $source estabelece uma conexão com o banco de dados Atlas, visando especificamente a coleção orders no banco de dados db. Ele permite que a sincronização inicial capture documentos no banco de dados no momento da ativação do processador e garante que o documento completo seja capturado por cada evento de changestream.

  2. O estágio $iceberg estabelece uma conexão com o bucket AWS S3, gravando em uma tabela chamada myTable no caminho iceberg-warehouse/. Ao especificar somente as operações insert, ele garante um fluxo de escrita somente para anexação e no estilo de 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"
}
}

O exemplo a seguir demonstra como espelhar uma coleção do Atlas em sua totalidade para uma tabela Apache Iceberg.

Antes de definir a agregação, defina a seguinte variável:

const isDeleteExpr = {$eq: [{$meta: "stream.source.operationType"}, "delete"]};

A agregação a seguir adiciona, atualizar e excluir entradas de tabela Apache Iceberg em sincronização com as alterações na coleção de origem do Atlas. Tem quatro etapas:

  1. O estágio $source estabelece uma conexão com o banco de dados Atlas, visando especificamente a coleção orders no banco de dados db. Ele permite que a sincronização inicial capture documentos no banco de dados no momento da ativação do processador e garante que o documento completo seja capturado por cada evento de changestream.

  2. A etapa $match filtra operationType para que apenas documentos com declarações de tipo de operação válidas sejam processados.

  3. A fase $replaceRoot altera a raiz do documento dependendo do tipo de operação.

    • Para operações de exclusão, ele altera a raiz do documento para a chave do documento. Isso resulta em um registro de que o documento foi excluído, mas exclui seu conteúdo em processamento posterior.

    • Para todas as outras operações, ele altera a raiz do documento para o fullDocument, passando o conteúdo do documento para em processamento adicional, excluindo metadados de fluxo de alteração.

  4. O estágio $iceberg estabelece uma conexão com o bucket do AWS S3, gravando em uma tabela do Apache icebergue chamada myTable no caminho iceberg-warehouse/. No modocdc, esse estágio determina a operação a ser executada na tabela do Apache Geophyllum, lendo o campo de metadados 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"
}
}

No exemplo a seguir, o Atlas Stream Processing usa expressões dinâmicas para rotear dinamicamente documentos para vários destinos de saída.

  1. O estágio $source estabelece uma conexão com o banco de dados Atlas, visando especificamente as coleções a, b e c no banco de dados db. Ele permite que a sincronização inicial capture documentos no banco de dados no momento da ativação do processador e garante que o documento completo seja capturado por cada evento de changestream.

  2. O estágio $match filtra documentos cujo operationType é um de "insert", "update", "delete" ou "replace".

  3. A fase $replaceRoot altera a raiz do documento dependendo do tipo de operação.

    • Para operações de exclusão, ele altera a raiz do documento para a chave do documento. Isso resulta em um registro de que o documento foi excluído, mas exclui seu conteúdo em processamento posterior.

    • Para todas as outras operações, ele altera a raiz do documento para o fullDocument, passando o conteúdo do documento para em processamento adicional, excluindo metadados de fluxo de alteração.

  4. O estágio $iceberg estabelece uma conexão com o bucket do AWS S3 chamado myData, gravando em uma tabela Apache icebergue no caminho iceberg-warehouse/. Ele determina o nome da tabela de acordo com o nome da origem collection conforme recuperado dos metadados do documento . Ele também determina a operação a ser executada de acordo com os metadados do 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"
}
}