Obtenha dados de streaming no lakehouse e acesse com o ponto de extremidade de análise SQL

Este quickstart explica como criar uma definição de trabalho Spark que contenha código Python com Spark Structured Streaming para transportar dados numa casa de lago e depois servi-los através de um endpoint de análise SQL. Depois de concluir este quickstart, terá uma definição de trabalho no Spark que corre continuamente e o endpoint de análise SQL poderá visualizar os dados recebidos.

Criar um script Python

Utilize o seguinte script Python para criar uma tabela Delta de streaming num armazém de dados em nuvem utilizando o Apache Spark. O script lê um fluxo de dados gerados (uma linha por segundo) e grava-o no modo de acréscimo em uma tabela Delta chamada streamingtable. Ele armazena os dados e as informações do ponto de verificação na casa do lago especificada.

  1. Use o seguinte código Python que usa o streaming estruturado do Spark para obter dados em uma tabela lakehouse.

    from pyspark.sql import SparkSession
    
    if __name__ == "__main__":
     # Start Spark session
     spark = SparkSession.builder \
         .appName("RateStreamToDelta") \
         .getOrCreate()
    
     # Table name used for logging
     tableName = "streamingtable"
    
     # Define Delta Lake storage path
     deltaTablePath = f"Tables/{tableName}"
    
     # Create a streaming DataFrame using the rate source
     df = spark.readStream \
         .format("rate") \
         .option("rowsPerSecond", 1) \
         .load()
    
     # Write the streaming data to Delta
     query = df.writeStream \
         .format("delta") \
         .outputMode("append") \
         .option("path", deltaTablePath) \
         .option("checkpointLocation", f"{deltaTablePath}/_checkpoint") \
         .start()
    
     # Keep the stream running
     query.awaitTermination()
    
  2. Salve seu script como arquivo Python (.py) em seu computador local.

Criar uma casa no lago

Use as seguintes etapas para criar uma casa no lago:

  1. Faça login no Portal Fabric.

  2. Navegue até o espaço de trabalho desejado ou crie um novo, se necessário.

  3. Para criar uma Lakehouse, selecione Novo item no espaço de trabalho e, em seguida, selecione Lakehouse no painel que se abre.

    Captura de ecrã mostrando a nova caixa de diálogo lakehouse.

  4. Digite o nome da sua casa do lago e selecione Criar.

Criar uma definição de trabalho do Spark

Use os seguintes passos para criar uma definição de trabalho Spark:

  1. No mesmo espaço de trabalho em que você criou uma casa no lago, selecione Novo item.

  2. No painel que se abre, em Obter Dados, selecione Definição de Tarefa do Spark.

  3. Introduza o nome da definição do seu trabalho no Spark e selecione Criar.

  4. Selecione Upload e selecione o arquivo Python que você criou na etapa anterior.

  5. Em Lakehouse Reference escolhe a lakehouse que criaste.

Definir a política de Retentativas para a definição de funções Spark

Use as seguintes etapas para definir a política de repetição para sua definição de trabalho do Spark:

  1. No menu superior, selecione o ícone Configuração .

    Captura de tela mostrando o ícone de configurações do Spark Job Definition.

  2. Abra o separador Otimização e defina o gatilho Política de Retentativa como Ativado.

    Captura de ecrã mostrando o separador de otimização de Definição de Trabalho do Spark.

  3. Defina o máximo de tentativas ou marque Permitir tentativas ilimitadas.

  4. Especifique o tempo entre cada tentativa de repetição e selecione Aplicar.

Nota

Há um limite vitalício de 90 dias para a configuração da política de novas tentativas. Quando a política de repetição estiver ativada, o trabalho será reiniciado de acordo com a política dentro de 90 dias. Após esse período, a política de repetição deixará automaticamente de funcionar e o trabalho será encerrado. Os usuários precisarão reiniciar manualmente o trabalho, o que, por sua vez, reativará a política de repetição.

Executar e monitorizar a definição do trabalho Spark

  1. No menu do topo, selecione o ícone Executar.

    Captura de ecrã a mostrar o ícone de execução da Definição de Trabalho do Spark.

  2. Verifique se a definição do Spark Job foi enviada com êxito e está em execução.

Visualizar dados utilizando um endpoint de análise SQL

Após a execução do script, uma tabela chamada streamingtable com as colunas timestamp e valor é criada no lakehouse. Você pode visualizar os dados usando o endpoint de análise SQL.

  1. A partir do espaço de trabalho, abra a sua casa do lago.

  2. Mude para o ponto de extremidade de análise SQL no canto superior direito.

  3. No painel de navegação esquerdo, expanda Schemas > dbo >Tables, selecione streamingtable para visualizar os dados.