11 Streaming

La sezione descrive l'uso dei dati in streaming o dei dati prodotti in modo continuo in Oracle AI Data Platform Workbench.

Informazioni sullo streaming

Puoi elaborare i dati in streaming o i dati prodotti in modo continuo quasi in tempo reale in Oracle AI Data Platform Workbench utilizzando la funzionalità di streaming strutturato di Apache Spark.

Sia i notebook che i flussi di lavoro supportano lo streaming strutturato di Apache Spark. È possibile utilizzare le seguenti origini e sink per leggere i dati dei flussi, scrivere i dati dei flussi in e per le posizioni dei checkpoint.

Tabella 11-1 Sorgenti e lavelli supportati

Origine o lavandino Supportato?
Percorso del volume (/Volume/bronze/bucket1) Supportato per tutti i formati
Percorso area di lavoro (/Workspace/folder1/) Supportato per tutti i formati
Tabelle nei cataloghi con tre nomi di parte (catalog.schema.table) Supportato solo per il formato Delta

Non supportato per i formati Parquet, CSV, JSON e ORC

Esempio 1: codice supportato

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

Esempio 2: codice non supportato

  • spark.readStream.option("withEventTimeOrder", "true").format("format") .table("stdcatalog.stdschema.samplecsv")
Kafka Supportato per qualsiasi flusso compatibile con Kafka senza convenzione di denominazione in tre parti

Non supportato per il catalogo basato su Kafka in base alla convenzione di denominazione in tre parti

Servizio di streaming OCI Supportato
Percorso di storage degli oggetti OCI (utilizzando OCI://) Non supportata
Oracle Autonomous AI Lakehouse, Oracle AI Database, Oracle Autonomous AI Transaction Processing Non supportato per lo streaming (readStream o writeStream)

Streaming strutturato tramite notebook

È possibile scrivere codice Python per elaborare i dati di flusso in un notebook. I percorsi del volume o dell'area di lavoro sono validi come posizione di checkpoint, ma i percorsi di storage degli oggetti (formato oci://) non sono supportati come posizione di checkpoint. Si consiglia di utilizzare i percorsi di volume come posizione di checkpoint.


Esempio di codice di streaming in una cella notebook del workbench di AI Data Platform


Esempio di codice Python utilizzato per elaborare i dati di flusso in un notebook AI Data Platform Workbench

È possibile visualizzare gli eventi correlati allo streaming di Apache Spark, ad esempio la velocità di input, la velocità di elaborazione e la durata del batch dalla scheda Dashboard nel notebook durante l'esecuzione del codice di streaming.


Scheda Dashboard in un notebook aperto per visualizzare i dati di streaming

È inoltre possibile visualizzare gli eventi relativi allo streaming raw dalla scheda Dati raw mentre si sviluppa il codice in modo incrementale.


Scheda Dati non elaborati aperta in un notebook che visualizza eventi correlati allo streaming

Configurazione dello streaming strutturato Spark mediante i flussi di lavoro

È possibile configurare un task di streaming all'interno di un flusso di lavoro per l'elaborazione continua dei dati di flusso.

Prima è necessario creare un job, quindi aggiungere un blocco note o un task Python a tale job per iniziare a utilizzare i flussi di lavoro con lo streaming in Oracle AI Data Platform Workbench.
  1. Passare all'area di lavoro e fare clic su Flusso di lavoro.
  2. Fare clic su Icona Crea clusterCrea job.
  3. Fornire nome e descrizione per il job.
  4. Fare clic su Sfoglia e selezionare la posizione in cui salvare il job nel workbench di AI Data Platform. Fare clic su Seleziona.
  5. Immettere 1 per Numero massimo di esecuzioni concorrenti.
  6. Fare clic su Crea.
  7. Fare clic sul job appena creato.
  8. Fare clic su Aggiungi task.
  9. Fornire un nome per il task.
  10. Selezionare Notebook o Python per Tipo di task.
  11. Fare clic su Sfoglia e passare allo script Notebook o Python che si desidera aggiungere come task di streaming. Fare clic su Seleziona.
  12. Selezionare un cluster di computazione per il task Notebook o Python, se non è già collegato.
  13. Selezionare la casella di controllo Streaming. Se si seleziona Streaming, il timeout di esecuzione e le dipendenze dei task vengono disabilitati come opzioni.

    Pagina Crea dettagli task aperta con la casella di controllo Streaming selezionata

  14. Selezionare il numero di nuovi tentativi che un task deve tentare in caso di errore. Se si selezionano più di 0, è necessario specificare anche il tempo di attesa dell'esecuzione del job tra i nuovi tentativi e se i nuovi tentativi devono essere eseguiti al timeout.

    Opzioni di nuovo tentativo task quando il numero di nuovi tentativi è pari o superiore a 1

  15. Fare clic su Esegui ora.
Dopo l'avvio di un task di streaming, l'esecuzione continua finché non viene interrotta manualmente. Durante la normale manutenzione mensile, l'attività di streaming viene arrestata e riavviata dal servizio senza richiedere alcuna azione da parte tua.