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

$source Estágio (Stream Processing)

$source

O estágio especifica $source uma conexão no Registro de conexão para transmitir dados. Os seguintes tipos de conexão são suportados:

  • Apache Kafka agente

  • Fluxo de alterações da collection do MongoDB

  • Fluxo de alteração do banco de dados MongoDB

  • Fluxo de alteração do cluster MongoDB

  • Fluxo de dados do AWS Kinesis

  • array de documentos

  • Cronograma

Para operar em dados de transmissão de um intermediário do Apache Kafka, o $source estágio possui o seguinte formato de protótipo:

{
"$source": {
"connectionName": "<registered-connection>",
"topic" : ["<source-topic>", ...],
"timeField": {
$toDate | $dateFromString: <expression>
},
"partitionIdleTimeout": {
"size": <duration-number>,
"unit": "<duration-unit>"
},
"schemaRegistry": {
"connectionName": "<schema-registry-name>",
},
"config": {
"auto_offset_reset": "<start-event>",
"group_id": "<group-id>",
"keyFormat": "<deserialization-type>",
"keyFormatError": "<error-handling>"
},
}
}

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

Campo
Tipo
necessidade
Descrição

connectionName

string

Obrigatório

Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados.

topic

cadeia de caracteres ou matriz de cadeias de caracteres

Obrigatório

Nome de um ou mais tópicos do Apache Kafka a partir dos quais transmitir mensagens. Se você deseja transmitir mensagens de mais de um tópico, especifique-os em uma array.

timeField

documento

Opcional

documento que define um carimbo de data/hora oficial para mensagens recebidas.

Se você usar timeField, deverá defini-lo como um dos seguintes:

  • uma : expressão:$toDate que utiliza um campo de mensagem de origem como um argumento.

  • uma : expressão:$dateFromString que utiliza um campo de mensagem de origem como um argumento.

Se você não declarar um timeField, o Atlas Stream Processing criará um carimbo de data/hora a partir do carimbo de data/hora da mensagem fornecido pela origem.

partitionIdleTimeout

documento

Opcional

documento que especifica a quantidade de tempo que uma partição pode ficar ociosa antes de ser ignorada nos cálculos de marca d'agua.

Este campo é desabilitado por padrão. Para lidar com partições que não progridem devido à inatividade, atribua um valor a este campo.

partitionIdleTimeout.size

inteiro

Opcional

Número que especifica a duração do tempo limite de inatividade da partição.

partitionIdleTimeout.unit

string

Opcional

Unidade de tempo para a duração do tempo limite de inatividade da partição.

O valor de unit pode ser um dos seguintes:

  • "ms" (milésimo de segundo)

  • "second"

  • "minute"

  • "hour"

  • "day"

schemaRegistry

documento

Opcional

Documento que habilita o uso de um Registro de Esquema para suportar a leitura de uma fonte serializada Avro.

Para habilitar esse recurso, você deve criar uma conexão de registro de esquema.

schemaRegistry.connectionName

string

Condicional

Nome da conexão do Registro de Esquema a ser usada para desserialização Avro.

config

documento

Opcional

documento que contém campo que substituem vários valores padrão.

config.auto_offset_reset

string

Opcional

Especifica com qual evento no tópico de origem do Apache Kafka iniciar a ingestão. auto_offset_reset usa os seguintes valores:

  • end, latest ou largest : para iniciar a ingestão a partir do evento mais recente no tópico no momento em que a agregação é inicializada.

  • earliest, beginning ou smallest : para iniciar a ingestão a partir do evento mais antigo no tópico.

Padrão é latest.

config.group_id

string

Opcional

ID do grupo de consumidores Kafka a ser associado ao processador de fluxo. Se omitido, o Atlas Stream Processing associa o espaço de trabalho do processamento de fluxos a um ID gerado automaticamente no seguinte formato:

asp-${streamProcessorId}-consumer

O Atlas Stream Processing gera automaticamente um valor para este parâmetro para todos os processadores de stream persistentes. Para processadores de fluxo efêmero definidos com start.process(), esse parâmetro será definido somente se você definí-lo manualmente.

config.enable_auto_commit

booleano

Condicional

Sinalizador que determina a política de confirmação para compensações de partição do Kafka Corretor. O Atlas Stream Processing suporta duas políticas de confirmação:

  • Se você definir esse parâmetro como true, o Atlas Stream Processing realiza a confirmação dos deslocamentos toda vez que o estágio $source passar os dados para o próximo operador.

  • Se você definir esse parâmetro como false, os processadores de stream farão a confirmação dos deslocamentos de partição quando o Atlas Stream Processing fizer um checkpoint.

Para um processador de fluxo efêmero definido com sp.process(), esse parâmetro é padrão como false, a menos que você defina group_id. Caso contrário, o padrão é true.

Para obter mais informações sobre deslocamentos ao usar o Kafka como $source, consulte Fontes de Kafka e Deslocamentos de grupo de consumidores.

config.keyFormat

string

Opcional

Tipo de dados usado para desserializar dados-chave do Apache Kafka. Deve ser um dos seguintes valores:

  • "binData"

  • "string"

  • "json"

  • "int"

  • "long"

Padrão é binData.

config.keyFormatError

string

Opcional

Como lidar com erros encontrados ao desserializar os dados-chave do Apache Kafka. Deve ser um dos seguintes valores:

Observação

O Atlas Stream Processing requer que os documentos no fluxo de dados de origem sejam json ou ejson válidos. O Atlas Stream Processing define os documentos que não atendem a esse requisito na sua fila de mensagens não entregues, se você tiver configurado uma.

Um fluxo de alterações de coleção do Atlas permite que aplicativos acessem mudanças de dados em tempo real em uma única coleção. Para aprender como abrir um fluxo de alterações em uma coleção, veja Change Streams.

Ao usar um fluxo de alterações $source, configure o cluster de origem com uma oplog window de pelo menos 24 horas.

Para ler o fluxo de alterações, o Atlas Stream Processing verifica a coleção de oplog. Como resultado, você poderá ver COLLSCAN avisos em seus logs. Esses avisos indicam um comportamento normal e não sinalizam um erro.

Se você configurar config.fullDocument ou config.fullDocumentBeforeChange para required, habilite changeStreamPreAndPostImages em cada coleção antes de qualquer operação de gravação que deseja capturar. Se uma pós-imagem estiver indisponível para um evento porque você não habilitou o recurso quando a gravação ocorreu ou porque a pós-imagem expirou, o processador de fluxo falhará. Para saber como ativar pré e pós-imagens, consulte Fluxo de alterações com pré e pós-imagens de documentos.

