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

$source Etapa (Stream Processing)

$source

La etapa $source especifica una conexión en el Registro de Conexiones desde la cual transmitir datos. Se admiten los siguientes tipos de conexión:

  • Apache Kafka broker

  • flujo de cambios de colección de MongoDB

  • Flujo de cambios de la base de datos MongoDB

  • Flujo de cambios del clúster de MongoDB

  • AWS Kinesis flujo de datos

  • Arreglo de documentos

  • Programación Cron

Para operar con datos en tiempo real desde un broker de Apache Kafka, la $source etapa tiene la siguiente forma prototipo:

{
"$source": {
"connectionName": "<registered-connection>",
"topic" : ["<source-topic>", ...],
"timeField": {
$toDate | $dateFromString: <expression>
},
"partitionIdleTimeout": {
"size": <duration-number>,
"unit": "<duration-unit>"
},
"schemaRegistry": {
"connectionName": "<schema-registry-name>",
},
"config": {
"auto_offset_reset": "<start-event>",
"group_id": "<group-id>",
"keyFormat": "<deserialization-type>",
"keyFormatError": "<error-handling>"
},
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

connectionName

string

Requerido

Etiqueta que identifica la conexión en el Registro de conexiones, para ingerir datos de esta.

topic

string o arreglo de strings

Requerido

Nombre de uno o más temas de Apache Kafka de los que transmitir mensajes. Si desea transmitir mensajes de más de un tema, especifíquelos en un arreglo.

timeField

Documento

Opcional

Documento que define una marca de tiempo autorizada para los mensajes entrantes.

Si utilizas timeField, debes definirlo como uno de los siguientes:

  • un :expresión:$toDate que toma un campo de mensaje de origen como argumento.

  • un :expresión:$dateFromString que toma un campo de mensaje de origen como argumento.

Si no declaras un timeField, Atlas Stream Processing creará una marca de tiempo a partir de la marca de tiempo del mensaje proporcionada por la fuente.

partitionIdleTimeout

Documento

Opcional

Documento especificando la cantidad de tiempo que se permite que una partición esté inactiva antes de que se ignore en los cálculos de marcas de agua.

Este campo está desactivado por defecto. Para gestionar las particiones que no avanzan debido a la inactividad, establece un valor para este campo.

partitionIdleTimeout.size

entero

Opcional

Número que especifica la duración del tiempo de espera de inactividad de la partición.

partitionIdleTimeout.unit

string

Opcional

Unidad de tiempo para la duración del tiempo de espera por inactividad de la partición.

El valor de unit puede ser uno de los siguientes:

  • "ms" (milisegundo)

  • "second"

  • "minute"

  • "hour"

  • "day"

schemaRegistry

Documento

Opcional

Documento que habilita el uso de un Registro de esquemas para apoyar la lectura desde una fuente serializada en Avro.

Para habilitar esta funcionalidad, debe crear una conexión con el Registro de esquemas.

schemaRegistry.connectionName

string

Condicional

Nombre de la conexión del Registro de esquemas que se utilizará para la deserialización de Avro.

config

Documento

Opcional

Documento que contiene campos que sobrescriben diversos valores por defecto.

config.auto_offset_reset

string

Opcional

Especifica con qué evento del tema de origen de Apache Kafka se debe comenzar la ingesta. auto_offset_reset toma los siguientes valores:

  • end, latest, o largest : para comenzar la ingesta desde el evento más reciente en el tema en el momento en que se inicializa la agregación.

  • earliest, beginning, o smallest : para comenzar la ingestión desde el evento más temprano en el tema.

Se establece por defecto en latest.

config.group_id

string

Opcional

ID del grupo de consumidores de Kafka para asociar con el procesador de flujo. Si se omite, Atlas Stream Processing asocia el espacio de trabajo de Stream Processing con una ID generada automáticamente en el siguiente formato:

asp-${streamProcessorId}-consumer

Atlas Stream Processing genera automáticamente un valor para este parámetro en todos los procesadores de flujos persistentes. Para los procesadores de flujo efímeros definidos con sp.process(), este parámetro solo se configura si lo define manualmente.

config.enable_auto_commit

booleano

Condicional

Indicador que determina la política de confirmaciones para los desfases de partición de los intermediarios de Kafka. Atlas Stream Processing admite dos políticas de confirmación:

  • Si configuras este parámetro en true, Atlas Stream Processing confirma los desplazamientos cada vez que la etapa $source pasa datos al siguiente operador.

  • Si estableces este parámetro en false, los procesadores de flujo confirman los desplazamientos de partición cuando Atlas Stream Processing realiza un punto de control.

Para un procesador de flujos efímero definido con sp.process(), este parámetro se asigna por defecto a false a menos que establezca group_id. De lo contrario, el valor por defecto es true.

Para obtener más información sobre los offsets cuando se utiliza Kafka como $source, consulte Fuentes de Kafka y offsets del grupo de consumidores.

config.keyFormat

string

Opcional

Tipo de dato utilizado para deserializar los datos de clave de Apache Kafka. Debe ser uno de los siguientes valores:

  • "binData"

  • "string"

  • "json"

  • "int"

  • "long"

Se establece por defecto en binData.

config.keyFormatError

string

Opcional

Cómo gestionar los errores encontrados al deserializar los datos clave de Apache Kafka. Debe ser uno de los siguientes valores:

  • dlq, que guarda el documento en tu fila de letra muerta.

  • passThrough, que envía el documento a la siguiente etapa sin datos clave.

Nota

Atlas Stream Processing requiere que los documentos en el flujo de datos de origen sean válidos json o ejson. Atlas Stream Processing establece que los documentos que no cumplen este requisito se asignen a tu fila de letra muerta si has configurado una.

Un flujo de cambios de la colección de Atlas permite que las aplicaciones accedan a los cambios de datos en tiempo real en una sola colección. Para aprender a abrir un flujo de cambios en una colección, consulta Change Streams.

Al usar un flujo de cambios $source, configura el clúster de origen con una oplog window de al menos 24 horas.

Para leer el flujo de cambios, Atlas Stream Processing escanea la colección oplog. Como resultado, es posible que veas COLLSCAN advertencias en tus registros. Estas advertencias indican un comportamiento normal y no señalan un error.

Si configuras config.fullDocument o config.fullDocumentBeforeChange en required, activa changeStreamPreAndPostImages en cada colección antes de cualquier operación de guardado que desees capturar. Si una postimagen no está disponible para un evento porque no se habilitó la funcionalidad cuando ocurrió el guardado, o porque la postimagen ha caducado, el procesador de flujos falla. Para obtener instrucciones sobre cómo habilitar las preimágenes y las postimágenes, consulta Change Streams con pre/postimágenes de documentos.

Importante

Solo se puede reanudar una secuencia de cambios mientras la colección de oplog contenga la operación que identifica su token de reanudación. Para garantizar que sus aplicaciones puedan reanudar sus secuencias de cambios, configure una ventana mínima de oplog mayor que la interrupción más larga prevista.

Para operar sobre datos en transmisión de un flujo de cambios de una colección de Atlas, la etapa $source tiene la siguiente forma prototipo:

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"db" : "<source-db>",
"coll" : ["<source-coll>",...],
"initialSync": {
"enable": <boolean>,
"parallelism": <integer>
},
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}],
"maxAwaitTimeMS": <time-ms>,
}
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

connectionName

string

Condicional

Etiqueta que identifica la conexión en el Registro de conexiones, para ingerir datos de esta.

timeField

Documento

Opcional

Documento que define una marca de tiempo autorizada para los mensajes entrantes.

Si utilizas timeField, debes definirlo como uno de los siguientes:

  • una expresión $toDate que toma un campo de mensaje de origen como argumento

  • una expresión $dateFromString que toma un campo de mensaje de origen como argumento.

Si no declaras un timeField, Atlas Stream Processing creará una marca de tiempo a partir de la marca de tiempo del mensaje proporcionada por la fuente.

db

string

Requerido

Nombre de la base de datos MongoDB alojada en la instancia de Atlas especificada por connectionName. El flujo de cambios de esta base de datos actúa como la fuente de datos de transmisión.

coll

string o arreglo de strings

Requerido

Nombre de una o más colecciones de MongoDB alojadas en la instancia de Atlas especificada por connectionName. El flujo de cambios de estas colecciones actúa como origen de datos en tiempo real. Si omite este campo, su procesador de flujos se abastecerá de un flujo de cambios de la base de datos de MongoDB.

initialSync

Documento

Opcional

Documento que contiene initialSync parámetros de configuración.

El procesamiento de flujos de Atlas initialSync permite ingerir documentos preexistentes en una colección de Atlas como si fueran documentos changeEvent de inserción. Si habilita initialSync, al iniciar el procesador de flujos, este primero ingiere y procesa todos los documentos existentes en la colección antes de proceder a ingerir y procesar los nuevos documentos changeEvent entrantes. Una vez completado el initialSync, no se repite.

Si coll nombra más de una colección, Atlas Stream Processing sincroniza todas las colecciones de la lista. Para obtener información sobre el orden de sincronización de varias colecciones, los límites de nivel y los ejemplos, consulte Configurar la sincronización inicial de varias colecciones.

Atlas Stream Processing fija el conjunto de colecciones a sincronizar cuando se inicia el procesador de transmisiones. No sincroniza las colecciones que cree después de ese punto, y si falla la sincronización de una colección, falla todo el procesador de transmisiones.

Si habilitas initialSync, no puedes utilizar las etapas $hoppingWindow, $sessionWindow o $tumblingWindow en tu pipeline.

IMPORTANTE: Consulte Limitaciones para conocer las consideraciones y restricciones de initialSync.

initialSync.enable

booleano

Condicional

Determina si se debe habilitar initialSync. Si declara un campo initialSync, debe establecer este campo.

initialSync.parallelism

entero

Opcional

Determina el nivel de paralelismo con el que se procesará la operación initialSync. Si no especificas un valor, este es por defecto 1.

Si coll nombra más de una colección, este valor se aplica a la lista de colecciones en su conjunto en lugar de a cada colección.

initialSync aplica este valor a colecciones con cualquier tipo _id compatible, incluidos los tipos que Atlas Stream Processing no puede leer en un orden predecible.

El valor máximo que puede establecer depende del nivel de su procesador de transmisión. Para obtener más información, consulte Guía de selección de nivel de Atlas Stream Processing.

Cada procesador de flujo tiene un valor máximo de paralelismo acumulado determinado por su nivel. El paralelismo acumulativo de un procesador de flujo se calcula de la siguiente manera:

parallelism total - parallelized stages

Donde parallelism total es la suma de todos los valores parallelism mayores que 1 a través de las etapas $source, $lookup, $merge, $emit y $externalFunction, y parallelized stages es la cantidad de estas etapas con valores parallelism mayores que 1.

Por ejemplo, si tu etapa $source establece un valor parallelism de 4, tu etapa $lookup no define un valor parallelism (por lo tanto, se toma el valor por defecto 1), y tu etapa $merge define un valor parallelism de 2, entonces tiene dos parallelized stages y el paralelismo acumulativo de su procesador de flujo se calcula como (4 + 2) - 2.

Si un procesador de flujo supera el paralelismo acumulativo máximo para su nivel, Atlas Stream Processing genera un error e informa del nivel mínimo de procesador requerido para el nivel de paralelismo previsto. Debe escalar el procesador a un nivel superior o reducir los valores de paralelismo de sus etapas para resolver el error. Para aprender más, consulte Procesamiento de flujos.

readPreference

string

Opcional

Preferencia de lectura para flujo de cambios y initialSync operaciones.

Se establece por defecto en primary.

readPreferenceTags

arreglo

Opcional

Leer etiquetas de preferencia de lectura para el flujo de cambios y las operaciones initialSync.

config

Documento

Opcional

Documento que contiene campos que sobrescriben diversos valores por defecto.

config.startAfter

