这个快速入门介绍了如何创建一个包含 Python 代码和 Spark 结构化流的 Spark 作业定义,将数据落在湖屋中,然后通过 SQL 分析端点提供数据。 完成这个快速入门后,你会有一个持续运行的Spark作业定义,SQL分析端点可以查看接收的数据。
创建 Python 脚本
使用以下 Python 脚本在 Lakehouse 中使用 Apache Spark 创建流式增量表。 该脚本读取生成的数据流(每秒一行),并将其以追加模式写入名为 streamingtableDelta 表。 它将数据和检查点信息存储在指定的湖仓中。
使用以下使用 Spark 结构化流的 Python 代码获取湖屋表中的数据。
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()在本地计算机中将脚本另存为 Python 文件 (.py)。
创建湖屋
使用以下步骤创建湖屋:
登录到 Fabric 门户。
导航到所需的工作区,或根据需要创建一个新工作区。
若要创建“Lakehouse”,请从工作区中选择“新建项”,然后在打开的面板中选择“Lakehouse”。
输入湖屋的名称,然后选择“创建”。
创建 Spark 作业定义
请使用以下步骤创建 Spark 作业定义:
在创建湖屋的同一工作区中,选择“新建项”。
在打开的面板中,在“获取数据”下,选择 Spark 作业定义。
输入你的 Spark 作业定义名称,然后选择 创建。
选择“上传”并选择你在上一步中创建的 Python 文件。
在“湖屋引用”下,选择你创建的湖屋。
为 Spark 作业定义设置重试策略
使用以下步骤为 Spark 作业定义设置重试策略:
注意
重试策略设置的生存期限制为 90 天。 启用重试策略后,作业将在 90 天内根据策略重启。 在此期限之后,重试策略将自动停止运行,作业将被终止。 然后,用户需要手动重启作业,而这又会重新启用重试策略。
执行并监控Spark作业定义
使用 SQL 分析终结点查看数据
脚本运行后,将在 Lakehouse 中创建一个名为 流式处理表 的表,其中包含 时间戳 和 值 列。 可以使用 SQL 分析终结点查看数据:
从工作区打开你的湖边别墅。
从右上角切换到 SQL 分析终结点 。
在左侧导航窗格中,展开 架构 > dbo >表,选择 流式处理表 以预览数据。