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.
Menu Docs

Configurar sincronização inicial de múltiplas coleções

initialSync permite que um processador de fluxo do Atlas Stream Processing ingira documentos preexistentes em uma coleção do Atlas como se fossem inserir changeEvent documentos antes de começar a seguir o fluxo de alterações. Um único estágio agora pode $source executar initialSync em mais de uma collection ao mesmo tempo.

Esta página compara as características de sincronização de uma única collection e de várias collections, orienta a configuração de uma sincronização de várias collections e demonstra um pipeline de replicação entre clusters criado sobre ela.

O coll campo do estágio aceita um único nome de coleção ou uma array de nomes de coleção. Cada coleção que você nomeia na array deve pertencer ao banco de dados especificado $source no db campo. Não é possível sincronizar coleções de mais de um banco de dados em um único estágio. O valor fornecido determina o escopo da $source sincronização:

Escopo
$source Configuração
Comportamento

Coleção única

db: "app", coll: "orders"

O Atlas Stream Processing sincroniza a collection nomeada.

Explicit list

db: "app", coll: ["a", "b", "c"]

O Atlas Stream Processing sincroniza cada collection na lista na ordem em que você as lista.

Uma operação multi-collection initialSync tem as seguintes características:

  • Sincroniza coleções na ordem da lista. coll: ["a", "b", "c"]Com, o Atlas Stream Processing drena a as partições de antes b de e b antes c de, exceto quando o transbordamento preenche a capacidade ociosa com partições de uma collection posterior.

  • Retoma sem ressincronizar as coleções concluídas. O Atlas Stream Processing acompanha quais collections foram concluídas, portanto, um stream operators reiniciado ignora as collections concluídas e retoma qualquer collection em andamento do último documento sincronizado. Ele pode devolver documentos já emitidos, portanto, projete seu coletor para lidar com documentos repetidos de forma idempotente.

  • Expõe o progresso granular da sincronização no nível da collection. stats().stats.operatorStats[0].targetStats os relatórios sincronizam o progresso e as contagens de documento para cada coleção de destino. Para saber mais, consulte Verificar o progresso da sincronização inicial.

A multicollection initialSync difere da sincronização de collection única de maneiras que afetam o modo como você configura e opera o processador de stream:

  • O paralelismo é global, não por collection. A configuração limita initialSync.parallelism o número de leituras simultâneas de partição em todo o conjunto de coleções de destino, não por coleção.

  • O escopo da recuperaçãosegue o escopo da origem. Após a fase de sincronização, o Atlas Stream Processing abre um change stream que corresponde ao escopo da origem: um change stream no nível do banco de dados para uma origem de banco de dados inteiro, um change stream no nível do cluster para uma origem de todo o cluster ou um change stream filtrado para um lista explícita de coleções. Você não precisa configurar isso separadamente.

  • Uma falha de sincronização da coleção falha em todo o processador. Assim como na sincronização de coleção única, se o Atlas Stream Processing não sincronizar nenhuma coleção na lista, todo o processador de fluxo falhará.

  • Algumas configurações do lado da origem não são replicadas para o destino. initialSync não replica índices, opções de collection, como capped collections, coleções de séries temporais, validadores ou agrupamento padrão, ou definições de visualização. Crie-os no destino antes de iniciar o processador de fluxo. O Atlas Stream Processing também não descobre as collections que você cria após o início do processador de stream; o conjunto de coleções a serem sincronizadas é fixado nesse ponto.

O procedimento a seguir usa duas collections do conjunto de dados sample_analytics:, customers que contém perfis de cliente , e, que contém as contas financeiros accounts de cada cliente. O procedimento configura um processador de fluxo que copia ambas as coleções de um Atlas cluster para outro e, em seguida, as mantém atualizadas à medida que os documentos são alterados.

Antes de concluir este procedimento, você deve ter:

Execute os comandos nesta seção em seu cluster. Para se conectar, consulte Conectar-se a um cluster via mongosh.

Conclua as etapas a seguir para preparar as coleções de origem:

1

Este procedimento utiliza as customers accounts collections e do conjunto de dados sample_analytics. Para saber como carregar dados de amostra, consulte Importar dados de amostra para seu sistema do Atlas .

2

Ative as imagens pré e pós customers nas accounts collections e , porque o estágio $source fullDocument: "required" do procedimento define. Execute os seguintes comandos no cluster:

db.getSiblingDB("sample_analytics").runCommand({
collMod: "customers",
changeStreamPreAndPostImages: { enabled: true }
})
{ ok: 1, ... }
db.getSiblingDB("sample_analytics").runCommand({
collMod: "accounts",
changeStreamPreAndPostImages: { enabled: true }
})
{ ok: 1, ... }

Para saber mais fullDocument: "required" sobre, consulte Fluxo de alteração de coleção do MongoDB .

Selecione a UI do Atlas ou para configurar o processador de mongosh fluxo.

Depois de iniciar o processador de fluxo, chame sp.<processor-name>.stats() com a opção verbose para visualizar o progresso da sincronização de cada coleção de destino. A entrada do operador de origem em stats.operatorStats relata uma array targetStats com um elemento por coleção:

sp.replicate_analytics_sp.stats(
{ options: { verbose: true } }
).stats.operatorStats[0].targetStats
[
{
db: 'sample_analytics',
coll: 'customers',
inputMessageCount: Long('500'),
inputMessageSize: 672326,
initialSync: {
estimatedDocs: Long('500'),
copiedDocs: Long('251'),
status: 'in_progress'
}
},
{
db: 'sample_analytics',
coll: 'accounts',
inputMessageCount: Long('1746'),
inputMessageSize: 1883681,
initialSync: {
estimatedDocs: Long('1746'),
copiedDocs: Long('1746'),
status: 'completed'
}
}
]

O campo initialSync.status de cada destino passa por pending, in_progress e completed à medida que o Atlas Stream Processing sincroniza essa coleção ou relata failed se a sincronização falhar. estimatedDocs e copiedDocs informam as contagens de documento estimadas e concluídas para essa collection.

Para visualizar o progresso total em cada collection de destino, use o campoinitialSync do próprio operador de origem, que resume estimatedDocs e copiedDocs em todas as collections:

sp.replicate_analytics_sp.stats(
{ options: { verbose: true } }
).stats.operatorStats[0].initialSync
{ estimatedDocs: Long('2246'), copiedDocs: Long('1997') }

O procedimento anterior replica duas collections entre Atlas clusters. O exemplo a seguir aplica o mesmo pipeline a um conjunto diferente de collections. Somente os db coll valores e no estágio $source mudam.

Combine a sincronização inicial de várias coleções com o roteamento dinâmico de namespace para replicar um conjunto de coleções.

O pipeline sincroniza documentos existentes em orders e customers e, em seguida, continua a replicar as alterações em ambas as coleções à medida que ocorrem. A agregação tem três estágios:

  1. O estágio $source sincroniza orders e customers em <source-connection-name> e, em seguida, abre um fluxo de alterações com escopo dessas duas collections.

  2. O estágio seleciona a chave do documento para excluir eventos e o documento completo para todos os outros evento. Como os metadados não fazem parte do corpo do documento ,$replaceRoot osstream.source.* metadados sobrevivem a esse estágio.

  3. O estágio grava cada evento na coleção $merge correspondente <destination-connection-name> em, usando os metadados do namespace do evento de origem para rotear a gravação e o tipo de operação para escolher entre inserir, substituir ou excluir o documento de destino.

sp = db.createStreamProcessor("replicate-app-db", [
{ $source: {
connectionName: "<source-connection-name>",
db: "app",
coll: ["orders", "customers"],
config: { fullDocument: "required" },
initialSync: { enable: true }
} },
{ $replaceRoot: {
newRoot: {
$cond: {
if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] },
then: { $meta: "stream.source.documentKey" },
else: "$fullDocument"
}
}
} },
{ $merge: {
into: {
connectionName: "<destination-connection-name>",
db: { $meta: "stream.source.ns.db" },
coll: { $meta: "stream.source.ns.coll" }
},
on: "_id",
whenMatched: { $cond: {
if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] },
then: "delete",
else: "replace"
} },
whenNotMatched: { $cond: {
if: { $eq: [{ $meta: "stream.source.operationType" }, "delete"] },
then: "discard",
else: "insert"
} }
} }
]);
sp.start();

Antes de iniciar esse processador de fluxo, crie quaisquer índices, opções de coleção ou visualizações que orders e customers precisem no <destination-connection-name>. initialSync não os replica.

Para aprender mais sobre os estágios e conceitos neste guia, consulte os seguintes recursos: