Obtenga datos de streaming en el almacén de lago y acceda con el punto de conexión de análisis SQL

Este inicio rápido explica cómo crear una definición de trabajo en Spark que contenga código en Python con Spark Structured Streaming para aterrizar datos en una casa de lago y luego servirlos a través de un endpoint de análisis SQL. Tras completar este inicio rápido, tendrás una definición de trabajo en Spark que se ejecuta continuamente y el endpoint de analítica SQL podrá ver los datos que llegan.

Creación de un script de Python

Utiliza el siguiente script de Python para crear una tabla Delta de transmisión en un lakehouse mediante Apache Spark. El script lee una secuencia de datos generados (una fila por segundo) y la escribe en modo de anexión en una tabla Delta denominada streamingtable. Almacena los datos y la información del punto de control en el lago especificado.

  1. Use el siguiente código de Python que usa el streaming estructurado de Spark para obtener datos en una tabla de 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. Guarde el script como archivo de Python (.py) en el equipo local.

Creación de un almacén de lago

Use los siguientes pasos para crear un cliente:

  1. Inicie sesión en el portal de Fabric.

  2. Vaya al área de trabajo deseada o cree una nueva si es necesario.

  3. Para crear un Lakehouse, seleccione Nuevo elemento en el área de trabajo y, a continuación, seleccione Lakehouse en el panel que se abre.

    Captura de pantalla que muestra el diálogo de la nueva instancia de Lakehouse.

  4. Escriba el nombre de su instancia de Lakehouse y seleccione Crear.

Creación de una definición de trabajo de Spark

Utiliza los siguientes pasos para crear una definición de trabajo en Spark:

  1. En el mismo espacio de trabajo donde creó un lakehouse, seleccione Nuevo elemento.

  2. En el panel que se abre, en Obtener datos, seleccione Definición de trabajo de Spark.

  3. Introduce el nombre de la definición de tu puesto en Spark y selecciona Crear.

  4. Seleccione Cargar y seleccione el archivo Python que creó en el paso anterior.

  5. En Referencia de Lakehouse , elija la instancia de Lakehouse que ha creado.

Establecer política de reintentos para la definición de tareas de Spark

Siga estos pasos para establecer la directiva de reintento para la definición del trabajo de Spark:

  1. En el menú superior, seleccione el icono Configuración .

    Captura de pantalla que muestra el icono de configuración de definición de trabajo de Spark.

  2. Abra la pestaña Optimización y establezca el desencadenador de la política de reintentos en Activado.

    Captura de pantalla que muestra la pestaña Optimización de definición de trabajos de Spark.

  3. Defina el número máximo de reintentos o active Permitir intentos ilimitados.

  4. Especifique el tiempo entre cada intento de reintento y seleccione Aplicar.

Nota:

Hay un límite de duración de 90 días para la configuración de la directiva de reintento. Una vez habilitada la directiva de reintento, el trabajo se reiniciará según la directiva en un plazo de 90 días. Después de este período, la directiva de reintento dejará de funcionar automáticamente y el trabajo se finalizará. A continuación, los usuarios deberán reiniciar manualmente el trabajo, lo que, a su vez, volverá a habilitar la directiva de reintento.

Ejecutar y monitorizar la definición de tareas de Spark

  1. En el menú superior, seleccione el icono Ejecutar.

    Captura de pantalla que muestra el icono de ejecución de la Definición de Trabajo de Spark.

  2. Compruebe si la definición del trabajo de Spark se envió correctamente y se ejecutó.

Visualización de datos mediante un punto de conexión de análisis SQL

Una vez que se ejecuta el script, se crea una tabla denominada streamingtable con marca de tiempo y columnas de valor en lakehouse. Puede ver los datos mediante el endpoint de SQL Analytics:

  1. Desde el espacio de trabajo, abre tu casa del lago.

  2. Cambie al punto de conexión de SQL Analytics desde la esquina superior derecha.

  3. En el panel de navegación izquierdo, expanda Esquemas > dbo >Tablas, seleccione streamingtable para obtener una vista previa de los datos.