Consideraciones de producción para Structured Streaming

Ejecute cargas de trabajo de Structured Streaming de producción como trabajos programados de Lakeflow en Azure Databricks. Consulte Trabajos de Lakeflow.

Databricks recomienda configurar siempre lo siguiente:

  • Quite el código innecesario de los cuadernos que devolverían resultados, como display y count.
  • No ejecute cargas de trabajo de Structured Streaming con cómputo de uso general. Programe siempre los flujos como Jobs de Lakeflow usando el cómputo de Jobs.
  • Programar trabajos de Lakeflow mediante el Continuousmodo. Esto hace referencia a la característica de programación de trabajos de Azure Databricks, no al intervalo structured Streaming trigger.
  • No habilite el escalado automático para el proceso para trabajos de Structured Streaming.

Algunas cargas de trabajo se benefician de lo siguiente:

Databricks introdujo las canalizaciones de Lakeflow para reducir la complejidad de gestionar la infraestructura de producción para las cargas de trabajo de Structured Streaming. Databricks recomienda utilizar canalizaciones de Lakeflow para las nuevas canalizaciones de Structured Streaming. Consulte Spark Declarative Pipelines.

Nota

El escalado automático de proceso tiene limitaciones al reducir verticalmente el tamaño del clúster para cargas de trabajo de Structured Streaming. Databricks recomienda usar canalizaciones declarativas de Spark en Lakeflow con escalado automático mejorado para cargas de trabajo de streaming. Consulte Optimizar la utilización del clúster de canalización de Lakeflow mediante escalado automático.

:::note Computación sin servidor

En la computación sin servidor, solo Trigger.AvailableNow() y Trigger.Once() se admiten. Databricks recomienda Trigger.AvailableNow().

Para la transmisión continua en informática sin servidor, use desencadenado frente al modo de canalización continua en modo continuo.

Consulte Limitaciones de streaming.

:::

Reducir la latencia para la transmisión operativa

Las cargas de trabajo de streaming operativas ingieren, transforman y actúan sobre los datos casi en tiempo real. Ejemplos comunes incluyen la detección de fraude, la detección de anomalías, la personalización y la monitorización y alertas en tiempo real, donde el procesamiento retrasado afecta directamente a los resultados del negocio. La baja latencia para estas cargas de trabajo suele significar decenas o cientos de milisegundos, aunque muchos equipos establecen acuerdos de nivel de servicio (SLA) en el rango de segundos para tener en cuenta la variabilidad en percentiles altos.

Para la latencia de extremo a extremo más baja, se utiliza el modo en tiempo real, que logra una latencia de extremo a extremo inferior a un segundo en la cola y alrededor de 300 milisegundos en casos comunes. Véase conceptos de modo en tiempo real.

