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

Operaciones de procesamiento de flujos de Atlas

Esta guía describe las mejores prácticas para la gestión del procesamiento de flujos de Atlas, incluidos los procedimientos específicos para los flujos de trabajo esenciales.

Cada procesador de Stream Processing depende de conexiones a fuentes de datos y sumideros. Atlas Stream Processing cifra todos los datos en tránsito mediante TLS/SSL, pero las conexiones básicas siguen pasando datos a través de Internet público. Para lograr un rendimiento óptimo y seguridad con Atlas Stream Processing, tenga en cuenta las prácticas descritas en esta sección.

Para garantizar una comunicación segura entre Atlas Stream Processing y sistemas externos privados, como clústeres Apache Kafka autogestionados o API privadas, utilice conexiones de interconexión VPC o VNet. Estas conexiones protegen sus datos de la exposición a internet. Tenga en cuenta que, además de establecer la conexión de interconexión, debe configurar explícitamente las tablas de enrutamiento de la VPC de su aplicación para dirigir el tráfico al bloque CIDR de su VPC de Atlas. Para obtener más información sobre cómo configurar la tabla de enrutamiento, consulte la documentación de VPC de su proveedor externo.

Para enrutar datos de transmisión entre diferentes proyectos de Atlas dentro de la misma organización, Atlas Stream Processing admite Conexiones entre proyectos. Las conexiones entre proyectos transmiten datos sin exponer los clústeres a la Internet privada, lo que garantiza la privacidad sin la necesidad de configurar manualmente el emparejamiento de VPC.

Atlas Stream Processing admite las siguientes integraciones de AWS:

  • S3

  • Flujos de datos de Kinesis

  • Lambda (como $externalFunction)

Cada una de estas integraciones es compatible con el Acceso Unificado a AWS. Configurar el Acceso Unificado elimina la necesidad de almacenar identificadores de acceso estáticos o credenciales secretas en MongoDB, y es el modelo de autenticación recomendado para todas las integraciones de AWS.

Atlas Stream Processing admite la autenticación con brokers de Apache Kafka mediante OIDC con tokens web JSON. Sin embargo, esta función actualmente solo es compatible con proveedores de identidad accesibles a través de redes públicas. Esta función no es compatible con proveedores de identidad dentro de la VPC del cliente.

Atlas Stream Processing proporciona diversos mecanismos para garantizar un servicio fiable y robusto en escenarios de fallos comunes, incluidos puntosde control, colas de mensajes no entregados y procesadores de conmutación por error.

Atlas Stream Processing supervisa el estado del procesador de transmisión mediante un latido de 30segundos. Si un procesador de transmisión no envía un latido durante más de 30 segundos, Atlas Stream Processing reinicia automáticamente el procesador.

Atlas Stream Processing gestiona automáticamente los cambios de topología para integraciones fragmentadas como AWS Kinesis.

Puedes configurar Atlas Stream Processing para que dirija los datos no procesados ​​a una cola de mensajes fallidos (DLQ). A continuación, puedes ver los registros no procesados ​​en la DLQ para solucionar los fallos de procesamiento.

Al examinar su DLQ, busque los siguientes errores de procesamiento comunes:

  • Análisis de carga útil

    El registro entrante está mal formado y no se puede analizar. Esto suele deberse a un JSON no válido.

  • Evaluación de expresiones

    Una expresión en el proceso de análisis no se evalúa correctamente.

    Ejemplo

    Su pipeline incluye el campo de configuración $emit.config.tombstoneWhen para borrar registros de Kafka, y un documento dado no tiene el campo a evaluar.

  • Límite de tamaño

    El registro saliente excede el límite de tamaño de su destino, como 1 MB para AWS Kinesis.

  • Generación de claves

    El procesador encuentra un error al generar una clave de partición para su destino.

  • Datos tardíos

    Su procesador utiliza Windows y el documento llegó después de que expirara el período de retraso permitido configurado.

