Docs Menu
Docs Home
/ /

sp.process() (método mongosh)

sp.process()

Crea un efímero Procesador de flujo en el espacio de trabajo de procesamiento de flujo actual.

Este método es compatible en espacios de trabajo de Atlas Stream Processing.

El método sp.process() tiene la siguiente sintaxis:

sp.process(
[
<pipeline>
]
)

sp.process() toma estos campos:

Campo
Tipo
Necesidad
Descripción

pipeline

arreglo

Requerido

La canalización de agregación de flujos es el proceso que el procesador de flujos aplica a sus datos de transmisión.

options

Objeto

Opcional

Objeto que define la configuración opcional para su procesador de flujo.

options.limit

entero

Opcional

Número máximo de documentos que se devolverán a su terminal. sp.process() finaliza la sesión después de devolver este número de documentos.

options.tier

string

Opcional

Nivel que se asignará al procesador de flujo. Los nombres de los niveles no distinguen entre mayúsculas y minúsculas. Si omite este campo, el procesador de flujo utilizará el nivel predeterminado para su espacio de trabajo de procesamiento de flujo. Debe ser uno de los siguientes valores:

  • SP2

  • SP5

  • SP10

  • SP30

  • SP50

Para obtener más información, consulta Niveles.

sp.process() Crea un procesador de flujo efímero y sin nombre en el espacio de trabajo de procesamiento de flujo actual y lo inicializa inmediatamente. Este procesador de flujo solo existe mientras se ejecuta. Si finaliza un procesador de flujo efímero, deberá crearlo de nuevo para poder utilizarlo.

Un procesador de flujo efímero se detiene cuando ocurre cualquiera de las siguientes situaciones:

  • El procesador devuelve el número de documentos establecido por options.limit.

  • 10 minutes elapse.

  • Detienes el procesador.

El usuario que ejecuta sp.process() debe tener el atlasAdmin rol.

El siguiente ejemplo crea un procesador de flujo efímero que ingiere datos de la sample_stream_solar conexión. El procesador excluye todos los documentos donde el valor del device_id campo sea,device_8 y pasa el resto a una ventana deslizante con una 10duración de segundos. Cada ventana agrupa los documentos que recibe y luego devuelve diversas estadísticas de cada grupo. El procesador de flujo fusiona estos registros en a solar_db.solar_coll través de la mongodb1 conexión.

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

Volver

sp.listWorkspaceDefaults

En esta página