Definição
sp.process()Cria um processador de fluxo efêmero no espaço de trabalho de Stream Processingatual.
Você só pode invocar este comando enquanto estiver conectado a um espaço de trabalho de stream processing.
Este comando requer
mongoshversão ≥ 2.0.
Compatibilidade
Esse método é suportado em Atlas Stream Processing Workspaces.
Sintaxe
O método sp.process() tem a seguinte sintaxe:
sp.process( [ <pipeline> ], { <options> } )
Campos de comando
sp.process() usa estes campos:
Campo | Tipo | necessidade | Descrição |
|---|---|---|---|
| array | Obrigatório | Pipeline de agregação de streaming que o processador de streaming aplica aos seus dados de streaming. |
| objeto | Opcional | Objeto que define configurações opcionais para o seu processador de fluxo. |
| inteiro | Opcional | Número máximo de documentos a serem retornados ao seu terminal. |
| string | Opcional | Tier para atribuir ao processador de fluxo. Os nomes de tier não diferenciam maiúsculas de minúsculas. Se você omitir este campo, o processador de fluxo usará o tier padrão para seu Stream Processing Workspace. Deve ser um dos seguintes valores:
Para saber mais, consulte Níveis. |
Comportamento
sp.process() cria um processador de fluxo efêmero e sem nome no espaço de trabalho de Stream Processing atual e o inicializa imediatamente. Este processador de fluxo persiste apenas enquanto é executado. Se você encerrar um processador de fluxo efêmero, deverá criá-lo novamente para usá-lo.
Um processador de fluxo efêmero para quando qualquer um dos seguintes ocorre:
O processador retorna o número de documentos definidos por
options.limit.10 minutes elapse.
Você interrompe o processador.
Controle de acesso
O usuário que executa o sp.process() deve ter a função atlasAdmin.
Exemplo
O exemplo a seguir cria um processador de fluxo efêmero que ingere dados da conexão sample_stream_solar. O processador exclui todos os documentos onde o valor do campo device_id é device_8, passando o restante para uma janela de tombamento com uma duração de 10segundos. Cada janela agrupa os documentos que recebe, e então retorna várias estatísticas de cada grupo. O processador de fluxo então mescla esses registros para solar_db.solar_coll sobre a conexão mongodb1.
sp.process( [ { $source: { connectionName: 'sample_stream_solar', timeField: { $dateFromString: { dateString: '$timestamp' } } } }, { $match: { $expr: { $ne: [ "$device_id", "device_8" ] } } }, { $tumblingWindow: { interval: { size: Int32(10), unit: "second" }, "pipeline": [ { $group: { "_id": { "device_id": "$device_id" }, "max_temp": { $max: "$obs.temp" }, "max_watts": { $max: "$obs.watts" }, "min_watts": { $min: "$obs.watts" }, "avg_watts": { $avg: "$obs.watts" }, "median_watts": { $median: { input: "$obs.watts", method: "approximate" } } } } ] } }, { $merge: { into: { connectionName: "mongodb1", db: "solar_db", coll: "solar_coll" }, on: ["_id"] } } ] )