Para operar na transmissão de dados de um change stream do Atlas, o estágio $source tem o seguinte formato de protótipo:

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"db" : "<source-db>",
"coll" : ["<source-coll>",...],
"initialSync": {
"enable": <boolean>,
"parallelism": <integer>
},
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}],
"maxAwaitTimeMS": <time-ms>,
}
}
}

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

Campo
Tipo
necessidade
Descrição

connectionName

string

Condicional

Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados.

timeField

documento

Opcional

documento que define um carimbo de data/hora oficial para mensagens recebidas.

Se você usar timeField, deverá defini-lo como um dos seguintes:

  • uma expressão $toDate que usa um campo de mensagem de origem como argumento

  • uma expressão $dateFromString que usa um campo de mensagem de origem como argumento.

Se você não declarar um timeField, o Atlas Stream Processing criará um carimbo de data/hora a partir do carimbo de data/hora da mensagem fornecido pela origem.

db

string

Obrigatório

Nome de um banco de dados MongoDB hospedado na instância do Atlas especificado por connectionName. O change stream desse banco de dados atua como a fonte de dados de streaming.

coll

cadeia de caracteres ou matriz de cadeias de caracteres

Obrigatório

Nome de uma ou mais collections MongoDB hospedadas na instância do Atlas especificada por connectionName. O change stream dessas collections atua como a fonte de dados de streaming . Se você omitir este campo, seu processador de stream fonterá a partir de um Change Stream do Banco de Dados MongoDB .

initialSync

documento

Opcional

Documento contendo parâmetros de configuração initialSync.

O Atlas Stream Processing initialSync permite a você ingestão de documentos preexistentes em uma collection do Atlas como se fossem documentos de inserção changeEvent. Se você habilitar initialSync, ao iniciar o processador de fluxo, ele primeiro ingeri e processa todos os documentos existentes na coleção antes de prosseguir para a ingestão e processamento dos novos documentos changeEvent recebidos. Depois que o initialSync estiver concluído, ele não se repetirá.

Se coll nomear mais de uma collection, o Atlas Stream Processing sincroniza todas as collections na lista. Para saber mais sobre a ordem de sincronização de várias collections, limites de nível e exemplos, consulte Configurar sincronização inicial de várias collections.

O Atlas Stream Processing corrige o conjunto de coleções a serem sincronizadas quando o processador de streams é iniciado. Ele não sincroniza as coleções que você cria depois desse ponto e, se a sincronização de qualquer coleção falhar, todo o processador de fluxo falhará.

Se você habilitar initialSync, não poderá usar os estágios $hoppingWindow, $sessionWindow ou $tumblingWindow em seu pipeline.

IMPORTANTE: consulte Limitações para considerações e restrições initialSync.

initialSync.enable

booleano

Condicional

Determina se initialSync deve ser ativado. Se você declarar um campo initialSync, deverá definir este campo.

initialSync.parallelism

inteiro

Opcional

Determina o nível de paralelismo com o qual processar a operação initialSync. Se você não especificar um valor, o padrão será 1.

Se coll nomear mais de uma collection, esse valor se aplicará à lista de collection como um todo e não a cada collection.

initialSync aplica esse valor a collections com qualquer tipo _id compatível, incluindo tipos que o Atlas Stream Processing não consegue ler em uma ordem previsível.

O valor máximo que você pode definir depende do tier do seu processador de stream. Para saber mais, consulte o Guia de seleção do tier do Atlas Stream Processing.

Cada processador de fluxo tem um valor máximo de paralelismo cumulativo determinado por seu nível. O paralelismo cumulativo de um processador de fluxo é calculado da seguinte forma:

parallelism total - parallelized stages

Em que parallelism total é a soma de todos os parallelism valores maiores que 1 nos $source estágios, $lookup, $merge, $emit e $externalFunction, e parallelized stages é o número desses estágios com parallelism valores superiores 1 a.

Por exemplo, se o seu estágio $source define um valor parallelism de 4, o seu estágio $lookup não define nenhum valor parallelism (portanto, o padrão é 1) e o seu estágio $merge define um parallelism valor de 2, então você tem dois parallelized stages, e o paralelismo cumulativo do seu processador de fluxo é calculado como (4 + 2) - 2.

Se um processador de fluxo exceder o paralelismo cumulativo máximo para seu nível, o Atlas Stream Processing lançará um erro e o avisará sobre o nível mínimo de processador necessário para o nível pretendido de paralelismo. Você deve dimensionar o processador para um nível superior ou reduzir os valores de paralelismo de seus estágios para resolver o erro. Para saber mais, consulte Stream Processing.

readPreference

string

Opcional

Read preference para change stream e initialSync operações.

Padrão é primary.

readPreferenceTags

array

Opcional

Marcações de read preference para change stream e initialSync operações.

config

documento

Opcional

documento que contém campo que substituem vários valores padrão.

config.startAfter

token

Condicional

O evento de alteração após o qual a fonte começa a relatar. Isso assume a forma de um resume token.

Você pode usar apenas um entre config.startAfter ou config.startAtOperationTime.

config.startAtOperationTime

carimbo de data/hora | data

Condicional

O tempo de operação após o qual a fonte deve começar a relatar.

Você pode usar apenas um entre config.startAfter ou config.startAtOperationTime.

Aceita $date $timestamp valores do MongoDB Extended JSON ou.

config.fullDocument

string

Condicional

Configuração que controla se uma fonte de fluxo de alterações deve retornar um documento completo ou apenas as alterações quando ocorrer uma atualização. Deve ser um dos seguintes:

  • default : não retorna um documento completo para operações update.

  • updateLookup : retorna apenas as alterações na atualização.

  • required : Deve retornar um documento completo . Se um documento completo não estiver disponível, o processador de fluxo falhará. Isso não se aplica a operações de exclusão, que ainda criam eventos de exclusão se o documento completo não estiver disponível.

  • whenAvailable : retorna um documento completo sempre que houver um disponível, caso contrário, retorna alterações.

Para usar este campo com um fluxo de alterações de coleção, você deve habilitar o fluxo de alterações pré e pós-imagens nessa coleção.

config.fullDocumentOnly

booleano

Condicional

Configuração que controla se uma change stream retorna todo o documento do evento, incluindo todos os metadados, ou somente o conteúdo de fullDocument. Se definido como true, a origem retornará somente o conteúdo de fullDocument.

Para usar este campo com um fluxo de alterações de coleção, você deve habilitar o fluxo de alterações pré e pós-imagens nessa coleção.

