11 Streaming

In diesem Abschnitt wird die Verwendung von Streamingdaten oder kontinuierlich erzeugten Daten in Oracle AI Data Platform Workbench behandelt.

Informationen zu Streaming

Mit der Apache Spark Structured Streaming-Funktion können Sie Streamingdaten oder kontinuierlich erzeugte Daten nahezu in Echtzeit in Oracle AI Data Platform Workbench verarbeiten.

Sowohl Notizbücher als auch Workflows unterstützen strukturiertes Streaming von Apache Spark. Sie können die folgenden Quellen und Senken zum Lesen von Streamdaten, Schreiben von Streamdaten in und für Checkpoint-Speicherorte verwenden.

Tabelle 11-1: Unterstützte Quellen und Sinks

Quelle oder Sink Unterstützt?
Volume-Pfad (/Volume/bronze/bucket1) Unterstützt für alle Formate
Workspace-Pfad (/Workspace/folder1/) Unterstützt für alle Formate
Tabellen in Katalogen mit drei Teilenamen (catalog.schema.table) Nur für Delta-Format unterstützt

Nicht unterstützt für Parquet-, CSV-, JSON-, ORC-Formate

Beispiel 1: Unterstützter Code

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

Beispiel 2: Nicht unterstützter Code

  • spark.readStream.option("withEventTimeOrder", "true").format("format") .table("stdcatalog.stdschema.samplecsv")
Kafka Wird für alle Kafka-kompatiblen Streams ohne dreiteilige Benennungskonvention unterstützt

Nicht unterstützt für Kafka-basierten Katalog nach dreiteiligen Benennungskonventionen)

OCI Streaming-Service Unterstützt
OCI-Objektspeicherpfad (mit OCI://) Nicht unterstützt
Oracle Autonomous AI Lakehouse, Oracle AI Database, Oracle Autonomous AI Transaction Processing Nicht unterstützt für Streaming (readStream oder writeStream)

Strukturiertes Streaming mit Notizbüchern

Sie können Python-Code schreiben, um Streamdaten in einem Notizbuch zu verarbeiten. Volume-Pfade oder Workspace-Pfade sind als Checkpoint-Speicherort gültig, Object Storage-Pfade (oci://-Format) werden jedoch nicht als Checkpoint-Speicherort unterstützt. Wir empfehlen, Volume-Pfade als Checkpoint-Speicherort zu verwenden.


Beispiel für das Streaming von Code in einer Notizbuchzelle der AI Data Platform Workbench


Beispiel für Python-Code zur Verarbeitung von Streamdaten in einem AI Data Platform Workbench-Notizbuch

Sie können Apache Spark-Streamingereignisse wie Eingaberate, Verarbeitungsrate und Batchdauer auf der Registerkarte Dashboard in Ihrem Notizbuch anzeigen, während Sie Streamingcode ausführen.


Dashboard-Registerkarte in einem geöffneten Notizbuch zur Anzeige von Streamingdaten

Sie können die Raw Streaming-bezogenen Ereignisse auch auf der Registerkarte Raw-Daten anzeigen, während Sie Ihren Code inkrementell entwickeln.


Registerkarte "Rohdaten" in einem Notizbuch geöffnet, in dem Streaming-bezogene Ereignisse angezeigt werden

Strukturiertes Spark-Streaming mit Workflows konfigurieren

Sie können eine Streamingaufgabe in einem Workflow für die kontinuierliche Verarbeitung von Streamdaten konfigurieren.

Sie müssen zunächst einen Job erstellen und dann diesem Job eine Notizbuch- oder Python-Aufgabe hinzufügen, um Workflows mit Streaming in Oracle AI Data Platform Workbench zu verwenden.
  1. Navigieren Sie zu Ihrem Workspace, und klicken Sie auf Workflow.
  2. Klicken Sie auf Cluster erstellen (Symbol)Job erstellen.
  3. Geben Sie einen Namen und die Beschreibung für Ihren Job an.
  4. Klicken Sie auf Durchsuchen, und wählen Sie den Speicherort aus, in dem der Job in der AI Data Platform Workbench gespeichert werden soll. Klicken Sie auf Auswählen.
  5. Geben Sie 1 für Max. gleichzeitige Ausführungen ein.
  6. Klicken Sie auf Create.
  7. Klicken Sie auf den gerade erstellten Job.
  8. Klicken Sie auf Aufgabe hinzufügen.
  9. Geben Sie einen Namen für Ihre Aufgabe an.
  10. Wählen Sie Notizbuch oder Python als Aufgabentyp aus.
  11. Klicken Sie auf Durchsuchen, und navigieren Sie zu dem Notizbuch- oder Python-Skript, das Sie als Streamingaufgabe hinzufügen möchten. Klicken Sie auf Auswählen.
  12. Wählen Sie ein Compute-Cluster für die Notizbuch- oder Python-Aufgabe aus, wenn noch kein Compute-Cluster angehängt ist.
  13. Aktivieren Sie das Kontrollkästchen Streaming. Wenn Sie Streaming auswählen, werden Ausführungstimeout und Aufgabenabhängigkeiten als Optionen deaktiviert.

    Seite "Aufgabendetails erstellen" mit aktiviertem Kontrollkästchen "Streaming" geöffnet

  14. Wählen Sie die Anzahl der Wiederholungsversuche, die eine Aufgabe bei einem Fehler versuchen soll. Wenn Sie mehr als 0 auswählen, müssen Sie auch angeben, wie lange der Joblauf zwischen Wiederholungen warten soll und ob Wiederholungen bei Timeout versucht werden sollen.

    Optionen für Aufgabenwiederholung, wenn die Anzahl der Wiederholungen 1 oder höher ist

  15. Klicken Sie auf Jetzt ausführen.
Nachdem eine Streamingaufgabe gestartet wurde, wird sie weiter ausgeführt, bis Sie sie manuell stoppen. Während der regelmäßigen monatlichen Wartung wird die Streamingaufgabe vom Service gestoppt und neu gestartet, ohne dass eine Aktion von Ihrem Ende aus erforderlich ist.