Cuando el modo en tiempo real no se adapta a tu carga de trabajo, las siguientes mejores prácticas reducen la latencia para el streaming estructurado micro-batch:

  • Modo de salida: Usa el modo de actualización donde tus operadores de consulta y el sink lo soporten. El modo de actualización emite filas actualizadas tras cada activación y sigue actualizándolas hasta que expira la marca de agua, así que haz que tu sumidero de destino sea idempotente para poder manejar las actualizaciones. Utiliza el modo de anexado para las cargas de trabajo que el modo de actualización no admite, como las uniones entre flujos, o cuando puedes descartar los datos que llegan tarde. No uses el modo completo para baja latencia. Consulte Seleccionar un modo de salida para Structured Streaming.
  • Disparador: Usa un processingTime disparador con intervalo 0 , que inicia el siguiente microlote tan pronto como termina el anterior y hay nuevos datos disponibles. Esto ofrece la menor latencia de microlote, pero aumenta los costes de las API de almacenamiento en la nube. No use AvailableNow, Once ni Continuous para cargas de trabajo operativas. Consulte Configurar intervalos del desencadenador de Structured Streaming.
  • Watermark: Configura el watermark con una duración suficiente para incluir los datos que llegan con retraso y que tu carga de trabajo no debe descartar. La marca de agua controla durante cuánto tiempo la consulta acepta datos con tiempo de evento fuera de orden antes de descartarlos y eliminar el estado, por lo que una marca de agua demasiado corta descarta silenciosamente registros tardíos válidos. Dentro de esa limitación, una marca de agua más corta reduce la latencia y mantiene menos estado, mientras que una marca de agua más larga tolera más datos que llegan con retraso, a costa de una mayor latencia y de más estado. Un múltiplo pequeño de tu SLA de latencia, como 2x, es un punto de partida razonable para el ajuste. Consulte Aplicar marcas de agua para controlar los umbrales de procesamiento de datos.
  • Fuentes y sumideros: Lee fuentes de baja latencia como buses de mensajes (Apache Kafka, Amazon Kinesis, Apache Pulsar o Google Cloud Pub/Sub) o cambia fuentes de datos de las tablas Delta Lake e Apache Iceberg. Escriba en destinos de baja latencia y alto rendimiento, como sistemas de mensajería, bases de datos operativas o destinos foreach. Diseña las operaciones de sumidero para que sean idempotentes, de modo que los consumidores aguas abajo gestionen duplicados y datos que llegan tarde.
  • Estado y puntos de control: Para consultas con estado, utiliza el almacén de estados RocksDB, que es necesario tanto para la creación de puntos de control del registro de cambios como para la creación asíncrona de puntos de control del estado. Habilita la creación de puntos de control del registro de cambios para persistir únicamente los cambios incrementales del estado. Cuando el punto de control de estado suponga el cuello de botella en la duración de tu lote, activa el punto de control de estado asíncrono para solapar las escrituras de puntos de control con el siguiente microlote, después de revisar las advertencias sobre la recuperación ante fallos y el cambio de tamaño del clúster. Asigna a cada consulta su propio directorio de puntos de control en almacenamiento en la nube duradero. Consulta Configurar el almacén de estado de RocksDB en Azure Databricks, Puntos de control de estado asíncronos para consultas con estado y Puntos de control de streaming estructurado.
  • Gestión de offsets: Para reducir la latencia causada por la creación de puntos de control de offsets en flujos continuos, habilite el seguimiento asíncrono del progreso, que actualiza los registros de offsets y de confirmación sin bloquear el procesamiento de datos. No es compatible con los desencadenadores AvailableNow o Once. Consulte Seguimiento de progreso asincrónico.
  • Saltos de almacenamiento: Mantener el cálculo dentro de una única tubería de streaming siempre que sea posible. Dividir la lógica entre varios trabajos o pipelines añade saltos de almacenamiento que aumentan la latencia.

Diseño de cargas de trabajo de streaming para esperar un error

Databricks recomienda configurar siempre los trabajos de streaming para reiniciarse automáticamente en caso de error. Algunas funcionalidades, incluida la evolución del esquema, requieren que las cargas de trabajo de Structured Streaming vuelvan a intentarlo automáticamente. Consulte Configuración de trabajos de Structured Streaming para que reinicien las consultas de streaming en caso de error.

Algunas operaciones como foreachBatch proporcionan al menos una vez en lugar de garantías exactamente una vez. Para estas operaciones, asegúrese de que la canalización de procesamiento sea idempotente. Consulte Uso de foreachBatch para escribir en receptores de datos arbitrarios.

Nota

Cuando se reinicia una consulta, se procesa el microlote planificado durante la ejecución anterior. Si su trabajo falló debido a un error de falta de memoria o canceló manualmente un trabajo debido a un microlote sobredimensionado, es posible que necesite escalar verticalmente el proceso informático para procesar correctamente el microlote.

Si cambia las configuraciones entre ejecuciones, estas configuraciones se aplican al primer lote nuevo planeado. Consulte Recuperación después de los cambios en una consulta de Structured Streaming.