config.fullDocumentBeforeChange

string

Opcional

Especifica se uma fonte de change stream deve incluir o documento completo em seu estado original "antes das alterações" na saída. Deve ser um dos seguintes:

  • off : omite o campo fullDocumentBeforeChange .

  • required : deve retornar um documento completo em seu estado anterior às alterações. Se um documento completo em seu estado anterior às alterações não estiver disponível, o processador de fluxo falhará.

  • whenAvailable : retorna um documento completo em seu estado anterior às alterações sempre que um estiver disponível, caso contrário, omite o campo fullDocumentBeforeChange .

Se você não especificar um valor para fullDocumentBeforeChange, o padrão será off.

Para usar este campo com um fluxo de alterações de coleção, você deve habilitar o fluxo de alterações pré e pós-imagens nessa coleção.

config.pipeline

documento

Opcional

Especifica um pipeline de agregação para filtrar a saída do fluxo de alterações antes de passá-la para processamento adicional. Este pipeline deve estar em conformidade com os parâmetros descritos em Modificar Saída do Fluxo de Alterações.

IMPORTANTE: cada Alterar evento inclui os campos wallTime e clusterTime. Os estágios de Atlas Stream Processing após $source esperam receber esses campos conforme o processador os ingeriu. Para garantir o em processamento adequado dos dados de Change Stream, não modifique esses campos em $source.config.pipeline.

config.maxAwaitTimeMS

inteiro

Opcional

Tempo máximo, em milissegundos, para aguardar que novas alterações de dados sejam relatadas ao cursor do fluxo de alterações antes de retornar um lote vazio.

Padrão é 1000.

Um fluxo de alterações do banco de dados Atlas permite que aplicativos acessem mudanças de dados em tempo real em um único banco de dados. Para saber como abrir um fluxo de alteração em um banco de dados, consulte Fluxos de alteração.

Ao usar um fluxo de alterações $source, configure o cluster de origem com uma oplog window de pelo menos 24 horas.

Para ler o fluxo de alterações, o Atlas Stream Processing verifica a coleção de oplog. Como resultado, você poderá ver COLLSCAN avisos em seus logs. Esses avisos indicam um comportamento normal e não sinalizam um erro.

Se você configurar config.fullDocument ou config.fullDocumentBeforeChange para required, habilite changeStreamPreAndPostImages em cada coleção antes de qualquer operação de gravação que deseja capturar. Se uma pós-imagem estiver indisponível para um evento porque você não habilitou o recurso quando a gravação ocorreu ou porque a pós-imagem expirou, o processador de fluxo falhará. Para saber como ativar pré e pós-imagens, consulte Fluxo de alterações com pré e pós-imagens de documentos.

Para operar na transmissão de dados de um fluxo de alteração de banco de dados do Atlas, o estágio $source tem o seguinte formato de protótipo:

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"db" : "<source-db>",
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}]
},
}
}

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

Campo
Tipo
necessidade
Descrição

connectionName

string

Condicional

Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados.

timeField

documento

Opcional

documento que define um carimbo de data/hora oficial para mensagens recebidas.

Se você usar timeField, deverá defini-lo como um dos seguintes:

  • uma expressão $toDate que usa um campo de mensagem de origem como argumento

  • uma expressão $dateFromString que usa um campo de mensagem de origem como argumento.

Se você não declarar um timeField, o Atlas Stream Processing criará um carimbo de data/hora a partir do carimbo de data/hora da mensagem fornecido pela origem.

db

string

Obrigatório

Nome de um banco de dados MongoDB hospedado na instância do Atlas especificado por connectionName. O change stream desse banco de dados atua como a fonte de dados de streaming.

readPreference

string

Opcional

Read preference para operações de change stream.

Padrão é primary.

readPreferenceTags

array

Opcional

Tags de read preference para operações de change stream.

config

documento

Opcional

documento que contém campo que substituem vários valores padrão.

config.startAfter

token

Condicional

O evento de alteração após o qual a fonte começa a relatar. Isso assume a forma de um resume token.

Você pode usar apenas um entre config.startAfter ou config.startAtOperationTime.

config.startAtOperationTime

carimbo de data/hora | data

Condicional

O tempo de operação após o qual a fonte deve começar a relatar.

Você pode usar apenas um entre config.startAfter ou config.startAtOperationTime.

Aceita $date $timestamp valores do MongoDB Extended JSON ou.

config.fullDocument

string

Condicional

Configuração que controla se uma fonte de fluxo de alterações deve retornar um documento completo ou apenas as alterações quando ocorrer uma atualização. Deve ser um dos seguintes:

  • default : usa o comportamento padrão do servidor. Não retorna um documento completo para operações update.

  • updateLookup : retorna apenas as alterações na atualização.

  • required : Deve retornar um documento completo . Se um documento completo não estiver disponível, o processador de fluxo falhará. Isso não se aplica a operações de exclusão, que ainda criam eventos de exclusão se o documento completo não estiver disponível.

  • whenAvailable : retorna um documento completo sempre que houver um disponível, caso contrário, retorna alterações.

Se você não especificar um valor para fullDocument, o padrão será default.

Para usar esse campo com um fluxo de alteração de banco de dados, você deve ativar a alteração do fluxo de pré e pós-imagens em cada coleção nesse banco de dados.

config.fullDocumentOnly

booleano

Condicional

Configuração que controla se uma change stream retorna todo o documento do evento, incluindo todos os metadados, ou somente o conteúdo de fullDocument. Se definido como true, a origem retornará somente o conteúdo de fullDocument.

Para usar esse campo com um fluxo de alteração de banco de dados, você deve ativar a alteração do fluxo de pré e pós-imagens em cada coleção nesse banco de dados.

config.fullDocumentBeforeChange

string

Opcional

Especifica se uma fonte de change stream deve incluir o documento completo em seu estado original "antes das alterações" na saída. Deve ser um dos seguintes:

  • off : omite o campo fullDocumentBeforeChange .

  • required : deve retornar um documento completo em seu estado anterior às alterações. Se um documento completo em seu estado anterior às alterações não estiver disponível, o processador de fluxo falhará.

  • whenAvailable : retorna um documento completo em seu estado anterior às alterações sempre que um estiver disponível, caso contrário, omite o campo fullDocumentBeforeChange .

Se você não especificar um valor para fullDocumentBeforeChange, o padrão será off.

Para usar esse campo com um fluxo de alteração de banco de dados, você deve ativar a alteração do fluxo de pré e pós-imagens em cada coleção nesse banco de dados.

