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.
Comparação com a initial sync de collection única
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 |
| O Atlas Stream Processing sincroniza a collection nomeada. |
Explicit list |
| O Atlas Stream Processing sincroniza cada collection na lista na ordem em que você as lista. |
Traços de sincronização inicial de múltiplas coleções
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 drenaaas partições de antesbde ebantescde, 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].targetStatsos 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.
Diferenças de comportamento a serem consideradas
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.parallelismo 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.
initialSyncnã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.
Configurar sincronização inicial de múltiplas coleções
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.
Pré-requisitos
Antes de concluir este procedimento, você deve ter:
Um espaço de trabalho de processamento de fluxo com conexões Atlas para os clusters de origem e destino. Para adicionar uma conexão, consulte Gerenciar conexões.
Uma camada de processador que ofereça suporte ao número de coleções que você planeja listar. Para saber mais, consulte Alocação de recursos.
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:
Carregue o conjunto de dados de amostra.
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 .
Ative a pré e publicação-imagens do fluxo de alterações.
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 .
Procedimento
Selecione a UI do Atlas ou para configurar o processador de mongosh fluxo.
Verificar o progresso da sincronização inicial
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') }
Exemplos
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.
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:
O estágio
$sourcesincronizaordersecustomersem<source-connection-name>e, em seguida, abre um fluxo de alterações com escopo dessas duas collections.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 ,
$replaceRootosstream.source.*metadados sobrevivem a esse estágio.O estágio grava cada evento na coleção
$mergecorrespondente<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.
Informações adicionais
Para aprender mais sobre os estágios e conceitos neste guia, consulte os seguintes recursos:
Para saber mais sobre o
$sourceestágio e seusinitialSynccampos, consulte MongoDB Collection Change Stream.Para saber como o nível do processador afeta o número de coleções que você pode sincronizar, consulte Alocação de recursos.
Para saber como o Atlas Stream Processing verifica uma sincronização em andamento, consulte Pontos de verificação.
Para saber como criar, iniciar, parar e monitorar processadores de fluxo, consulte Desenvolver e gerenciar processadores de fluxo.