token

Condicional

El evento de cambio después del cual la fuente comienza a reporte. Esto toma la forma de un token de currículum.

Puedes utilizar solo uno de config.startAfter o config.startAtOperationTime.

config.startAtOperationTime

timestamp | fecha

Condicional

El operation time después del cual la fuente debe comenzar a reportar.

Puedes utilizar solo uno de config.startAfter o config.startAtOperationTime.

Acepta valores MongoDB Extended JSON $date u $timestamp.

config.fullDocument

string

Condicional

Configuración que controla si una fuente de flujo de cambios debe devolver un documento completo o sólo las modificaciones cuando se produce una actualización. Debe ser uno de los siguientes:

  • default : No devuelve un documento completo para las operaciones de update.

  • updateLookup : Devuelve la versión actual del documento completo que cuenta con la aprobación mayoritaria, además de los cambios realizados por la operación de actualización.

  • required Debe devolver el documento completo. Si no se dispone del documento completo, el procesador de flujo falla. Esto no se aplica a las operaciones de eliminación, que siguen generando eventos de eliminación aunque no se disponga del documento completo.

  • whenAvailable : Devuelve un documento completo cada vez que haya uno disponible, en caso contrario, devuelve los cambios.

Para usar required o whenAvailable con un flujo de cambios de colección, debe habilitar las imágenes previas y posteriores del flujo de cambios en esa colección.

config.fullDocumentOnly

booleano

Condicional

Configuración que controla si una fuente de flujos de cambios devuelve todo el documento del evento de cambio, incluyendo todos los metadatos, o solo el contenido de fullDocument. Si se establece en true, la fuente devuelve únicamente el contenido de fullDocument.

Si fullDocument está configurado como required o whenAvailable, debe habilitar las imágenes previas y posteriores del flujo de cambios en esa colección.

config.fullDocumentBeforeChange

string

Opcional

Especifica si una fuente de flujo de cambios debe incluir el documento completo en su estado original "antes de los cambios" en la salida. Debe ser uno de los siguientes:

  • off : Omite el campo fullDocumentBeforeChange.

  • required : Debe devolverse un documento completo en su estado previo a los cambios. Si no hay disponibilidad de un documento completo en su estado previo a los cambios, el procesador de transmisión falla.

  • whenAvailable : devuelve un documento completo en su estado previo a los cambios siempre que esté disponible; de lo contrario, omite el campo fullDocumentBeforeChange.

Si no se especifica un valor para fullDocumentBeforeChange, se usa por defecto off.

Para utilizar este campo con una colección de flujo de cambios, debe habilitar la Preimágenes y Postimágenes de flujo de cambios en esa colección.

config.pipeline

Documento

Opcional

Especifica un pipeline de agregación para filtrar la salida del flujo de cambios antes de pasarla para procesamiento adicional. Este pipeline debe respetar los parámetros que se describen en Modifica la salida del change stream.

IMPORTANTE: Cada evento de cambio incluye los campos wallTime y clusterTime. Las etapas de Atlas Stream Processing después de $source esperan recibir estos campos tal como fueron ingeridos por el procesador. Para garantizar el procesamiento adecuado de los datos de Change Stream, no modifiques estos campos en $source.config.pipeline.

config.maxAwaitTimeMS

entero

Opcional

Tiempo máximo, en milisegundos, para esperar a que se reporten nuevos cambios de datos al cursor del flujo de cambios antes de devolver un lote vacío.

Se establece por defecto en 1000.

Un flujo de cambios de base de datos de Atlas permite que las aplicaciones accedan a cambios de datos en tiempo real en una única base de datos. Para aprender a abrir un flujo de cambios en una base de datos, consulta Change Streams.

Al usar un flujo de cambios $source, configura el clúster de origen con una oplog window de al menos 24 horas.

Para leer el flujo de cambios, Atlas Stream Processing escanea la colección oplog. Como resultado, es posible que veas COLLSCAN advertencias en tus registros. Estas advertencias indican un comportamiento normal y no señalan un error.

Si configuras config.fullDocument o config.fullDocumentBeforeChange en required, activa changeStreamPreAndPostImages en cada colección antes de cualquier operación de guardado que desees capturar. Si una postimagen no está disponible para un evento porque no se habilitó la funcionalidad cuando ocurrió el guardado, o porque la postimagen ha caducado, el procesador de flujos falla. Para obtener instrucciones sobre cómo habilitar las preimágenes y las postimágenes, consulta Change Streams con pre/postimágenes de documentos.

Importante

Solo se puede reanudar una secuencia de cambios mientras la colección de oplog contenga la operación que identifica su token de reanudación. Para garantizar que sus aplicaciones puedan reanudar sus secuencias de cambios, configure una ventana mínima de oplog mayor que la interrupción más larga prevista.

Para operar sobre datos en transmisión de un flujo de cambios de la base de datos Atlas, la etapa $source tiene la siguiente forma de prototipo:

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"db" : "<source-db>",
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}]
},
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

connectionName

string

Condicional

Etiqueta que identifica la conexión en el Registro de conexiones, para ingerir datos de esta.

timeField

Documento

Opcional

Documento que define una marca de tiempo autorizada para los mensajes entrantes.

Si utilizas timeField, debes definirlo como uno de los siguientes:

  • una expresión $toDate que toma un campo de mensaje de origen como argumento

  • una expresión $dateFromString que toma un campo de mensaje de origen como argumento.

Si no declaras un timeField, Atlas Stream Processing creará una marca de tiempo a partir de la marca de tiempo del mensaje proporcionada por la fuente.

db

string

Requerido

Nombre de la base de datos MongoDB alojada en la instancia de Atlas especificada por connectionName. El flujo de cambios de esta base de datos actúa como la fuente de datos de transmisión.

readPreference

string

Opcional

preferencia de lectura para operaciones de flujo de cambios.

Se establece por defecto en primary.

readPreferenceTags

arreglo

Opcional

Lea las etiquetas de preferencia de lectura para las operaciones de flujo de cambios.

config