config.pipeline

documento

Opcional

Especifica um pipeline de agregação para filtrar a saída do fluxo de alterações no ponto de origem. Este pipeline deve estar em conformidade com os parâmetros descritos em Modificar Saída do Fluxo de Alterações.

IMPORTANTE: cada Alterar evento inclui os campos wallTime e clusterTime. Os estágios de Atlas Stream Processing após $source esperam receber esses campos conforme o processador os ingeriu. Para garantir o em processamento adequado dos dados de Change Stream, não modifique esses campos em $source.config.pipeline.

config.maxAwaitTimeMS

inteiro

Opcional

Tempo máximo, em milissegundos, para aguardar que novas alterações de dados sejam relatadas ao cursor do fluxo de alterações antes de retornar um lote vazio.

Padrão é 1000.

Ao usar um fluxo de alterações $source, configure o cluster de origem com uma oplog window de pelo menos 24 horas.

Para ler o fluxo de alterações, o Atlas Stream Processing verifica a coleção de oplog. Como resultado, você poderá ver COLLSCAN avisos em seus logs. Esses avisos indicam um comportamento normal e não sinalizam um erro.

Se você configurar config.fullDocument ou config.fullDocumentBeforeChange para required, habilite changeStreamPreAndPostImages em cada coleção antes de qualquer operação de gravação que deseja capturar. Se uma pós-imagem estiver indisponível para um evento porque você não habilitou o recurso quando a gravação ocorreu ou porque a pós-imagem expirou, o processador de fluxo falhará. Para saber como ativar pré e pós-imagens, consulte Fluxo de alterações com pré e pós-imagens de documentos.

Para operar em dados de streaming de um fluxo de mudança de cluster inteiro do Atlas, o estágio $source tem a seguinte forma de protótipo:

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}]
},
}
}

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

Campo
Tipo
necessidade
Descrição

connectionName

string

Condicional

Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados.

timeField

documento

Opcional

documento que define um carimbo de data/hora oficial para mensagens recebidas.

Se você usar timeField, deverá defini-lo como um dos seguintes:

  • uma expressão $toDate que usa um campo de mensagem de origem como argumento

  • uma expressão $dateFromString que usa um campo de mensagem de origem como argumento.

Se você não declarar um timeField, o Atlas Stream Processing criará um carimbo de data/hora a partir do carimbo de data/hora da mensagem fornecido pela origem.

readPreference

string

Opcional

Read preference para operações de change stream.

Padrão é primary.

readPreferenceTags

array

Opcional

Tags de read preference para operações de change stream.

config

documento

Opcional

documento que contém campo que substituem vários valores padrão.

config.startAfter

token

Condicional

O evento de alteração após o qual a fonte começa a relatar. Isso assume a forma de um resume token.

Você pode usar apenas um entre config.startAfter ou config.startAtOperationTime.

config.startAtOperationTime

data | registro de data e hora

Condicional

O tempo de operação após o qual a fonte deve começar a relatar.

Você pode usar apenas um entre config.startAfter ou config.startAtOperationTime.

Aceita $date $timestamp valores do MongoDB Extended JSON ou.

config.fullDocument

string

Condicional

Configuração que controla se uma fonte de fluxo de alterações deve retornar um documento completo ou apenas as alterações quando ocorrer uma atualização. Deve ser um dos seguintes:

  • default : usa o comportamento padrão do servidor. Não retorna um documento completo para operações update.

  • updateLookup : retorna apenas as alterações na atualização.

  • required : Deve retornar um documento completo . Se um documento completo não estiver disponível, o processador de fluxo falhará. Isso não se aplica a operações de exclusão, que ainda criam eventos de exclusão se o documento completo não estiver disponível.

  • whenAvailable : retorna um documento completo sempre que houver um disponível, caso contrário, retorna alterações.

Se você não especificar um valor para fullDocument, o padrão será default.

Para usar esse campo com um fluxo de alteração de cluster, você deve ativar a pré e pós-imagens do fluxo de alterações em cada coleção nesse cluster.

config.fullDocumentOnly

booleano

Condicional

Configuração que controla se uma change stream retorna todo o documento do evento, incluindo todos os metadados, ou somente o conteúdo de fullDocument. Se definido como true, a origem retornará somente o conteúdo de fullDocument.

Para usar esse campo com um fluxo de alteração de banco de dados, você deve ativar a alteração do fluxo de pré e pós-imagens em cada coleção nesse banco de dados.

config.fullDocumentBeforeChange

string

Opcional

Especifica se uma fonte de change stream deve incluir o documento completo em seu estado original "antes das alterações" na saída. Deve ser um dos seguintes:

  • off : omite o campo fullDocumentBeforeChange .

  • required : deve retornar um documento completo em seu estado anterior às alterações. Se um documento completo em seu estado anterior às alterações não estiver disponível, o processador de fluxo falhará.

  • whenAvailable : retorna um documento completo em seu estado anterior às alterações sempre que um estiver disponível, caso contrário, omite o campo fullDocumentBeforeChange .

Se você não especificar um valor para fullDocumentBeforeChange, o padrão será off.

Para usar esse campo com um fluxo de alteração de banco de dados, você deve ativar a alteração do fluxo de pré e pós-imagens em cada coleção nesse banco de dados.

config.pipeline

documento

Opcional

Especifica um pipeline de agregação para filtrar a saída do fluxo de alterações no ponto de origem. Este pipeline deve estar em conformidade com os parâmetros descritos em Modificar Saída do Fluxo de Alterações.

Observe que o Atlas Stream Processing espera receber os campos wallTime e clusterTime de cada evento de alteração ingerido. Para garantir o processamento adequado dos dados do Change Stream, não modifique esses campos no $source.config.pipeline.

config.maxAwaitTimeMS

inteiro

Opcional

Tempo máximo, em milissegundos, para aguardar que novas alterações de dados sejam relatadas ao cursor do fluxo de alterações antes de retornar um lote vazio.

Padrão é 1000.

O Atlas Stream Processing permite criar conexões de link privado para streams do AWS Kinesis. Para aprender mais, consulte Adicionar uma conexão de link privado do Kinesis.

Para operar em dados de um fluxo de dados do AWS Kinesis, o estágio $source possui o seguinte formato de protótipo:

{
"$source": {
"connectionName": "<registered-connection>",
"stream": "<stream-name>",
"region": "<aws-region>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"tsFieldName": "<field-name>",
"shardIdleTimeout": {
"size": <duration-number>,
"unit": "<duration-unit>"
},
"config": {
"consumerARN": "<aws-arn>",
"initialPosition": <initial-position>,
reshardDetectionIntervalSecs: <interval>
}
}
}

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