Tras analizar las causas de los fallos de procesamiento, modifique sus fuentes de datos anteriores para corregir los datos antes de que entren en su canalización de Atlas Stream Processing, o bien modifique su canalización para gestionar los casos excepcionales.

Considere las siguientes opciones de remediación para los registros en la DLQ después de utilizarlos para guiar las correcciones a su configuración de Atlas Stream Processing:

  • Descartar

    Si un registro está muy mal formado o es irrelevante—como en el caso de datos muy atrasados en sistemas muy sensibles al tiempo—bórrelo de la colección.

  • Corrección manual

    Para un registro con errores JSON menores o uno que sea de vital importancia, ajuste manualmente los campos e insértelos directamente en su destino.

  • Reprocesamiento automatizado:

    Para grandes volúmenes de datos recuperables y relevantes, cree un script secundario o un procesador de flujo adicional que lea desde la cola de mensajes no entregados (DLQ) y aplique una transformación en función de las condiciones de fallo conocidas antes de enrutar los datos limpios al destino.

Atlas Stream Processing proporciona procesadores de conmutación por error para evitar la interrupción del servicio en caso de fallo del sistema en toda la región.

Los siguientes principios definen una configuración de conmutación por error sólida para proteger sus cargas de trabajo de Atlas Stream Processing.

  • Implemente clústeres de Atlas en varias regiones.

    Para garantizar el funcionamiento continuo de sus procesadores de flujo en caso de una conmutación por error regional, cualquier clúster de Atlas al que se conecte su procesador de flujo debe ser un clúster multirregión. Este principio garantiza que sus conexiones de origen o destino permanezcan activas cuando la región principal no esté disponible.

  • Cree espacios de trabajo de Atlas Stream Processing para cada región.

    Tras implementar los clústeres multirregión, cree un espacio de trabajo para cada región. Este principio minimiza la latencia de las conexiones entre los procesadores de flujo y las bases de datos Atlas que los soportan.

  • Rellenar el Registro de conexiones con conexiones de copia de seguridad.

    Configure las conexiones para cada origen y destino en el espacio de trabajo de cada región. Asegúrese de configurar las conexiones no solo para las conexiones básicas, sino también para las conexiones de VPC y de punto final privado.

  • Configurar la recuperación ante desastres específica de la conexión

    Para garantizar la continuidad del servicio para un procesador de flujos, cualquier proveedor externo del que dependa también debe tener habilitada la conmutación por error. Asegúrese de que cada origen o destino de proveedor externo esté configurado correctamente según su documentación.

    Ejemplo

    Si su procesador de flujo utiliza un clúster de Kafka como origen o destino, configure el clúster para que conserve el desplazamiento del consumidor, lo que permitirá a Atlas Stream Processing reanudar el procesamiento desde un desplazamiento de grupo de consumidores.

  • Realizar simulacros de rutina.

    Verifique su configuración de conmutación por error realizando simulacros periódicamente. Estos simulacros sirven para asegurar que su configuración se conecte a las redes y cargue los puntos de control correctamente, y para entrenar el flujo de trabajo de conmutación por error, facilitando así su ejecución cuando sea necesario. Considere establecer una periodicidad regular para los simulacros de conmutación por error y realice uno cada vez que modifique su configuración.

Una vez que se asegure de que las mejores prácticas están implementadas, cree un procesador de transmisión y verifique su configuración de conmutación por error.

Atlas Stream Processing admite la conmutación por error regional automatizada únicamente para procesadores con las siguientes configuraciones:

  • Origen y destino de Atlas

  • Fuente de Atlas y sumidero de Apache Kafka

Para obtener más información sobre la conmutación por error automatizada, consulte Procesador de conmutación por error.

Sin embargo, puede configurar una arquitectura de Atlas Stream Processing que admita la conmutación por error manual para cualquier combinación de conexión de origen y destino.

