Definición
New in version 7.0:
Crea un procesador de flujos en el entorno de trabajo de procesamiento de flujoactual.
Compatibilidad
Este método es compatible con Atlas Stream Processing Workspaces.
Sintaxis
The sp.createStreamProcessor() method has the following syntax:
sp.createStreamProcessor( <name>, [ <pipeline> ], { <options> } )
Campos de comandos
sp.createStreamProcessor() toma estos campos:
Campo | Tipo | Necesidad | Descripción |
|---|---|---|---|
| string | Requerido | Nombre lógico para el procesador de flujos. Esto debe ser único dentro del espacio de trabajo de stream processing. |
| arreglo | Requerido | Transmisión pipeline de agregación que desea aplicar a sus datos de transmisión. |
| Objeto | Opcional | Objeto que define varias configuraciones opcionales para tu procesador de flujo. |
| Objeto | Condicional | Asignar una fila de letra muerta para tu espacio de trabajo de stream processing. Este campo es necesario si se define el campo |
| string | Condicional | Etiqueta que identifica una conexión en tu registro de conexiones. Esta conexión debe hacer referencia a un clúster de Atlas. Este campo es necesario si se define el campo |
| string | Condicional | Nombre de una base de datos de Atlas en el clúster especificado en |
| string | Condicional | Nombre de una colección en la base de datos especificada en |
| string | Opcional | El nivel al que Atlas Stream Processing asigna el procesador. Si no declara esta opción, Atlas Stream Processing asigna el procesador al nivel del Espacio de trabajo de procesamiento de flujos. Debe ser uno de los siguientes:
Para obtener más información, consulta Niveles. |
Comportamiento
sp.createStreamProcessor() crea un procesador de flujo persistente y nombrado en el espacio de trabajo actual de stream processing. Puedes inicializar este procesador de flujo con sp.processor.start(). Si intenta crear un procesador de flujo con el mismo nombre que un procesador de flujo existente, mongosh devolverá un error.
Control de acceso
El usuario que ejecuta sp.createStreamProcessor() debe tener el rol atlasAdmin.
Ejemplo
El siguiente ejemplo crea un procesador de flujo llamado solarDemo 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 útiles 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.createStreamProcessor( 'solarDemo', [ { $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"] } } ] )