Campo
Tipo
necessidade
Descrição

connectionName

string

Obrigatório

Etiqueta que identifica a conexão no Registro de conexão a partir da qual consumir dados.

config.consumerARN

string

Obrigatório

ARN correspondente a um consumidor Kinesis. Seu consumidor deve usar distribuição aprimorada.

stream

string

Obrigatório

Fluxo de dados do Kinesis do qual transmitir mensagens.

region

string

Condicional

Região da AWS na qual o fluxo especificado existe. O Kinesis oferece suporte a vários fluxos de dados com o mesmo nome em diferentes regiões. Se você usar o mesmo nome para streams de dados em duas ou mais regiões dentro da mesma conexão, deverá usar esse campo para especificar qual combinação de nome e região usar.

timeField

documento

Opcional

documento que define um carimbo de data/hora oficial para mensagens recebidas.

Se você usar timeField, deverá defini-lo como um dos seguintes:

  • uma expressão $toDate que usa um campo de mensagem de origem como argumento

  • uma expressão $dateFromString que usa um campo de mensagem de origem como argumento.

Se você não declarar um timeField, o Atlas Stream Processing criará um carimbo de data/hora a partir do carimbo de data/hora da mensagem fornecido pela origem.

tsFieldName

string

Opcional

Nome do campo para o carimbo de data/hora no documento projetado. Use este campo para substituir o nome do campo de carimbo de data/hora padrão.

shardIdleTimeout

documento

Opcional

Documento que especifica o tempo de espera que um fragmento pode ficar ocioso antes de ser ignorado nos cálculos de marca d'água.

Este campo está desabilitado por padrão. Para lidar com fragmentos que não avançam devido à ociosidade, defina um valor para esse campo.

shardIdleTimeout.size

documento

Opcional

Número que especifica a duração do tempo limite ocioso do fragmento.

shardIdleTimeout.unit

documento

Opcional

Unidade de tempo durante o tempo limite ocioso do fragmento.

O valor de unit pode ser um dos seguintes:

  • "ms" (milésimo de segundo)

  • "second"

  • "minute"

  • "hour"

  • "day"

config

documento

Opcional

documento que contém campo que substituem vários valores padrão.

config.initialPosition

string

Opcional

Posição no histórico do fluxo de dados do Kinesis a partir do qual iniciar a ingestão de mensagens. Deve ser um dos seguintes:

  • "TRIM_HORIZON": comece a ingestão a partir da mensagem mais antiga do fragmento.

  • "LATEST": Comece a ingestão a partir da mensagem mais recente no fragmento.

  • "AT_TIMESTAMP": comece a ingestão a partir de um carimbo de data/hora específico. Requer config.atTimestamp.

  • "AT_SEQUENCE_NUMBER": Comece a ingerir em um número de sequência específico.

  • "AFTER_SEQUENCE_NUMBER": Comece a ingerir após um número de sequência específico.

Padrão é "LATEST".

config.atTimestamp

data

Condicional

O carimbo de data/hora a partir do qual começar a ingerir mensagens. Obrigatório quando config.initialPosition é "AT_TIMESTAMP".

config.reshardDetectionIntervalSecs

inteiro

Opcional

Intervalo, em segundos, entre as verificações da taxa de fluxo de dados pelo fluxo do Kinesis para fins de refragmentação.

tem como padrão 300 segundos.

Para operar em uma array de documentos, o estágio $source tem o seguinte formato de protótipo:

{
"$source": {
"timeField": {
$toDate | $dateFromString: <expression>
},
"documents" : [{source-doc},...] | <expression>
}
}

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

Campo
Tipo
necessidade
Descrição

timeField

documento

Opcional

documento que define um carimbo de data/hora oficial para mensagens recebidas.

Se você usar timeField, deverá defini-lo como um dos seguintes:

  • uma expressão $toDate que usa um campo de mensagem de origem como argumento

  • uma expressão $dateFromString que usa um campo de mensagem de origem como argumento.

Se você não declarar um timeField, o Atlas Stream Processing criará um carimbo de data/hora a partir do carimbo de data/hora da mensagem fornecido pela origem.

documents

array

Condicional

Array de documentos para usar como fonte de dados de streaming. O valor deste campo pode ser uma matriz de objetos ou uma expressão que avalia para uma matriz de objetos. Não utilize este campo ao utilizar o campo connectionName .

Para gerar um documento em uma programação recorrente em vez de ler a partir de uma conexão, o estágio $source tem a seguinte forma de protótipo:

{
"$source": {
"schedule": "<cron-expression>",
"tsFieldName": "<timestamp-field-name>"
}
}

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

Campo
Tipo
necessidade
Descrição

schedule

string

Obrigatório

expressão cron de seis campos que determina quando o estágio emite um documento. Os campos são, em ordem: segundo, minuto, hora, dia do mês, mês e dia da semana. O Atlas Stream Processing avalia a expressão em UTC.

Cada campo aceita * para todos os valores, um único valor, um intervalo como 1-5, uma etapa como 0/15 ou uma lista separada por vírgula como 1,15,30. Você pode usar os nomes JAN a DEC para meses e SUN a SAT para dias da semana. Não inclua um prefixo CRON_TZ.

Para saber como o Atlas Stream Processing resolve uma expressão que restringe o dia do mês e o dia da semana, consulte Comportamento.

tsFieldName

string

Opcional

Nome do campo no qual o estágio projeta o carimbo de data/hora programado. O padrão é _ts.

$source deve ser o primeiro estágio de qualquer pipeline em que aparece. Você pode usar apenas um estágio por $source pipeline.

Para os estágios Kafka $source, o Atlas Stream Processing lê em paralelo a partir de várias partições dentro do tópico de origem. O limite de partição é determinado pelo nível do processador. Para saber mais, veja a referência de faturamento do Stream Processing.

Para estágios cron $source, cada horário agendado produz um único documento vazio que contém o carimbo de data/hora agendado. Use estágios posteriores do pipeline para preencher o documento.

Quando uma expressão schedule restringe o dia do mês e o dia da semana, o Atlas Stream Processing emite um documento somente nos horários agendados que satisfazem ambos os campos. Isso difere das implementações do cron que emitem um documento quando um dos campo corresponde.

Se um processador de fluxo com um estágio cron $source for reiniciado, ele emitirá um documento para cada horário agendado que foi perdido enquanto não estava em execução. Ele não emite documentos para horários programados que foram perdidos depois que você o interrompeu.