Documento

Opcional

Documento que contiene campos que sobrescriben diversos valores por defecto.

config.startAfter

token

Condicional

El evento de cambio después del cual la fuente comienza a reporte. Esto toma la forma de un token de currículum.

Puedes utilizar solo uno de config.startAfter o config.startAtOperationTime.

config.startAtOperationTime

timestamp | fecha

Condicional

El operation time después del cual la fuente debe comenzar a reportar.

Puedes utilizar solo uno de config.startAfter o config.startAtOperationTime.

Acepta valores MongoDB Extended JSON $date u $timestamp.

config.fullDocument

string

Condicional

Configuración que controla si una fuente de flujo de cambios debe devolver un documento completo o sólo las modificaciones cuando se produce una actualización. Debe ser uno de los siguientes:

  • default : Utiliza el comportamiento por defecto del servidor. No devuelve un documento completo para operaciones update.

  • updateLookup : Devuelve la versión actual del documento completo que cuenta con la aprobación mayoritaria, además de los cambios realizados por la operación de actualización.

  • required Debe devolver el documento completo. Si no se dispone del documento completo, el procesador de flujo falla. Esto no se aplica a las operaciones de eliminación, que siguen generando eventos de eliminación aunque no se disponga del documento completo.

  • whenAvailable : Devuelve un documento completo cada vez que haya uno disponible, en caso contrario, devuelve los cambios.

Si no se especifica un valor para fullDocument, se usa por defecto default.

Para usar required o whenAvailable con un flujo de cambios de base de datos, debe habilitar las imágenes previas y posteriores del flujo de cambios en cada colección de esa base de datos.

config.fullDocumentOnly

booleano

Condicional

Configuración que controla si una fuente de flujos de cambios devuelve todo el documento del evento de cambio, incluyendo todos los metadatos, o solo el contenido de fullDocument. Si se establece en true, la fuente devuelve únicamente el contenido de fullDocument.

Si fullDocument está configurado como required o whenAvailable, debe habilitar las imágenes previas y posteriores del flujo de cambios en cada colección de esa base de datos.

config.fullDocumentBeforeChange

string

Opcional

Especifica si una fuente de flujo de cambios debe incluir el documento completo en su estado original "antes de los cambios" en la salida. Debe ser uno de los siguientes:

  • off : Omite el campo fullDocumentBeforeChange.

  • required : Debe devolverse un documento completo en su estado previo a los cambios. Si no hay disponibilidad de un documento completo en su estado previo a los cambios, el procesador de transmisión falla.

  • whenAvailable : devuelve un documento completo en su estado previo a los cambios siempre que esté disponible; de lo contrario, omite el campo fullDocumentBeforeChange.

Si no se especifica un valor para fullDocumentBeforeChange, se usa por defecto off.

Para usar este campo con un flujo de cambios de base de datos, debes habilitar imágenes anteriores y posteriores en cada colección de esa base de datos.

config.pipeline

Documento

Opcional

Especifica un pipeline de agregación para filtrar la salida del flujo de cambios en el punto de origen. Esta pipeline debe ajustarse a los parámetros descritos en Modificar la salida del flujo de cambios.

IMPORTANTE: Cada evento de cambio incluye los campos wallTime y clusterTime. Las etapas de Atlas Stream Processing después de $source esperan recibir estos campos tal como fueron ingeridos por el procesador. Para garantizar el procesamiento adecuado de los datos de Change Stream, no modifiques estos campos en $source.config.pipeline.

config.maxAwaitTimeMS

entero

Opcional

Tiempo máximo, en milisegundos, para esperar a que se reporten nuevos cambios de datos al cursor del flujo de cambios antes de devolver un lote vacío.

Se establece por defecto en 1000.

Al usar un flujo de cambios $source, configura el clúster de origen con una oplog window de al menos 24 horas.

Para leer el flujo de cambios, Atlas Stream Processing escanea la colección oplog. Como resultado, es posible que veas COLLSCAN advertencias en tus registros. Estas advertencias indican un comportamiento normal y no señalan un error.

Si configuras config.fullDocument o config.fullDocumentBeforeChange en required, activa changeStreamPreAndPostImages en cada colección antes de cualquier operación de guardado que desees capturar. Si una postimagen no está disponible para un evento porque no se habilitó la funcionalidad cuando ocurrió el guardado, o porque la postimagen ha caducado, el procesador de flujos falla. Para obtener instrucciones sobre cómo habilitar las preimágenes y las postimágenes, consulta Change Streams con pre/postimágenes de documentos.

Importante

Solo se puede reanudar una secuencia de cambios mientras la colección de oplog contenga la operación que identifica su token de reanudación. Para garantizar que sus aplicaciones puedan reanudar sus secuencias de cambios, configure una ventana mínima de oplog mayor que la interrupción más larga prevista.

Para operar datos de transmisión de todo un flujo de cambio del clúster Atlas, la etapa $source tiene la siguiente forma prototipo:

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}]
},
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

connectionName

string

Condicional

Etiqueta que identifica la conexión en el Registro de conexiones, para ingerir datos de esta.

timeField

Documento

Opcional

Documento que define una marca de tiempo autorizada para los mensajes entrantes.

Si utilizas timeField, debes definirlo como uno de los siguientes:

  • una expresión $toDate que toma un campo de mensaje de origen como argumento

  • una expresión $dateFromString que toma un campo de mensaje de origen como argumento.

Si no declaras un timeField, Atlas Stream Processing creará una marca de tiempo a partir de la marca de tiempo del mensaje proporcionada por la fuente.

readPreference

string

Opcional

preferencia de lectura para operaciones de flujo de cambios.

Se establece por defecto en primary.

readPreferenceTags

arreglo

Opcional

Lea las etiquetas de preferencia de lectura para las operaciones de flujo de cambios.

config

Documento

Opcional

Documento que contiene campos que sobrescriben diversos valores por defecto.

config.startAfter

token

Condicional

El evento de cambio después del cual la fuente comienza a reporte. Esto toma la forma de un token de currículum.

