Definición
Crea un efímero Procesador de flujo en el espacio de trabajo de procesamiento de flujo 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> ] )
Campos de comandos
sp.process() toma estos campos:
Campo | Tipo | Necesidad | Descripción |
|---|---|---|---|
| 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. |
| Objeto | Opcional | Objeto que define la configuración opcional para su procesador de flujo. |
| entero | Opcional | Número máximo de documentos que se devolverán a su terminal. |
| 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:
Para obtener más información, consulta Niveles. |
Comportamiento
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.
Control de acceso
El usuario que ejecuta sp.process() debe tener el
atlasAdmin rol.
Ejemplo
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"] } } ] )