11 Flujo

En esta sección se trata el uso de datos de transmisión o datos producidos continuamente en Oracle AI Data Platform Workbench.

Acerca de Streaming

Puede procesar datos de transmisión o datos producidos continuamente en casi tiempo real en Oracle AI Data Platform Workbench mediante la capacidad de transmisión estructurada de Apache Spark.

Tanto los blocs de notas como los flujos de trabajo admiten la transmisión estructurada de Apache Spark. Puede utilizar los siguientes orígenes y receptores para leer datos de flujos, escribir datos de flujos y para ubicaciones de puntos de control.

Tabla 11-1 Orígenes y fregaderos admitidos

Origen o receptor ¿Soportado?
Ruta de volumen (/Volume/bronze/bucket1) Soportado para todos los formatos
Ruta de espacio de trabajo (/Workspace/folder1/) Soportado para todos los formatos
Tablas en catálogos con tres nombres de parte (catalog.schema.table) Soportado solo para formato Delta

No soportado para formatos Parquet, CSV, JSON y ORC

Ejemplo 1: código soportado

  • streaming_df = spark.readStream.format("delta").table('stdcatalog.stdschema.deltatable')
  • streaming_df.writeStream.format("delta").outputMode("append").option("checkpointLocation", "/Volumes/checkpoints1/").toTable("stdcatalog.stdschema.deltatable")

Ejemplo 2: código no soportado

  • spark.readStream.option("withEventTimeOrder", "true").format("format") .table("stdcatalog.stdschema.samplecsv")
Kafka Soportado para cualquier flujo compatible con Kafka sin convención de nomenclatura de tres partes

No soportado para el catálogo basado en Kafka tras la convención de nomenclatura de tres partes)

OCI Streaming Soportados
Ruta de OCI Object Storage (mediante OCI://) No soportada
Oracle Autonomous AI Lakehouse, Oracle AI Database, Oracle Autonomous AI Transaction Processing No soportado para transmisión (readStream o writeStream)

Flujo estructurado mediante blocs de notas

Puede escribir código Python para procesar datos de flujo en un bloc de notas. Las rutas de acceso de volumen o de espacio de trabajo son válidas como ubicación de punto de control, pero las rutas de Object Storage (formatooci://) no están soportadas como ubicación de punto de control. Recomendamos utilizar rutas de volumen como ubicación de punto de control.


Ejemplo de código de transmisión en una celda de bloc de notas de AI Data Platform Workbench


Ejemplo de código Python utilizado para procesar datos de flujo en un bloc de notas de AI Data Platform Workbench

Puede ver los eventos relacionados con el flujo de Apache Spark, como el ratio de entrada, el ratio de procesamiento y la duración del lote, en el separador Panel de control de su bloc de notas mientras ejecuta el código de flujo.


Separador Dashboard (Panel de control) en un bloc de notas abierto para mostrar datos de transmisión

También puede ver los eventos relacionados con el flujo sin procesar desde el separador Datos sin procesar mientras desarrolla el código de forma incremental.


Separador Raw Data abierto en un bloc de notas que muestra eventos relacionados con la transmisión

Configuración de flujos estructurados de Spark mediante flujos de trabajo

Puede configurar una tarea de flujo dentro de un flujo de trabajo para el procesamiento continuo de datos de flujo.

Primero debe crear un trabajo y, a continuación, agregar una tarea de Notebook o Python a ese trabajo para empezar a utilizar flujos de trabajo con transmisión en Oracle AI Data Platform Workbench.
  1. Vaya al espacio de trabajo y haga clic en Flujo de trabajo.
  2. Haga clic en Icono Crear clusterCrear trabajo.
  3. Proporcione un nombre y la descripción para su trabajo.
  4. Haga clic en Examinar y seleccione la ubicación para guardar el trabajo en el área de trabajo de AI Data Platform. Haga clic en Seleccionar.
  5. Introduzca 1 para Máximo de ejecuciones simultáneas.
  6. Haga clic en Create.
  7. Haga clic en el trabajo que acaba de crear.
  8. Haga clic en Agregar tarea.
  9. Proporcione un nombre para la tarea.
  10. Seleccione Notebook o Python para Tipo de tarea.
  11. Haga clic en Examinar y navegue hasta el bloc de notas o script de Python que desea agregar como tarea de flujo. Haga clic en Seleccionar.
  12. Seleccione un cluster de recursos informáticos para la tarea de bloc de notas o Python, si aún no hay uno asociado.
  13. Active la casilla de control Flujo. Al seleccionar Streaming, se desactivan el timeout de ejecución y las dependencias de tareas como opciones.

    Se abre la página Create Task details con la casilla de control Streaming seleccionada

  14. Seleccione el número de reintentos que una tarea debe intentar en caso de fallo. Si selecciona más de 0, también debe especificar cuánto tiempo debe esperar la ejecución del trabajo entre reintentos y si se deben intentar reintentos con timeout.

    Opciones de reintento de tarea cuando el número de reintentos es 1 o mayor

  15. Haga clic en Ejecutar Ahora.
Después de iniciar una tarea de Streaming, se sigue ejecutando hasta que la pare manualmente. Durante el mantenimiento mensual regular, el servicio detiene y reinicia la tarea Streaming sin necesidad de realizar ninguna acción por su parte.