Definición
Crea un Stream Processor efímero en el Stream Processing Workspace.actual
Compatibilidad
Este método es compatible en espacios de trabajo de Atlas Stream Processing.
Sintaxis
El método sp.process() tiene la siguiente sintaxis:
sp.process( [ <pipeline> ], { <options> } )
Campos de comandos
sp.process() toma estos campos:
Campo | Tipo | Necesidad | Descripción |
|---|---|---|---|
| arreglo | Requerido | Transmisión pipeline de agregación que el procesador de transmisión aplica a sus datos de transmisión. |
| Objeto | Opcional | Objeto que define la configuración opcional para tu procesador de flujo. |
| entero | Opcional | Número máximo de documentos para devolver a su terminal. |
| 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 transmisión utilizará el nivel predeterminado para su espacio de trabajo de Stream Processing. Debe ser uno de los siguientes valores:
Para obtener más información, consulta Niveles. |
Comportamiento
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.
Control de acceso
El usuario que ejecuta sp.process() debe tener el rol atlasAdmin.
Ejemplo
El siguiente ejemplo crea un procesador de transmisión efímero que ingiere datos de la conexión sample_stream_solar. El procesador excluye todos los documentos donde el valor del campo device_id es device_8, pasando el resto a una ventana deslizante con una duración de 10segundos. Cada ventana agrupa los documentos que recibe y, a continuación, devuelve varias estadísticas de cada grupo. A continuación, el procesador de transmisión fusiona estos registros con solar_db.solar_coll a través de la conexión 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"] } } ] )