Definição
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
Sintaxe
Corretora Apache Kafka
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 | |
|---|---|---|---|---|
| string | Obrigatório | Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados. | |
| 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. | |
| documento | Opcional | documento que define um carimbo de data/hora oficial para mensagens recebidas. Se você usar
Se você não declarar um | |
| 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. | |
| inteiro | Opcional | Número que especifica a duração do tempo limite de inatividade da partição. | |
| string | Opcional | Unidade de tempo para a duração do tempo limite de inatividade da partição. O valor de
| |
| 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. | |
| string | Condicional | Nome da conexão do Registro de Esquema a ser usada para desserialização Avro. | |
| documento | Opcional | documento que contém campo que substituem vários valores padrão. | |
| string | Opcional | Especifica com qual evento no tópico de origem do Apache Kafka iniciar a ingestão.
Padrão é | |
| 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: 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. | |
| 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:
Para um processador de fluxo efêmero definido com sp.process(), esse parâmetro é padrão como Para obter mais informações sobre deslocamentos ao usar o Kafka como | |
| string | Opcional | Tipo de dados usado para desserializar dados-chave do Apache Kafka. Deve ser um dos seguintes valores:
Padrão é | |
| 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.
Fluxo de alterações da coleção do MongoDB
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 |
|---|---|---|---|
| string | Condicional | Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados. |
| documento | Opcional | documento que define um carimbo de data/hora oficial para mensagens recebidas. Se você usar
Se você não declarar um |
| string | Obrigatório | Nome de um banco de dados MongoDB hospedado na instância do Atlas especificado por |
| 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 |
| documento | Opcional | Documento contendo parâmetros de configuração O Atlas Stream Processing Se 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 IMPORTANTE: consulte Limitações para considerações e restrições |
| booleano | Condicional | Determina se |
| inteiro | Opcional | Determina o nível de paralelismo com o qual processar a operação Se
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:
Em que Por exemplo, se o seu estágio 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. |
| string | Opcional | Read preference para change stream e Padrão é |
| array | Opcional | Marcações de read preference para change stream e |
| documento | Opcional | documento que contém campo que substituem vários valores padrão. |
| 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 |
| carimbo de data/hora | data | Condicional | |
| 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:
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. |
| 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 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. |
| 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:
Se você não especificar um valor para 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. |
| 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 |
| 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 é |
change stream do MongoDB database
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 |
|---|---|---|---|
| string | Condicional | Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados. |
| documento | Opcional | documento que define um carimbo de data/hora oficial para mensagens recebidas. Se você usar
Se você não declarar um |
| string | Obrigatório | Nome de um banco de dados MongoDB hospedado na instância do Atlas especificado por |
| string | Opcional | Read preference para operações de change stream. Padrão é |
| array | Opcional | Tags de read preference para operações de change stream. |
| documento | Opcional | documento que contém campo que substituem vários valores padrão. |
| 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 |
| carimbo de data/hora | data | Condicional | |
| 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:
Se você não especificar um valor para 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. |
| 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 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. |
| 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:
Se você não especificar um valor para 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. |
| 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 |
| 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 é |
Fonte de fluxo de alterações em todo o cluster do MongoDB
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 |
|---|---|---|---|
| string | Condicional | Etiqueta que identifica a conexão no Registro de conexão, para ingestão de dados. |
| documento | Opcional | documento que define um carimbo de data/hora oficial para mensagens recebidas. Se você usar
Se você não declarar um |
| string | Opcional | Read preference para operações de change stream. Padrão é |
| array | Opcional | Tags de read preference para operações de change stream. |
| documento | Opcional | documento que contém campo que substituem vários valores padrão. |
| 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 |
| data | registro de data e hora | Condicional | |
| 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:
Se você não especificar um valor para 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. |
| 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 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. |
| 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:
Se você não especificar um valor para 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. |
| 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 |
| 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 é |
AWS Kinesis Data Stream
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 |
|---|---|---|---|
| string | Obrigatório | Etiqueta que identifica a conexão no Registro de conexão a partir da qual consumir dados. |
| string | Obrigatório | ARN correspondente a um consumidor Kinesis. Seu consumidor deve usar distribuição aprimorada. |
| string | Obrigatório | Fluxo de dados do Kinesis do qual transmitir mensagens. |
| 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. |
| documento | Opcional | documento que define um carimbo de data/hora oficial para mensagens recebidas. Se você usar
Se você não declarar um |
| 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. |
| 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. |
| documento | Opcional | Número que especifica a duração do tempo limite ocioso do fragmento. |
| documento | Opcional | Unidade de tempo durante o tempo limite ocioso do fragmento. O valor de
|
| documento | Opcional | documento que contém campo que substituem vários valores padrão. |
| 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:
Padrão é |
| data | Condicional | O carimbo de data/hora a partir do qual começar a ingerir mensagens. Obrigatório quando |
| 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 |
Array do documento
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 |
|---|---|---|---|
| documento | Opcional | documento que define um carimbo de data/hora oficial para mensagens recebidas. Se você usar
Se você não declarar um |
| 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 |
Agendamento Cron
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 |
|---|---|---|---|
| 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 saber como o Atlas Stream Processing resolve uma expressão que restringe o dia do mês e o dia da semana, consulte Comportamento. |
| string | Opcional | Nome do campo no qual o estágio projeta o carimbo de data/hora programado. O padrão é |
Comportamento
$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.
Exemplos
Exemplo do Kafka
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:
O estágio estabelece uma
$sourceconexão com o intermediário do Apache Kafka, coletando esses relatórios em ummy_weatherdatató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-oingestionTimecomo.O estágio exclui documentos que tenham
$matchumdewPoint.valuemenor ou igual a5.0e passa os documentos comdewPoint.valuemaior que5.0para o próximo estágio.O estágio escreve o resultado em uma collection do Atlas
$mergechamadastreamnosample_weatherstreambanco 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.
Exemplo de fluxo de alteração
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:
O
$sourceestágio se conecta àcluster0-collectionorigem e abre um fluxo de alteração nadatacoleção nosample_weatherdatabanco de dados .O estágio grava os documentos do change stream filtrados em uma collection do Atlas
$mergechamadadata_changesnosample_weatherdatabanco 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') } ]
Exemplo de fluxo de alteração de banco de dados
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:
O
$sourceestágio se conecta àdb-change-stream-connectionorigem e abre um fluxo de alterações nosample_mflixbanco de dados de origem. Oconfig.startAtOperationTimecampo define o tempo após o qual a fonte começa a reportar. Este exemplo define ostartAtvalor para iniciar um minuto atrás.O estágio grava os documentos do change stream em uma collection do Atlas
$mergechamadadb_changesnosample_mflix_changesbanco 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.
Exemplo de fluxo de alterações em todo o cluster
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.
O estágio
$sourcese conecta à origemcluster-changestream-connectione abre um fluxo de alterações em todo o cluster de origem. O campoconfig.startAtOperationTimeespecifica 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.O estágio
$mergegrava os documentos do change stream em uma collection do Atlas chamadaeventsno banco de dadoscluster_changesno 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.
Exemplo de array de documentos
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:
O estágio define
$sourceumadocumentsarray in-line de observações meteorológicas como fonte de dados de transmissão e usatimeFieldpara designar otimestampcampo de cada documento como o carimbo de data/hora oficial.O estágio passa apenas os documentos
$matchcomdewPoint.valuemaior que5.0para o próximo estágio.O estágio escreve o resultado em uma collection do Atlas
$mergechamadastreamnosample_weatherstreambanco 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') } ]
Exemplo de Corrente Sol
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:
O estágio
$sourcese conecta à fontesample_stream_solar. O campotimeFieldconverte o campotimestampde cada relatório recebido em uma data usando$dateFromString.O estágio
$matchexclui documentos em quedevice_idédevice_8e passa relatórios de todos os outros dispositivos para o próximo estágio.O estágio
$mergeescreve o resultado na collectionsolarCollno banco de dadossolarDbdo clustersolar-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.
Exemplo de agendamento cron
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:
O estágio
$sourceemite um documento vazio no segundo0de cada quinto minuto, em UTC.O estágio
$projectrotula cada documento com um nome de tarefa e copia o carimbo de data/hora agendado do campo_tspara um camporunAt.O estágio escreve o resultado em uma collection do Atlas
$mergechamadaheartbeatsnosample_weatherstreambanco 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.