Para agentes de IA: hay un índice de documentación disponible en https://www.mongodb.com/es/docs/llms.txt — versiones en markdown de todas las páginas están disponibles agregando .md a cualquier ruta URL.
Docs Menu

sp.process() (método mongosh)

sp.process()

Crea un Stream Processor efímero en el Stream Processing Workspace.actual

Este método es compatible con Atlas Stream Processing Workspaces.

El método tiene la siguiente sp.process() 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 tu procesador de flujo.

options.limit

entero

Opcional

Número máximo de documentos para devolver 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 Stream Processing. Los nombres de nivel no distinguen entre mayúsculas y minúsculas. Si omite este campo, el procesador de Stream Processing utiliza el nivel por defecto para su espacio de trabajo de Stream Processing. Debe ser uno de los siguientes valores:

  • SP2

  • SP5

  • SP10

  • SP30

  • SP50

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

sp.process() crea un Stream Processor efímero y sin nombre en el Stream Processing Workspace actual y lo inicializa de inmediato. Este Stream Processor persiste solo mientras se ejecuta. Si termina un Stream Processor efímero, debe crearlo de nuevo para usarlo.

Un procesador de transmisión efímero se detiene cuando ocurre cualquiera de los siguientes eventos:

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

  • 10 minutes elapse.

  • Usted detiene el procesador.

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

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 cuyo valor en el 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 devuelve diversas estadísticas de cada grupo. A continuación, 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"]
}
}
]
)