Cuando se reintenta una tarea

Puede programar varias tareas como parte de un trabajo de Azure Databricks. Al configurar un trabajo mediante el desencadenador continuo, no se pueden establecer dependencias entre tareas.

Puede optar por programar varias secuencias en un solo trabajo mediante uno de los siguientes métodos:

  • Varias tareas: defina un trabajo con varias tareas que ejecutan cargas de trabajo de streaming mediante el desencadenador continuo.
  • Varias consultas: defina varias consultas de streaming en el código fuente para una sola tarea.

También puede combinar estas estrategias. En la siguiente tabla se comparan estos enfoques.

Estrategia Varias tareas Varias consultas
¿Cómo se comparte el proceso? Databricks recomienda implementar procesos de trabajo del tamaño adecuado para cada tarea de streaming. Opcionalmente, puede compartir el proceso entre tareas. Todas las consultas comparten el mismo proceso. Puede opcionalmente asignar consultas a grupos de planificación.
¿Cómo se controlan los reintentos? Todas las tareas deben producir un error antes de que el trabajo se reintente. La tarea vuelve a intentarlo si se produce un error en alguna consulta.

Para más información sobre cómo trabajar con varias tareas o consultas, consulte Ejecución de varias consultas de Structured Streaming en el mismo clúster.

Configurar trabajos de Structured Streaming para reiniciar las consultas de streaming en caso de error

Databricks recomienda configurar todas las cargas de trabajo de streaming mediante el desencadenador continuo. Consulte Ejecución de trabajos continuamente.

El desencadenador continuo tiene el siguiente comportamiento de forma predeterminada:

  • Impide que se ejecute más de una ejecución simultánea del trabajo.
  • Inicia una nueva ejecución cuando se produce un error en una ejecución anterior.
  • Usa retroceso exponencial para reintentos.

Databricks recomienda usar siempre el proceso de trabajos en lugar del proceso multiuso al programar flujos de trabajo. En caso de error y reintento del trabajo, se implementan nuevos recursos de proceso.

Nota

Databricks recomienda no usar streamingQuery.awaitTermination() ni spark.streams.awaitAnyTermination(). Consulte Cuándo usar awaitTermination().

Cuándo usar awaitTermination()

streamingQuery.awaitTermination() y spark.streams.awaitAnyTermination() bloquean el hilo actual hasta que finalice una consulta en streaming. Si se usan estas funciones depende del entorno de ejecución.

En el caso de los trabajos de Lakeflow, no use streamingQuery.awaitTermination() ni spark.streams.awaitAnyTermination(). Estas funciones no son necesarias porque el servicio Trabajos impide que una ejecución se complete automáticamente cuando una consulta de streaming esté activa. Ambas funciones impiden que las celdas del cuaderno finalicen e impidan que el servicio Trabajos realice el seguimiento de la consulta de streaming, lo que interrumpe las métricas de trabajos pendientes y las notificaciones de trabajo.

Use awaitTermination() en los casos siguientes:

Caso de uso Comportamiento
Cuadernos interactivos en cómputo de propósito general awaitTermination() mantiene la celda en ejecución, permite observar el estado de la consulta y garantiza que los errores se muestran en la salida del cuaderno.
Entornos locales y de desarrollo Cuando se ejecuta un programa spark localmente, el proceso se cierra cuando se completa el subproceso principal. Llame awaitTermination() para mantener el programa activo hasta que finalice o falle la consulta de transmisión.
Propagación de errores al controlador Sin awaitTermination(), es posible que un error de consulta de streaming en un contexto que no sea de trabajo no se propague al subproceso que realiza la llamada. La consulta puede producir errores de forma silenciosa, lo que dificulta la detección y el diagnóstico de errores. Llamar a awaitTermination() vuelve a generar la excepción de consulta en el controlador.