Um processador de stream com um estágio cron $source é executado continuamente entre os horários agendados, e o Atlas cobra no seu nível de processador durante todo esse período. Um agendamento pouco frequente não reduz o custo de execução do processador. Para saber mais, consulte a referência de faturamento do Stream Processing.

Uma fonte de dados de streaming gera relatórios meteorológicos detalhados a partir de vários locais, em conformidade com o esquema do Conjunto de dados meteorológicos de amostra. A seguinte agregação tem três estágios:

  1. O estágio estabelece uma $source conexão com o intermediário do Apache Kafka, coletando esses relatórios em um my_weatherdata tópico chamado, expondo cada registro à medida que ele é ingerido para os estágios de agregação subsequentes. Esse estágio também substitui o nome do campo de carimbo de data/hora que ele projeta, definindo-o ingestionTime como.

  2. O estágio exclui documentos que tenham $match um dewPoint.value menor ou igual a 5.0 e passa os documentos com dewPoint.value maior que 5.0 para o próximo estágio.

  3. O estágio escreve o resultado em uma collection do Atlas $merge chamada stream no sample_weatherstream banco de dados. Se não existir tal banco de dados ou coleção, o Atlas os criará.

[{
"$source": {
"connectionName": "sample_weatherdata",
"topic": "my_weatherdata"
}
},
{
"$match": { "dewPoint.value": { "$gt": 5 } }
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "stream"
}
}
}]

Para visualizar os documentos na coleção sample_weatherstream.stream resultante, conecte-se ao cluster Atlas e execute o seguinte comando:

db.getSiblingDB("sample_weatherstream").stream.find()
{
_id: ObjectId('66ad2edfd4fcac13b1a28ce3'),
airTemperature: { quality: '1', value: 27.7 },
atmosphericPressureChange: {
quantity24Hours: { quality: '9', value: 99.9 },
quantity3Hours: { quality: '1' },
tendency: { code: '1', quality: '1' }
},
atmosphericPressureObservation: {
altimeterSetting: { quality: '1', value: 1015.9 },
stationPressure: { quality: '1', value: 1021.9 }
},
callLetters: 'CGDS',
dataSource: '4',
dewPoint: { quality: '9', value: 25.7 },
elevation: 9999,
extremeAirTemperature: {
code: 'N',
period: 99.9,
quantity: '9',
value: -30.4
},
ingestionTime: ISODate('2024-08-02T19:09:18.071Z'),
liquidPrecipitation: {
condition: '9',
depth: 160,
period: 24,
quality: '2'
},
pastWeatherObservationManual: {
atmosphericCondition: { quality: '1', value: '8' },
period: { quality: '9', value: 3 }
},
position: { coordinates: [ 153.3, 50.7 ], type: 'Point' },
precipitationEstimatedObservation: {
discrepancy: '4',
estimatedWaterDepth: 4
},
presentWeatherObservationManual: { condition: '53', quality: '1' },
pressure: { quality: '1', value: 1016.3 },
qualityControlProcess: 'V020',
seaSurfaceTemperature: { quality: '9', value: 27.6 },
sections: [ 'AA2', 'SA1', 'MW1', 'AG1', 'GF1' ],
skyCondition: {
cavok: 'N',
ceilingHeight: { determination: 'C', quality: '1', value: 6900 }
},
skyConditionObservation: {
highCloudGenus: { quality: '1', value: '05' },
lowCloudGenus: { quality: '9', value: '03' },
lowestCloudBaseHeight: { quality: '9', value: 150 },
lowestCloudCoverage: { quality: '1', value: '05' },
midCloudGenus: { quality: '9', value: '08' },
totalCoverage: { opaque: '99', quality: '1', value: '06' }
},
skyCoverLayer: {
baseHeight: { quality: '9', value: 99999 },
cloudType: { quality: '9', value: '05' },
coverage: { quality: '1', value: '04' }
},
st: 'x+35700-027900',
type: 'SAO',
visibility: {
distance: { quality: '1', value: 4000 },
variability: { quality: '1', value: 'N' }
},
waveMeasurement: {
method: 'I',
seaState: { code: '99', quality: '9' },
waves: { height: 99.9, period: 14, quality: '9' }
},
wind: {
direction: { angle: 280, quality: '9' },
speed: { quality: '1', rate: 30.3 },
type: '9'
}
}

Observação

O exemplo anterior é representativo. Os dados de streaming não são estáticos e cada usuário vê documentos distintos.

A seguinte agregação ingere dados da fonte cluster0-collection, que se conecta a um cluster do Atlas carregado com o conjunto de dados de amostra. Para aprender a criar um espaço de trabalho do processamento de fluxos e adicionar uma conexão a um cluster do Atlas ao registro de conexão, veja Introdução ao Atlas Stream Processing. Esta agregação executa duas etapas para abrir um fluxo de alterações e registrar as alterações na coleção data no banco de dados sample_weatherdata:

  1. O $source estágio se conecta à cluster0-collection origem e abre um fluxo de alteração na data coleção no sample_weatherdata banco de dados .

  2. O estágio grava os documentos do change stream filtrados em uma collection do Atlas $merge chamada data_changes no sample_weatherdata banco de dados. Se essa coleção não existir, o Atlas a criará.

[{
"$source": {
"connectionName": "cluster0-connection",
"db": "sample_weatherdata",
"coll": "data"
}
},
{
"$merge": {
"into": {
"connectionName": "cluster0-connection",
"db": "sample_weatherdata",
"coll": "data_changes"
}
}
}]

O seguinte comando exclui mongosh um data documento :

db.getSiblingDB("sample_weatherdata").data.deleteOne(
{ _id: ObjectId("5553a99ae4b02cf715120e4b") }
)

Após a data exclusão do documento, o processador de fluxo grava o documento de evento de fluxo de alteração na sample_weatherdata.data_changes coleção. Para visualizar os documentos na sample_weatherdata.data_changes coleção resultante, utilize mongosh o para conectar ao seu cluster do Atlas e execute o seguinte comando:

db.getSiblingDB("sample_weatherdata").data_changes.find()
[
{
_id: {
_data: '8267A3D7A3000000012B042C0100296E5A1004800951B8EDE4430AB5C1B254BB3C96D6463C6F7065726174696F6E54797065003C64656C6574650046646F63756D656E744B65790046645F696400645553A99AE4B02CF715120E4B000004'
},
clusterTime: Timestamp({ t: 1738790819, i: 1 }),
documentKey: { _id: ObjectId('5553a99ae4b02cf715120e4b') },
ns: { db: 'sample_weatherdata', coll: 'data' },
operationType: 'delete',
wallTime: ISODate('2025-02-05T21:26:59.313Z')
}
]