Puede iniciar una conmutación por error forzada a nivel de espacio de trabajo o de procesador de flujo.

Para iniciar la conmutación por error del espacio de trabajo de Stream Processing, siga el procedimiento descrito en Activador de conmutación por error del espacio de trabajo de Stream Processing.

Para iniciar la conmutación por error de un procesador de flujo individual, siga el procedimiento descrito en Iniciar la conmutación por error de un procesador de flujo.

En caso de una interrupción del servicio regional:

1

Cuando observe síntomas de una interrupción del servicio, revise su informe de estado preferido para verificar el evento. Puede consultar esta información en la página de estado de MongoDB Cloud o en la página de estado de su proveedor de nube. Puede configurar alertas de procesamiento de flujos de Atlas. Anote la marca de tiempo del evento de fallo.

2

Siga el procedimiento descrito en Iniciar conmutación por error para un procesador de transmisión. Para asegurarse de que el servicio se reanude desde el último evento de transmisión consumido, pase la marca de tiempo que anotó anteriormente como parámetro startAtOperationTime.

Los procesadores que utilizan fuentes de Apache Kafka reanudan la ejecución según los desplazamientos de los grupos de consumidores.

3

Cuando la región principal vuelva a estar en línea, los procesadores principales se reiniciarán automáticamente, lo que provocará un procesamiento duplicado. Cree y ejecute un script para comprobar continuamente el estado de la región principal y detener todos los procesadores de dicha región cuando vuelva a estar en línea.

4

Tras completar la conmutación por error, compruebe periódicamente el estado de la región principal de su espacio de trabajo. Puede consultar esta información en la página de estado de MongoDB Cloud o en la página de estado de su proveedor de nube. También puede configurar las alertas de procesamiento de flujos de Atlas. Si utiliza la supresión automática del procesador principal, puede configurar el script para que envíe notificaciones cuando detecte que los procesadores principales se reinician automáticamente.

5

Detenga los procesadores de flujo de conmutación por error. Anote la marca de tiempo de cada detención.

6

Reinicie los procesadores de transmisión en su región principal. Para asegurarse de que el servicio se reanude desde el último evento de transmisión consumido, pase la marca de tiempo que anotó anteriormente como parámetro startAtOperationTime.

Los procesadores que utilizan fuentes de Apache Kafka reanudan la ejecución según los desplazamientos de los grupos de consumidores.

Tenga en cuenta lo siguiente al configurar sus sistemas para la conmutación por error regional o al activar la conmutación por error manual:

  • Cuando se produce una conmutación por error regional, Atlas Stream Processing promueve los procesadores de conmutación por error inactivos en su región de conmutación por error a procesadores activos. Estos procesadores reanudan el Stream Processing desde el último punto de controldisponible.

  • Se garantiza que los puntos de control de Atlas Stream Processing no tendrán más de 10 minutos de antigüedad. Sin embargo, los procesadores de conmutación por error reprocesan cualquier dato que el procesador principal haya procesado desde el último punto de control al inicio. La conmutación por error garantiza el procesamiento al menos una vez.

  • Si un procesador de conmutación por error no se inicia durante un evento de conmutación por error regional, seguirá asumiendo el rol de procesador activo. En ese caso, puede solucionar la causa subyacente y reiniciar el procesador con normalidad en la región de conmutación por error.

Atlas Stream Processing ofrece diversas herramientas de observabilidad para evaluar el rendimiento y el estado de sus espacios de trabajo, conexiones y procesadores. Además, para consultar el estado del servicio regional y las interrupciones, puede visitar la página de estado de MongoDB Cloud.

Atlas Stream Processing aprovisiona recursos por cada procesador de transmisión. Para obtener más información, consulta la Guía de selección de nivel.

Para los procesadores de Atlas Stream Processing que operan en clústeres de Atlas, asegúrese de escalar tanto sus procesadores como sus clústeres proporcionalmente a sus necesidades de E/S de canalización. En particular, evalúe los procesadores que realizan un alto volumen de operaciones $lookup o $merge como candidatos para la ampliación del clúster.