Puedes utilizar solo uno de config.startAfter o config.startAtOperationTime.

config.startAtOperationTime

fecha | marca temporal

Condicional

El operation time después del cual la fuente debe comenzar a reportar.

Puedes utilizar solo uno de config.startAfter o config.startAtOperationTime.

Acepta valores MongoDB Extended JSON $date u $timestamp.

config.fullDocument

string

Condicional

Configuración que controla si una fuente de flujo de cambios debe devolver un documento completo o sólo las modificaciones cuando se produce una actualización. Debe ser uno de los siguientes:

  • default : Utiliza el comportamiento por defecto del servidor. No devuelve un documento completo para operaciones update.

  • updateLookup : Devuelve la versión actual del documento completo que cuenta con la aprobación mayoritaria, además de los cambios realizados por la operación de actualización.

  • required Debe devolver el documento completo. Si no se dispone del documento completo, el procesador de flujo falla. Esto no se aplica a las operaciones de eliminación, que siguen generando eventos de eliminación aunque no se disponga del documento completo.

  • whenAvailable : Devuelve un documento completo cada vez que haya uno disponible, en caso contrario, devuelve los cambios.

Si no se especifica un valor para fullDocument, se usa por defecto default.

Para usar required o whenAvailable con un flujo de cambios de clúster, debe habilitar las imágenes previas y posteriores del flujo de cambios en cada colección de ese clúster.

config.fullDocumentOnly

booleano

Condicional

Configuración que controla si una fuente de flujos de cambios devuelve todo el documento del evento de cambio, incluyendo todos los metadatos, o solo el contenido de fullDocument. Si se establece en true, la fuente devuelve únicamente el contenido de fullDocument.

Si fullDocument está configurado como required o whenAvailable, debe habilitar las imágenes previas y posteriores del flujo de cambios en cada colección de ese clúster.

config.fullDocumentBeforeChange

string

Opcional

Especifica si una fuente de flujo de cambios debe incluir el documento completo en su estado original "antes de los cambios" en la salida. Debe ser uno de los siguientes:

  • off : Omite el campo fullDocumentBeforeChange.

  • required : Debe devolverse un documento completo en su estado previo a los cambios. Si no hay disponibilidad de un documento completo en su estado previo a los cambios, el procesador de transmisión falla.

  • whenAvailable : devuelve un documento completo en su estado previo a los cambios siempre que esté disponible; de lo contrario, omite el campo fullDocumentBeforeChange.

Si no se especifica un valor para fullDocumentBeforeChange, se usa por defecto off.

Para usar este campo con un flujo de cambios de base de datos, debes habilitar imágenes anteriores y posteriores en cada colección de esa base de datos.

config.pipeline

Documento

Opcional

Especifica un pipeline de agregación para filtrar la salida del flujo de cambios en el punto de origen. Esta pipeline debe ajustarse a los parámetros descritos en Modificar la salida del flujo de cambios.

Ten en cuenta que el Atlas Stream Processing espera recibir los campos wallTime y clusterTime de cada Evento de Cambio ingerido. Para garantizar el procesamiento adecuado de los datos de Change Stream, no modifiques estos campos en $source.config.pipeline.

config.maxAwaitTimeMS

entero

Opcional

Tiempo máximo, en milisegundos, para esperar a que se reporten nuevos cambios de datos al cursor del flujo de cambios antes de devolver un lote vacío.

Se establece por defecto en 1000.

Atlas Stream Processing admite la creación de conexiones Private Link con flujos AWS Kinesis. Para obtener más información, consulte Agregar una conexión privada de Kinesis.

Para operar con datos de un flujo de datos de AWS Kinesis, la etapa $source tiene el siguiente prototipo:

