11 Transmission en continu

Cette section traite de l'utilisation de données transmises en continu ou de données produites en continu dans Oracle AI Data Platform Workbench.

À propos de Streaming

Vous pouvez traiter des données de transmission en continu ou des données produites en continu en temps quasi réel dans Oracle AI Data Platform Workbench à l'aide de la fonctionnalité Apache Spark Structured Streaming.

Les blocs-notes et les workflows prennent en charge la diffusion en continu structurée Apache Spark. Vous pouvez utiliser les sources et les puits suivants pour lire les données de flux, écrire les données de flux vers et pour les emplacements de point de reprise.

Tableau 11-1 Sources et puits pris en charge

Source ou récepteur Pris en charge ?
Chemin du volume (/Volume/bronze/bucket1) Prise en charge pour tous les formats
Chemin de l'espace de travail (/Workspace/folder1/) Prise en charge pour tous les formats
Tables dans les catalogues avec trois noms de parties (catalog.schema.table) Pris en charge uniquement pour le format Delta

Non pris en charge pour les formats Parquet, CSV, JSON et ORC

Exemple 1 : code pris en charge

  • 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")

Exemple 2 : code non pris en charge

  • spark.readStream.option("withEventTimeOrder", "true").format("format") .table("stdcatalog.stdschema.samplecsv")
Kafka Pris en charge pour tous les flux compatibles avec Kafka sans convention de dénomination en trois parties

Non pris en charge pour le catalogue basé sur Kafka suivant une convention de dénomination en trois parties)

Service OCI Streaming Pris en charge
Chemin du stockage OCI Object (à l'aide d'oci ://) Non pris en charge
Oracle Autonomous AI Lakehouse, Oracle AI Database, Oracle Autonomous AI Transaction Processing Non pris en charge pour la diffusion en continu (readStream ou writeStream)

Streaming structuré à l'aide de blocs-notes

Vous pouvez écrire du code Python pour traiter les données de flux dans un bloc-notes. Les chemins de volume ou d'espace de travail sont valides en tant qu'emplacement de point de reprise, mais les chemins Object Storage (formatoci ://) ne sont pas pris en charge en tant qu'emplacement de point de reprise. Nous vous recommandons d'utiliser les chemins de volume comme emplacement de point de reprise.


Exemple de code de transmission en continu dans une cellule de bloc-notes AI Data Platform Workbench


Exemple de code Python utilisé pour traiter les données de flux dans un bloc-notes AI Data Platform Workbench

Vous pouvez voir les événements liés à la transmission en continu d'Apache Spark, tels que le taux d'entrée, le taux de traitement et la durée de batch, dans l'onglet Tableau de bord de votre bloc-notes lors de l'exécution du code de transmission en continu.


Onglet Tableau de bord d'un bloc-notes ouvert pour afficher les données de transmission en continu

Vous pouvez également afficher les événements bruts liés à la transmission en continu à partir de l'onglet Données brutes tout en développant votre code de manière incrémentielle.


Onglet Données brutes ouvert dans un bloc-notes affichant les événements liés à la diffusion en continu

Configurer Spark Structured Streaming à l'aide de workflows

Vous pouvez configurer une tâche de transmission en continu dans un workflow pour un traitement continu des données de flux.

Vous devez d'abord créer un travail, puis ajouter une tâche Bloc-notes ou Python à ce travail pour commencer à utiliser des workflows avec la transmission en continu dans Oracle AI Data Platform Workbench.
  1. Accédez à votre espace de travail et cliquez sur Workflow.
  2. Cliquez sur Icône Créer un clusterCréer un travail.
  3. Indiquez le nom et la description de votre travail.
  4. Cliquez sur Parcourir et sélectionnez l'emplacement où enregistrer le travail dans AI Data Platform Workbench. Cliquez sur Sélectionner.
  5. Entrez 1 pour Nombre maximal d'exécutions simultanées.
  6. Cliquez sur Créer.
  7. Cliquez sur le travail que vous venez de créer.
  8. Cliquez sur Ajouter une tâche.
  9. Indiquez le nom de votre tâche.
  10. Sélectionnez Bloc-notes ou Python pour le type de tâche.
  11. Cliquez sur Parcourir et accédez au bloc-notes ou au script Python à ajouter en tant que tâche de transmission en continu. Cliquez sur Sélectionner.
  12. Sélectionnez un cluster de calcul pour la tâche Bloc-notes ou Python, si elle n'est pas déjà attachée.
  13. Cochez la case Streaming. La sélection de Streaming désactive le délai d'expiration de l'exécution et les dépendances de tâche en tant qu'options.

    Page Créer des détails de tâche ouverte avec la case Diffusion en continu cochée

  14. Sélectionnez le nombre de tentatives qu'une tâche doit effectuer en cas d'échec. Si vous sélectionnez plus de 0, vous devez également indiquer la durée d'attente de l'exécution du travail entre les nouvelles tentatives et si les nouvelles tentatives doivent être tentées en cas d'expiration.

    Options de nouvelle tentative de tâche lorsque le nombre de nouvelles tentatives est supérieur ou égal à 1

  15. Cliquez sur Maintenant.
Une fois qu'une tâche Streaming est démarrée, elle continue de s'exécuter jusqu'à ce que vous l'arrêtiez manuellement. Au cours d'une maintenance mensuelle régulière, la tâche Streaming est arrêtée et redémarrée par le service sans aucune action de votre part.