A seguinte agregação ingere dados da db-change-stream-connection fonte, que se conecta a um Atlas cluster carregado com o conjunto de dados Sample Mflix Dataset. Para saber como criar um espaço de trabalho de processamento de fluxo e adicionar uma conexão a um cluster Atlas ao registro de conexão, consulte Introdução ao Atlas Stream Processing. Essa agregação executa duas etapas para abrir um change stream no sample_mflix banco de dados de origem e registrar as alterações na db_changes collection no sample_mflix_changes banco de dados coletor:

  1. O $source estágio se conecta à db-change-stream-connection origem e abre um fluxo de alterações no sample_mflix banco de dados de origem. O config.startAtOperationTime campo define o tempo após o qual a fonte começa a reportar. Este exemplo define o startAt valor para iniciar um minuto atrás.

  2. O estágio grava os documentos do change stream em uma collection do Atlas $merge chamada db_changes no sample_mflix_changes banco de dados do coletor.

const startAt = new Date(Date.now() - 60 * 1000);
const pipeline = [
{
$source: {
connectionName: "db-change-stream-connection",
db: "sample_mflix",
config: {
startAtOperationTime: startAt
}
}
},
{
$merge: {
into: {
connectionName: "db-change-stream-connection",
db: "sample_mflix_changes",
coll: "db_changes"
}
}
}
];

Execute os seguintes comandos em seu cluster de origem para verificar o comportamento do processador de stream. Para se conectar, consulte Conectar-se a um cluster via mongosh.

Insira um documento na coleção movies e um documento na coleção comments no banco de dados de origem do sample_mflix :

db.getSiblingDB("sample_mflix").movies.insertOne({
title: "The Stream Processor",
year: 2026
})
db.getSiblingDB("sample_mflix").comments.insertOne({
name: "Ada Lovelace",
text: "A fine film about data in motion."
})

Depois de inserir os documentos, o processador de fluxo grava um documento de evento de fluxo de alteração para cada inserção na coleção sample_mflix_changes.db_changes. Execute o seguinte comando em seu cluster para visualizar os documentos na coleção sample_mflix_changes.db_changes resultante:

db.getSiblingDB("sample_mflix_changes").db_changes.find(
{},
{ _id: 0, clusterTime: 1, ns: 1, operationType: 1, fullDocument: 1 }
)
[
{
clusterTime: Timestamp({ t: 1787778766, i: 1 }),
fullDocument: {
_id: ObjectId('6a8f56ce73c5d4d0ee569a0d'),
title: 'The Stream Processor',
year: 2026
},
ns: { db: 'sample_mflix', coll: 'movies' },
operationType: 'insert'
},
{
clusterTime: Timestamp({ t: 1787778771, i: 1 }),
fullDocument: {
_id: ObjectId('6a8f56d373c5d4d0ee569a0e'),
name: 'Ada Lovelace',
text: 'A fine film about data in motion.'
},
ns: { db: 'sample_mflix', coll: 'comments' },
operationType: 'insert'
}
]

Cada evento inclui um carimbo de data/hora, a alteração completa do documento , um operationType e um campo ns. O campo ns nomeia o banco de dados de origem e a coleção, para que você possa distinguir qual coleção produz cada alteração.

A agregação a seguir ingere dados da origem cluster-changestream-connection, que se conecta a um Atlas cluster carregado com o conjunto de dados Sample Mflix Dataset, e grava no destino cluster-changestream-sink-connection, que se conecta a um Atlas cluster separado. Para saber como criar um espaço de trabalho de processamento de fluxo e adicionar uma conexão a um cluster Atlas ao registro de conexão, consulte Introdução ao Atlas Stream Processing. Essa agregação executa dois estágios para abrir um change stream em todo o cluster no cluster de origem e registrar as alterações na collection events no banco de dados cluster_changes no cluster de destino. Esse estágio $source omite db e coll, portanto, ele relata alterações de cada banco de dados e coleção no cluster de origem, em vez de uma única coleção ou banco de dados.

  1. O estágio $source se conecta à origem cluster-changestream-connection e abre um fluxo de alterações em todo o cluster de origem. O campo config.startAtOperationTime especifica que o processador de fluxo deve começar a relatar as alterações que ocorrem em ou após um período de tempo especificado.

  2. O estágio $merge grava os documentos do change stream em uma collection do Atlas chamada events no banco de dados cluster_changes no cluster de destino. Se esse banco de dados ou collection não existir, o Atlas os criará.

[{
"$source": {
"connectionName": "cluster-changestream-connection",
"config": {
"startAtOperationTime": {"$date": "2024-08-19T18:00:00.000Z"}
}
}
},
{
"$merge": {
"into": {
"connectionName": "cluster-changestream-sink-connection",
"db": "cluster_changes",
"coll": "events"
}
}
}]

O comando mongosh a seguir insere um documento na collection movies no banco de dados sample_mflix no cluster de origem:

db.getSiblingDB("sample_mflix").movies.insertOne(
{ _id: ObjectId("66c1a1f1f1f1f1f1f1f1f1f1"), title: "Example Movie" }
)

Depois que o documento é inserido, o processador de stream grava o documento de evento de fluxo de alterações na coleção cluster_changes.events no cluster de destino.

Para visualizar os documentos na coleção cluster_changes.events resultante, execute os seguintes comandos no cluster de destino. Para se conectar, consulte Conectar-se a um cluster via mongosh.

db.getSiblingDB("cluster_changes").events.find()
[
{
_id: {
_data: '826A8F5E65000000012B042C0100296E5A10049CC3BB749CB641F4BF5A032A2ABE2FB1463C6F7065726174696F6E54797065003C696E736572740046646F63756D656E744B65790046645F6964006466C1A1F1F1F1F1F1F1F1F1F1000004'
},
clusterTime: Timestamp({ t: 1787780709, i: 1 }),
documentKey: { _id: ObjectId('66c1a1f1f1f1f1f1f1f1f1f1') },
fullDocument: { _id: ObjectId('66c1a1f1f1f1f1f1f1f1f1f1'), title: 'Example Movie' },
ns: { db: 'sample_mflix', coll: 'movies' },
operationType: 'insert',
wallTime: ISODate('2026-08-26T21:45:09.059Z')
}
]