{
"$source": {
"connectionName": "<registered-connection>",
"stream": "<stream-name>",
"region": "<aws-region>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"tsFieldName": "<field-name>",
"shardIdleTimeout": {
"size": <duration-number>,
"unit": "<duration-unit>"
},
"config": {
"consumerARN": "<aws-arn>",
"initialPosition": <initial-position>,
reshardDetectionIntervalSecs: <interval>
}
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

connectionName

string

Requerido

Etiqueta que identifica la conexión en el Registro de Conexiones, de la cual se obtiene la información.

config.consumerARN

string

Requerido

ARN correspondiente a un consumidor de Kinesis. Su consumidor debe usar Fan-out mejorado.

stream

string

Requerido

Flujo de datos de AWS Kinesis desde el cual transmitir mensajes.

region

string

Condicional

Región AWS en la que existe el flujo especificado. Kinesis admite varios flujos de datos con el mismo nombre en diferentes regiones. Si se usa el mismo nombre para los flujos de datos en dos o más regiones dentro de la misma conexión, se debe usar este campo para especificar qué combinación de nombre y región usar.

timeField

Documento

Opcional

Documento que define una marca de tiempo autorizada para los mensajes entrantes.

Si utilizas timeField, debes definirlo como uno de los siguientes:

  • una expresión $toDate que toma un campo de mensaje de origen como argumento

  • una expresión $dateFromString que toma un campo de mensaje de origen como argumento.

Si no declaras un timeField, Atlas Stream Processing creará una marca de tiempo a partir de la marca de tiempo del mensaje proporcionada por la fuente.

tsFieldName

string

Opcional

Nombre del campo para la marca de tiempo en el documento proyectado. Utiliza este campo para anular el nombre de campo de la marca de tiempo por defecto.

shardIdleTimeout

Documento

Opcional

Documento que especifica el tiempo que se permite que una partición esté inactiva antes de que se ignore en los cálculos de marcas de agua.

Este campo está desactivado por defecto. Para gestionar las particiones que no avanzan debido a la inactividad, establece un valor para este campo.

shardIdleTimeout.size

Documento

Opcional

Número que especifica la duración del tiempo de espera de inactividad de la partición.

shardIdleTimeout.unit

Documento

Opcional

Unidad de tiempo para la duración del tiempo de espera inactivo de la partición.

El valor de unit puede ser uno de los siguientes:

  • "ms" (milisegundo)

  • "second"

  • "minute"

  • "hour"

  • "day"

config

Documento

Opcional

Documento que contiene campos que sobrescriben diversos valores por defecto.

config.initialPosition

string

Opcional

Posición en el historial de la transmisión de datos de Kinesis desde la cual se comenzará a ingresar mensajes. Debe ser uno de los siguientes:

  • "TRIM_HORIZON": Comienza a ingerir desde el mensaje más antiguo en la partición.

  • "LATEST": Comienza a ingerir desde el mensaje más reciente de la partición.

  • "AT_TIMESTAMP"comenzar la ingesta a partir de un momento específico. Se requiere config.atTimestamp.

  • "AT_SEQUENCE_NUMBER": Comenzar la ingestión en un número de secuencia específico.

  • "AFTER_SEQUENCE_NUMBER": comienza la ingestión después de un número de secuencia específico.

Se establece por defecto en "LATEST".

config.atTimestamp

fecha

Condicional

La marca de tiempo desde la que comenzar a recibir mensajes. Requerido cuando config.initialPosition es "AT_TIMESTAMP".

config.reshardDetectionIntervalSecs

entero

Opcional

Intervalo, en segundos, entre chequeos de la velocidad del flujo de datos a través de tu flujo de Kinesis para fines de rebalanceo.

Por defecto son 300 segundos.

Para operar sobre un arreglo de documentos, la etapa $source tiene el siguiente formato prototipo:

{
"$source": {
"timeField": {
$toDate | $dateFromString: <expression>
},
"documents" : [{source-doc},...] | <expression>
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

timeField

Documento

Opcional

Documento que define una marca de tiempo autorizada para los mensajes entrantes.

Si utilizas timeField, debes definirlo como uno de los siguientes:

  • una expresión $toDate que toma un campo de mensaje de origen como argumento

  • una expresión $dateFromString que toma un campo de mensaje de origen como argumento.

Si no declaras un timeField, Atlas Stream Processing creará una marca de tiempo a partir de la marca de tiempo del mensaje proporcionada por la fuente.

documents

arreglo

Condicional

Arreglo de documentos para usar como fuente de datos en streaming. El valor de este campo puede ser un arreglo de objetos o una expresión que se evalúe como un arreglo de objetos. No utilice este campo cuando use el campo connectionName.

Para generar un documento de forma recurrente en lugar de leerlo desde una conexión, la etapa $source tiene el siguiente prototipo:

{
"$source": {
"schedule": "<cron-expression>",
"tsFieldName": "<timestamp-field-name>"
}
}

La etapa $source procesa un documento con los siguientes campos:

Campo
Tipo
Necesidad
Descripción

schedule

string

Requerido

Expresión cron de seis campos que determina cuándo la etapa emite un documento. Los campos son, en orden: segundo, minuto, hora, día del mes, mes y día de la semana. Atlas Stream Processing evalúa la expresión en UTC.

Cada campo acepta * para todos los valores, un único valor, un rango como 1-5, un intervalo como 0/15 o una lista separada por comas como 1,15,30. Puede usar los nombres JAN a DEC para los meses y SUN a SAT para los días de la semana. No incluya el prefijo CRON_TZ.

Para obtener información sobre cómo Atlas Stream Processing resuelve una expresión que restringe tanto el día del mes como el día de la semana, consulte la sección Comportamiento.

tsFieldName

string

Opcional

Nombre del campo en el que la etapa proyecta la marca de tiempo programada. Por defecto es _ts.

$source debe ser la primera etapa de cualquier pipeline en el que aparezca. Solo puedes usar una $source etapa por pipeline.

Para etapas de Kafka $source, Atlas Stream Processing lee en paralelo desde múltiples particiones dentro del tema de origen. El límite de partición lo determina tu nivel de procesador. Para obtener detalles adicionales, revise la referencia de facturación de Stream Processing.

Para las etapas cron $source, cada hora programada produce un único documento vacío que lleva la marca de hora programada. Utilice etapas posteriores en el pipeline para rellenar el documento.

Cuando una expresión schedule restringe tanto el día del mes como el día de la semana, Atlas Stream Processing emite un documento solo en los momentos programados que cumplen con ambos campos. Esto difiere de las implementaciones de cron, que emiten un documento cuando coincide cualquiera de los campos.

Si un procesador de flujo con una etapa cron $source se reinicia, emite un documento para cada hora programada que omitió mientras no estaba en ejecución. No emite documentos para las horas programadas que omitió después de que lo detuviera.

Un procesador de flujo con una etapa cron $source se ejecuta continuamente entre intervalos programados, y Atlas lo factura según su nivel de procesador para todo ese período. Una programación poco frecuente no reduce el costo de ejecución del procesador. Para obtener más información, consulte la referencia de facturación del procesamiento de flujo.

Una fuente de datos de transmisión genera informes meteorológicos detallados desde varias ubicaciones, cumpliendo con el esquema del Conjunto de Datos Meteorológicos de Muestra. La siguiente agregación tiene tres etapas:

  1. La etapa $source establece una conexión con el broker Apache Kafka que recopila estos informes en un tema llamado my_weatherdata, exponiendo cada registro a medida que se ingiere a las etapas de agregación subsiguientes. Esta etapa también sobrescribe el nombre del campo de marca de tiempo que proyecta, estableciéndolo como ingestionTime.

  2. La etapa $match excluye los documentos que tienen un dewPoint.value menor o igual a 5.0 y pasa los documentos con dewPoint.value mayor que 5.0 a la siguiente etapa.

  3. La fase $merge escribe la salida en una colección Atlas denominada stream en la base de datos sample_weatherstream. Si no existe tal base de datos o colección, Atlas los creará.

[{
"$source": {
"connectionName": "sample_weatherdata",
"topic": "my_weatherdata"
}
},
{
"$match": { "dewPoint.value": { "$gt": 5 } }
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "stream"
}
}
}]

Para ver los documentos de la colección sample_weatherstream.stream resultante, conéctate a un clúster de Atlas y ejecuta el siguiente comando:

db.getSiblingDB("sample_weatherstream").stream.find()

Nota

Lo anterior es un ejemplo representativo. Los datos en transmisión no son estáticos, y cada usuario ve documentos distintos.

La siguiente agregación ingiere datos de la fuente cluster0-collection, que se conecta a un clúster de Atlas cargado con el conjunto de datos de muestra. Para aprender cómo crear un espacio de trabajo de Stream Processing y agregar una conexión a un clúster de Atlas al registro de conexión, consulte Introducción a Atlas Stream Processing. Esta agregación ejecuta dos etapas para abrir un flujo de cambios y registrar los cambios en la colección data de la base de datos sample_weatherdata:

  1. La etapa $source se conecta a la fuente cluster0-collection y abre un flujo de cambios en la colección data de la base de datos sample_weatherdata.

  2. La etapa $merge escribe los documentos de flujo de cambios filtrados en una colección de Atlas denominada data_changes en la base de datos sample_weatherdata. Si no existe tal colección, Atlas la crea.

[{
"$source": {
"connectionName": "cluster0-connection",
"db": "sample_weatherdata",
"coll": "data"
}
},
{
"$merge": {
"into": {
"connectionName": "cluster0-connection",
"db": "sample_weatherdata",
"coll": "data_changes"
}
}
}]

El siguiente comando mongosh borra un documento data:

db.getSiblingDB("sample_weatherdata").data.deleteOne(
{ _id: ObjectId("5553a99ae4b02cf715120e4b") }
)

Después de que el documento data sea eliminado, el procesador de flujo registra el documento de evento del flujo de cambios en la colección sample_weatherdata.data_changes. Para ver los documentos en la colección resultante de sample_weatherdata.data_changes, usa mongosh para conectarte a tu clúster de Atlas y ejecutar el siguiente comando:

db.getSiblingDB("sample_weatherdata").data_changes.find()

La siguiente agregación ingiere datos de la fuente db-change-stream-connection, que se conecta a un clúster de Atlas cargado con el conjunto de datos Sample Mflix Dataset Collections. Para aprender a crear un espacio de trabajo de procesamiento de flujos y agregar una conexión a un clúster de Atlas al registro de conexiones, consulte Introducción al procesamiento de flujos de Atlas. Esta agregación ejecuta dos etapas para abrir un flujo de cambios contra la base de datos de origen sample_mflix y registrar los cambios en la colección db_changes en la base de datos de destino sample_mflix_changes:

  1. La etapa $source se conecta a la fuente db-change-stream-connection y abre un flujo de cambios contra la base de datos de origen sample_mflix. El campo config.startAtOperationTime establece el tiempo después del cual la fuente comienza a informar. En este ejemplo, el valor startAt comienza un minuto antes.

  2. La etapa $merge escribe los documentos del flujo de cambios en una colección de Atlas llamada db_changes en la base de datos de destino sample_mflix_changes.

const startAt = new Date(Date.now() - 60 * 1000);
const pipeline = [
{
$source: {
connectionName: "db-change-stream-connection",
db: "sample_mflix",
config: {
startAtOperationTime: startAt
}
}
},
{
$merge: {
into: {
connectionName: "db-change-stream-connection",
db: "sample_mflix_changes",
coll: "db_changes"
}
}
}
];

Ejecute los siguientes comandos en su clúster de origen para verificar el comportamiento del procesador de flujo. Para conectarse, consulte Conectarse a un clúster mediante mongosh.

Inserte un documento en la colección movies y un documento en la colección comments en la base de datos de origen sample_mflix:

db.getSiblingDB("sample_mflix").movies.insertOne({
title: "The Stream Processor",
year: 2026
})
db.getSiblingDB("sample_mflix").comments.insertOne({
name: "Ada Lovelace",
text: "A fine film about data in motion."
})

Después de insertar los documentos, el procesador de flujo guardar un documento de evento de flujo de cambios para cada inserción en la colección sample_mflix_changes.db_changes. Ejecute el siguiente comando en su clúster para ver los documentos en la colección sample_mflix_changes.db_changes resultante:

db.getSiblingDB("sample_mflix_changes").db_changes.find(
{},
{ _id: 0, clusterTime: 1, ns: 1, operationType: 1, fullDocument: 1 }
)

Cada evento incluye una marca de tiempo, el cambio completo del documento, un operationType y un campo ns. El campo ns nombra la base de datos y la colección de origen, para que pueda distinguir qué colección produjo cada cambio.

La siguiente agregación ingiere datos de la fuente cluster-changestream-connection, que se conecta a un clúster de Atlas cargado con el conjunto de datos Sample Mflix Dataset Collections, y escribe en el destino cluster-changestream-sink-connection, que se conecta a un clúster de Atlas independiente. Para aprender a crear un espacio de trabajo de procesamiento de flujos y agregar una conexión a un clúster de Atlas al registro de conexiones, consulte Introducción al procesamiento de flujos de Atlas. Esta agregación ejecuta dos etapas para abrir un flujo de cambios en todo el clúster en el clúster de origen y registrar los cambios en la colección events en la base de datos cluster_changes en el clúster de destino. Esta etapa $source omite db y coll, por lo que informa los cambios de todas las bases de datos y colecciones en el clúster de origen, en lugar de una sola colección o base de datos.

  1. La etapa $source se conecta a la fuente cluster-changestream-connection y abre un flujo de cambios para todo el clúster de origen. El campo config.startAtOperationTime especifica que el procesador de flujo debe comenzar a informar los cambios que ocurran en o después de una hora específica.

  2. La etapa $merge escribe los documentos del flujo de cambios en una colección de Atlas llamada events en la base de datos cluster_changes del clúster de destino. Si no existe dicha base de datos o colección, Atlas las crea.

[{
"$source": {
"connectionName": "cluster-changestream-connection",
"config": {
"startAtOperationTime": {"$date": "2024-08-19T18:00:00.000Z"}
}
}
},
{
"$merge": {
"into": {
"connectionName": "cluster-changestream-sink-connection",
"db": "cluster_changes",
"coll": "events"
}
}
}]

El siguiente comando mongosh inserta un documento en la colección movies de la base de datos sample_mflix en el clúster de origen:

db.getSiblingDB("sample_mflix").movies.insertOne(
{ _id: ObjectId("66c1a1f1f1f1f1f1f1f1f1f1"), title: "Example Movie" }
)

Una vez insertado el documento, el procesador de flujo guarda el documento de evento del flujo de cambios en la colección cluster_changes.events en el clúster de destino.

Para ver los documentos de la colección cluster_changes.events resultante, ejecute los siguientes comandos en su clúster de destino. Para conectarse, consulte Conectarse a un clúster mediante mongosh.

db.getSiblingDB("cluster_changes").events.find()

El campo ns en el documento de salida muestra que el cambio se origina en la inserción sample_mflix.movies en el clúster de origen. La canalización de agregación de flujo de cambios refleja este cambio en la colección cluster_changes.events en el clúster de destino.

La siguiente agregación utiliza un arreglo de documentos en línea como fuente de datos de transmisión, que contiene observaciones meteorológicas para tres ubicaciones. El arreglo utiliza el mismo esquema que el Conjunto de Datos de muestra de Meteorología. Esta agregación ejecuta tres etapas:

  1. La etapa $source define un arreglo documents en línea de observaciones meteorológicas como fuente de datos en streaming y utiliza timeField para designar el campo timestamp de cada documento como la marca temporal autorizada.

  2. La etapa $match remite solo los documentos con dewPoint.value mayor que 5.0 a la siguiente etapa.

  3. La fase $merge escribe la salida en una colección Atlas denominada stream en la base de datos sample_weatherstream. Si no existe tal base de datos o colección, Atlas los creará.

[{
"$source": {
"documents": [
{
"location": "New York",
"timestamp": ISODate('2024-01-15T08:00:00Z'),
"temp": 23.5,
"dewPoint": { "value": 6.2 }
},
{
"location": "Los Angeles",
"timestamp": ISODate('2024-01-15T08:05:00Z'),
"temp": 18.2,
"dewPoint": { "value": 4.8 }
},
{
"location": "Chicago",
"timestamp": ISODate('2024-01-15T08:10:00Z'),
"temp": 26.8,
"dewPoint": { "value": 7.5 }
}
]
}
},
{
"$match": { "dewPoint.value": { "$gt": 5.0 } }
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "stream"
}
}
}]

