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

b.process() (método mongosh)

sp.process()

Cria um processador de fluxo efêmero no espaço de trabalho de Stream Processingatual.

Esse método é suportado em Atlas Stream Processing Workspaces.

O método tem a seguinte sp.process() sintaxe:

sp.process(
[
<pipeline>
]
)

sp.process() usa estes campos:

Campo
Tipo
necessidade
Descrição

pipeline

array

Obrigatório

options

objeto

Opcional

Objeto que define configurações opcionais para o seu processador de fluxo.

options.limit

inteiro

Opcional

Número máximo de documentos a serem retornados ao seu terminal. sp.process() encerra a sessão após retornar esse número de documentos.

options.tier

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 workspace de Stream Processing. Deve ser um dos seguintes valores:

  • SP2

  • SP5

  • SP10

  • SP30

  • SP50

Para saber mais, consulte Níveis.

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.

O usuário que executa o sp.process() deve ter a função atlasAdmin.

O exemplo a seguir cria um processador de fluxo efêmero que ingere dados da sample_stream_solar conexão. O processador exclui todos os documentos em que o valor do device_id campo device_8 é, passando o restante para uma janela em cascata com 10duração de segundos. Cada janela agrupa os documentos que recebe e, em seguida, retorna várias estatísticas de cada grupo. Em seguida, o processador de fluxo mescla esses registros com solar_db.solar_coll pela mongodb1 conexão.

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"]
}
}
]
)