Atlas Stream Processing puede usar un clúster de Atlas $source para consumir eventos de flujo de cambios. Para garantizar un rendimiento y una fiabilidad óptimos para dichos procesadores de flujo, revise las siguientes consideraciones:

  • Tamaño de la ventana de Oplog

    Los procesadores de flujo utilizan el oplog del clúster $source para reanudar el servicio después de una interrupción. Si un procesador permanece inactivo durante un período mayor que el período de retención del oplog, el token de reanudación puede apuntar a datos truncados, lo que provoca un error ChangeStreamHistoryLost.

    Mitigue este riesgo ajustando el tamaño del oplog a un tiempo de inactividad aceptable para sus sistemas de procesamiento de flujos.

  • Opciones de reanudación

    Atlas Stream Processing admite reanudar el flujo de cambios $source``s from either the last stored checkpoint with ``resumeFromCheckpoint o un momento específico con startAtOperationTime. Para evitar lagunas o duplicaciones en los datos procesados, defina una política estandarizada sobre qué opciones de reanudación utiliza en un escenario determinado.

  • Ingestión de backlog

    Los procesadores de flujo que operan en un namespace grande pueden experimentar largos períodos iniciales de procesamiento de puesta al día. Considere reducir el alcance del namespace de su procesador de flujo o elegir un punto de inicio más reciente para reducir este período inicial de puesta al día.

  • Imágenes de transmisión de cambio

    Las auditorías, las transacciones compensatorias y las operaciones de deshacer requieren acceso al estado previo al procesamiento del documento. Al trabajar con Atlas como origen, habilite las preimágenes y las postimágenes de documentos para admitir estas operaciones. Al implementar la compatibilidad con preimágenes o postimágenes de documentos, verifique la cobertura y la estabilidad en todo el sistema.

  • La etapa ``$replaceRoot``

    Las fuentes de flujo de cambios devuelven documentos que agrupan metadatos operativos con el contenido del documento afectado. Escribir estos documentos en un receptor requiere que los consumidores trabajen con el esquema de eventos de cambio en lugar de la lógica empresarial del documento de origen. A menos que necesite los metadatos, utilice la etapa $replaceRoot para promover el subdocumento de contenido de origen antes de un procesamiento posterior.

  • Pipeline orientado a uniones

    Atlas Stream Processing admite la etapa $lookup para permitir el enriquecimiento de datos en tiempo real mediante operaciones de unión. Estas operaciones pueden aumentar la latencia de procesamiento. Considere limitar sus canalizaciones de Atlas Stream Processing únicamente a las operaciones $lookup que deben realizarse antes de escribir en un destino, reservando los enriquecimientos más complejos para vistas precalculadas.

  • Puntos de control obsoletos

    Los puntos de control contienen solo el historial de eventos, no la lógica de procesamiento. Cuando reanuda la operación desde un punto de control, su procesador de transmisión aplica la lógica de su pipeline actual a los eventos antiguos, lo que puede provocar una falta de coincidencia entre su estado descendente existente y su nueva salida. Considere la posibilidad de introducir un nuevo procesador de transmisión en lugar de reanudar desde un punto de control si necesita una semántica de procesamiento significativamente diferente.

  • Planificación de reinicio

    Durante el reinicio, un procesador de flujo podría experimentar un rendimiento reducido de eventos nuevos. Para aumentar la resiliencia del sistema, programe reinicios para tener en cuenta la capacidad temporalmente reducida.

  • Planificación de conmutación por error

    Desarrollar procedimientos estándar para el uso compartido de puntos de control del procesador de transmisión en todas las regiones, la elección de puntos de reanudación seguros y la validación de la coherencia de la salida para garantizar que los eventos de conmutación por error regionales no provoquen problemas de calidad de los datos.