Para ver los documentos de la colección sample_weatherstream.stream resultante, conéctate a un clúster de Atlas y ejecuta el siguiente comando:

db.getSiblingDB("sample_weatherstream").stream.find()

La siguiente agregación filtra informes de la fuente de datos de transmisión de muestra sample_stream_solar. Archiva los resultados en la conexión solar-cluster-connection. Para aprender cómo crear un espacio de trabajo de Stream Processing y agregar una conexión a su espacio de trabajo de Atlas Stream Processing, consulte Introducción a Atlas Stream Processing y Agregar una conexión de Atlas Stream Processing. Esta agregación ejecuta tres etapas para filtrar informes de la fuente sample_stream_solar y guardar los resultados en una colección denominada solarColl en la base de datos solarDb:

  1. La etapa $source se conecta a la fuente sample_stream_solar. El campo timeField convierte el campo timestamp de cada informe entrante a una fecha usando $dateFromString.

  2. La etapa $match excluye los documentos donde device_id es device_8 y pasa los informes de todos los demás dispositivos a la siguiente etapa.

  3. La etapa $merge escribe la salida en la colección solarColl en la base de datos solarDb del clúster solar-cluster-connection.

[{
"$source": {
"connectionName": "sample_stream_solar",
"timeField": {
"$dateFromString": { "dateString": "$timestamp" }
}
}
},
{
"$match": { "device_id": { "$ne": "device_8" } }
},
{
"$merge": {
"into": {
"connectionName": "solar-cluster-connection",
"db": "solarDb",
"coll": "solarColl"
}
}
}]