O campo ns no documento de saída mostra que a alteração se origina da inserção sample_mflix.movies no cluster de origem. O agregação pipeline do fluxo de alterações reflete essa alteração na coleção cluster_changes.events no cluster de destino.

A seguinte agregação usa uma array de documento in-line como fonte de dados de streaming, contendo observações meteorológicas para três locais. A array usa o mesmo esquema que o conjunto de dados meteorológicos de amostra. Essa agregação executa três estágios:

  1. O estágio define $source uma documents array in-line de observações meteorológicas como fonte de dados de transmissão e usa timeField para designar o timestamp campo de cada documento como o carimbo de data/hora oficial.

  2. O estágio passa apenas os documentos $match com dewPoint.value maior que 5.0 para o próximo estágio.

  3. O estágio escreve o resultado em uma collection do Atlas $merge chamada stream no sample_weatherstream banco de dados. Se não existir tal banco de dados ou coleção, o Atlas os criará.

[{
"$source": {
"documents": [
{
"location": "New York",
"timestamp": ISODate('2024-01-15T08:00:00Z'),
"temp": 23.5,
"dewPoint": { "value": 6.2 }
},
{
"location": "Los Angeles",
"timestamp": ISODate('2024-01-15T08:05:00Z'),
"temp": 18.2,
"dewPoint": { "value": 4.8 }
},
{
"location": "Chicago",
"timestamp": ISODate('2024-01-15T08:10:00Z'),
"temp": 26.8,
"dewPoint": { "value": 7.5 }
}
]
}
},
{
"$match": { "dewPoint.value": { "$gt": 5.0 } }
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "stream"
}
}
}]

Para visualizar os documentos na coleção sample_weatherstream.stream resultante, conecte-se ao cluster Atlas e execute o seguinte comando:

db.getSiblingDB("sample_weatherstream").stream.find()
[
{
_id: ObjectId('67a6c3df14fcac13b1a21a01'),
dewPoint: { value: 6.2 },
location: 'New York',
temp: 23.5,
timestamp: ISODate('2024-01-15T08:00:00.000Z')
},
{
_id: ObjectId('67a6c3df14fcac13b1a21a03'),
dewPoint: { value: 7.5 },
location: 'Chicago',
temp: 26.8,
timestamp: ISODate('2024-01-15T08:10:00.000Z')
}
]

Os seguintes relatórios de filtros de agregação da fonte de dados de streaming de amostra sample_stream_solar. Ele arquiva os resultados na conexão solar-cluster-connection. Para saber como criar um espaço de trabalho de processamento de fluxo e adicionar uma conexão ao seu espaço de trabalho de Atlas Stream Processing , consulte Introdução ao Atlas Stream Processing e Adicionar uma conexão de Atlas Stream Processing . Essa agregação executa três estágios para filtrar relatórios da origem sample_stream_solar e escrever os resultados em uma collection chamada solarColl no banco de dadossolarDb:

  1. O estágio $source se conecta à fonte sample_stream_solar. O campo timeField converte o campo timestamp de cada relatório recebido em uma data usando $dateFromString.

  2. O estágio $match exclui documentos em que device_id é device_8 e passa relatórios de todos os outros dispositivos para o próximo estágio.

  3. O estágio $merge escreve o resultado na collection solarColl no banco de dados solarDb do cluster solar-cluster-connection.

[{
"$source": {
"connectionName": "sample_stream_solar",
"timeField": {
"$dateFromString": { "dateString": "$timestamp" }
}
}
},
{
"$match": { "device_id": { "$ne": "device_8" } }
},
{
"$merge": {
"into": {
"connectionName": "solar-cluster-connection",
"db": "solarDb",
"coll": "solarColl"
}
}
}]

Observação

A origem sample_stream_solar é uma conexão somente de teste que produz documentos de amostra a cada segundo e é destinada à criação rápida de protótipos.

Execute os seguintes comandos no seu cluster solar-cluster-connection para verificar o comportamento do processador de stream. Para se conectar, consulte Conectar-se a um cluster via mongosh.

Para visualizar os documentos na coleção solarDb.solarColl resultante, execute:

db.getSiblingDB("solarDb").solarColl.find()
[
{
_id: ObjectId('6a8f5f0352c6ffe64c8884ad'),
device_id: 'device_3',
event_type: 0,
group_id: 8,
max_watts: 450,
obs: { watts: 428, temp: 10 },
timestamp: '2024-08-12T21:47:47.287+00:00'
}
]

Cada documento inclui os campos device_id, event_type, group_id, max_watts, obs e timestamp de um relatório de dispositivo solar. Você pode observar que nenhum documento na saída tem um device_id de device_8, porque o estágio $match exclui esses relatórios antes que o estágio $merge grave os documentos restantes em solarColl.

Observação

O exemplo anterior é representativo. Os dados de streaming não são estáticos e cada usuário vê documentos distintos.

A seguinte agregação gera um documento a cada cinco minutos em vez de ler a partir de uma conexão. Essa agregação executa três estágios:

  1. O estágio $source emite um documento vazio no segundo 0 de cada quinto minuto, em UTC.

  2. O estágio $project rotula cada documento com um nome de tarefa e copia o carimbo de data/hora agendado do campo _ts para um camporunAt .

  3. O estágio escreve o resultado em uma collection do Atlas $merge chamada heartbeats no sample_weatherstream banco de dados. Se não existir tal banco de dados ou coleção, o Atlas os criará.

[{
"$source": {
"schedule": "0 0/5 * * * *"
}
},
{
"$project": {
"job": "five-minute-heartbeat",
"runAt": "$_ts"
}
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "heartbeats"
}
}
}]

Para visualizar os documentos na coleção sample_weatherstream.heartbeats resultante, conecte-se ao cluster Atlas e execute o seguinte comando:

db.getSiblingDB("sample_weatherstream").heartbeats.find()
[
{
_id: ObjectId('68b8a1c014fcac13b1a21b01'),
job: 'five-minute-heartbeat',
runAt: ISODate('2026-09-03T14:00:00.000Z')
},
{
_id: ObjectId('68b8a2f414fcac13b1a21b02'),
job: 'five-minute-heartbeat',
runAt: ISODate('2026-09-03T14:05:00.000Z')
},
{
_id: ObjectId('68b8a42814fcac13b1a21b03'),
job: 'five-minute-heartbeat',
runAt: ISODate('2026-09-03T14:10:00.000Z')
}
]

Observação

O anterior é um exemplo representativo . Os carimbos de data/hora que você vê dependem de quando você inicia o processador de stream.