11 Streaming

A seção aborda o uso de streaming de dados ou dados produzidos continuamente no Oracle AI Data Platform Workbench.

Sobre o Streaming

Você pode processar dados de streaming ou dados produzidos continuamente quase em tempo real no Oracle AI Data Platform Workbench usando o recurso Apache Spark Structured Streaming.

Tanto notebooks quanto workflows suportam streaming estruturado do Apache Spark. Você pode usar as seguintes origens e sumidouros para ler dados de stream, gravar dados de stream e para locais de checkpoint.

Tabela 11-1 Origens e Pias Suportadas

Origem ou Dissipador Suportado?
Caminho do volume (/Volume/bronze/bucket1) Suportado para todos os formatos
Caminho do espaço de trabalho (/Workspace/folder1/) Suportado para todos os formatos
Tabelas em catálogos com três nomes de partes (catalog.schema.table) Suportado somente para formato Delta

Não suportado para formatos Parquet, CSV, JSON, ORC

Exemplo 1: Código suportado

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

Exemplo 2: código não suportado

  • spark.readStream.option("withEventTimeOrder", "true").format("format") .table("stdcatalog.stdschema.samplecsv")
Kafka Suportado para qualquer fluxo compatível com Kafka sem convenção de nomeação em três partes

Não suportado para catálogo baseado no Kafka após convenção de nomeação em três partes)

Serviço OCI Streaming Suportado
Caminho do OCI Object Storage (usando OCI://) Não Suportado
Oracle Autonomous AI Lakehouse, Oracle AI Database, Oracle Autonomous AI Transaction Processing Não suportado para streaming (readStream ou writeStream)

Streaming Estruturado Usando Notebooks

Você pode gravar código Python para processar dados de fluxo em um notebook. Os caminhos de volume ou de espaço de trabalho são válidos como um local de checkpoint, mas os caminhos de Armazenamento de objetos (formato oci://) não são suportados como um local de checkpoint. Recomendamos o uso de caminhos de volume como um local de checkpoint.


Exemplo de código de streaming em uma célula de notebook do AI Data Platform Workbench


Exemplo de código Python usado para processar dados de fluxo em um notebook AI Data Platform Workbench

Você pode ver eventos relacionados ao streaming do Apache Spark, como taxa de entrada, taxa de processamento e duração do batch na guia Painel de Controle do seu notebook ao executar o código de streaming.


Guia Painel de controle em um notebook aberto para exibir dados de streaming

Você também pode exibir os eventos brutos relacionados ao streaming na guia Dados Brutos enquanto desenvolve seu código de forma incremental.


Guia Dados Brutos aberta em um notebook exibindo eventos relacionados ao streaming

Configurando o Streaming Estruturado do Spark usando Workflows

Você pode configurar uma tarefa de streaming dentro de um fluxo de trabalho para processamento contínuo de dados de fluxo.

Primeiro você precisa criar um job e depois adicionar um Notebook ou uma tarefa Python a esse job para começar a usar workflows com streaming no Oracle AI Data Platform Workbench.
  1. Navegue até o seu espaço de trabalho e clique em Workflow.
  2. Clique em Ícone Criar clusterCriar Job.
  3. Forneça um nome e descrição para o seu job.
  4. Clique em Procurar e selecione o local para salvar o job no Workbench da Plataforma de Dados do AI. Clique em Selecionar.
  5. Informe 1 para Máximo de Execuções Concorrentes.
  6. Clique em Criar.
  7. Clique no job que você acabou de criar.
  8. Clique em Adicionar tarefa.
  9. Forneça um nome para a sua tarefa.
  10. Selecione Notebook ou Python para Tipo de tarefa.
  11. Clique em Procurar e navegue até o script Notebook ou Python que você deseja adicionar como tarefa de Streaming. Clique em Selecionar.
  12. Selecione um cluster de computação para a tarefa Notebook ou Python, se ainda não houver um anexado.
  13. Marque a caixa de seleção Streaming. A seleção de Streaming desativa o timeout de execução e as dependências de tarefa como opções.

    Página Criar Detalhes da Tarefa aberta com a caixa de seleção Streaming marcada

  14. Selecione o número de novas tentativas que uma tarefa deve tentar em caso de falha. Se você selecionar mais de 0, também deverá especificar quanto tempo a execução do job deverá aguardar entre as novas tentativas e se as novas tentativas deverão ser feitas no tempo limite.

    Opções de repetição de tarefa quando o número de repetições for 1 ou maior

  15. Clique em Executar Agora.
Depois que uma tarefa do Streaming é iniciada, ela continua a ser executada até que você a interrompa manualmente. Durante a manutenção mensal regular, a tarefa de Streaming é interrompida e reiniciada pelo serviço sem exigir nenhuma ação da sua extremidade.