Nota

La fuente sample_stream_solar es una conexión solo para pruebas que genera documentos de muestra cada segundo y está destinada a la creación rápida de prototipos.

Ejecute los siguientes comandos en su clúster solar-cluster-connection para verificar el comportamiento del procesador de transmisiones. Para conectarse, consulte Conectarse a un clúster mediante mongosh.

Para ver los documentos de la colección solarDb.solarColl resultante, ejecuta:

db.getSiblingDB("solarDb").solarColl.find()

Cada documento incluye los campos device_id, event_type, group_id, max_watts, obs y timestamp de un informe de dispositivo solar. Puede observar que ningún documento en la salida tiene un device_id de device_8, porque la etapa $match excluye esos informes antes de que la etapa $merge escriba los documentos restantes en solarColl.

Nota

Lo anterior es un ejemplo representativo. Los datos en transmisión no son estáticos, y cada usuario ve documentos distintos.

La siguiente agregación genera un documento cada cinco minutos en lugar de leerlo desde una conexión. Esta agregación se ejecuta en tres etapas:

  1. La etapa $project etiqueta cada documento con un nombre de trabajo y copia la marca de tiempo programada del campo _ts en un campo runAt.

  2. La fase $merge escribe la salida en una colección Atlas denominada heartbeats en la base de datos sample_weatherstream. Si no existe tal base de datos o colección, Atlas los creará.

[{
"$source": {
"schedule": "0 0/5 * * * *"
}
},
{
"$project": {
"job": "five-minute-heartbeat",
"runAt": "$_ts"
}
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "heartbeats"
}
}
}]

Para ver los documentos de la colección sample_weatherstream.heartbeats resultante, conéctate a un clúster de Atlas y ejecuta el siguiente comando:

db.getSiblingDB("sample_weatherstream").heartbeats.find()

Nota

Lo anterior es un ejemplo representativo. Las marcas de tiempo que ve dependen de cuándo inicia el procesador de flujos.