Menu Docs
Página inicial do Docs
/ /

b.process() (método mongosh)

sp.process()

Cria um Stream Processor efêmero no espaço de trabalho de processamento de stream atual.

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

O método sp.process() tem a seguinte 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 seu processador de stream.

options.limit

inteiro

Opcional

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

options.tier

string

Opcional

Camada a ser atribuída ao processador de stream. Os nomes das camadas não diferenciam maiúsculas de minúsculas. Se você omitir este campo, o processador de stream utilizará a camada padrão para seu espaço de trabalho de processamento de stream. 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 processamento de fluxo atual e o inicializa imediatamente. Este processador de stream 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 é interrompido quando ocorre qualquer uma das seguintes situações:

  • O processador retorna o número de documentos definidos por options.limit.

  • 10 minutes elapse.

  • Você para 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"]
}
}
]
)

Voltar

sp.listWorkspaceDefaults

Nesta página