十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

awesome-copilot 技能库实战:Microsoft Fabric Lakehouse 的 PySpark 开发指南

awesome-copilot 技能库实战:Microsoft Fabric Lakehouse 的 PySpark 开发指南 awesome-copilot 技能库实战Microsoft Fabric Lakehouse 的 PySpark 开发指南【免费下载链接】awesome-copilotCommunity-contributed instructions, agents, skills, and configurations to help you make the most of GitHub Copilot.项目地址: https://gitcode.com/GitHub_Trending/aw/awesome-copilot本文以 awesome-copilot 开源仓库中 fabric-lakehouse 技能 自带的 PySpark 参考文档 为主体系统讲解在 Microsoft Fabric Lakehouse 中使用 PySpark 完成数据读取、Delta 表写入与 CRUD、Schema 定义、SQL Magic、表优化V-Order / OPTIMIZE / VACUUM以及增量加载与 SCD Type 2 等关键开发模式。读者学完后可以在 Fabric 数据工程与 AI 场景中直接复刻这些可运行片段并理解每个配置项与命令背后的存储与计算原理。技能定位这份 PySpark 参考在仓库中的角色在 awesome-copilot 中fabric-lakehouse 是一个面向“设计、构建、优化 Lakehouse 解决方案”的 Agent 技能SKILL.md。它由三部分组成SKILL.md技能主指令介绍 Lakehouse 核心概念Delta Tables、Files、SQL Endpoint、Shortcuts、Materialized Views、Spark Views、安全模型工作区角色与 OneLake 数据访问、性能优化与数据沿袭Lineage。PySpark 代码参考本文的主体覆盖从会话配置到增量加载的完整代码模式。数据接入参考说明如何通过 Data Factory 管道把外部数据搬进 Lakehouse。使用该技能的场景包括生成包含 Fabric Lakehouse 能力说明的文档、按最佳实践设计/构建/优化 Lakehouse 方案、理解 Lakehouse 中表数据与非表数据的管理方式。本指南即围绕其中的 PySpark 部分展开并结合其余两处文档做纵深补充。Spark 会话配置最佳实践在 Notebook 中正式处理数据前建议先对 Spark 会话启用 Fabric 的两项关键优化pyspark.md# Enable Fabric optimizations spark.conf.set(spark.sql.parquet.vorder.enabled, true) spark.conf.set(spark.microsoft.delta.optimizeWrite.enabled, true)两个配置项的作用如下配置键作用适用场景spark.sql.parquet.vorder.enabled启用 V-Order 写入优化对 Parquet/Delta 数据在写出时进行排序与聚类需要语义模型 / Power BI 读取加速的表spark.microsoft.delta.optimizeWrite.enabled启用 Delta 写入时的文件自动合并减少小文件数量频繁以小批量 append 的流式或微批写入其中 V-Order 是 Fabric 在引擎层提供的物理存储优化如 SKILL.md 所述它会在写入时预先对数据进行排序使数据按常见访问模式排列从而提升查询性能。由于这两个开关影响的是写入端行为建议在 Notebook 起始单元格统一设置保证后续所有saveAsTable/write操作都受益。读取数据从 Bronze 层文件到 SQL 端点Lakehouse 将表数据存放在Tables目录、非表文件存放在Files目录见 SKILL.md。在 PySpark 中Files/前缀路径即指向 Lakehouse 的Files目录常见做法是把原始数据按 bronze/silver/gold 分层存放。# Read CSV file df spark.read.format(csv) \ .option(header, true) \ .option(inferSchema, true) \ .load(Files/bronze/data.csv) # Read JSON file df spark.read.format(json).load(Files/bronze/data.json) # Read Parquet file df spark.read.format(parquet).load(Files/bronze/data.parquet) # Read Delta table df spark.read.table(my_delta_table) # Read from SQL endpoint df spark.sql(SELECT * FROM lakehouse.my_table)CSVheadertrue将首行作为列名inferSchematrue让 Spark 自动推断列类型。若数据量较大或类型推断不稳定建议改用显式 Schema见下文“使用 StructType 定义显式 Schema”。JSON / ParquetSpark 原生支持自描述格式无需指定 Schema 即可读取。Delta 表spark.read.table(my_delta_table)直接按逻辑表名读取 Lakehouse 中的托管 Delta 表。SQL 端点Lakehouse 自动生成只读 SQL 分析端点SKILL.md在 PySpark 中可通过spark.sql以lakehouse.table的库表形式访问。写入 Delta 表Lakehouse 中表的主要格式是 Delta其他格式如 CSV、Parquet 仅可用于 Spark 查询见 SKILL.md。写入 Delta 表的标准方式是df.write.format(delta).saveAsTable(...)# Write DataFrame as managed Delta table df.write.format(delta) \ .mode(overwrite) \ .saveAsTable(silver_customers) # Write with partitioning df.write.format(delta) \ .mode(overwrite) \ .partitionBy(year, month) \ .saveAsTable(silver_transactions) # Append to existing table df.write.format(delta) \ .mode(append) \ .saveAsTable(silver_events)关键点说明mode(overwrite)整表覆盖写入适合全量重算的场景如每日全量快照。partitionBy(year, month)按时间列分区数据在存储层按分区目录组织。后续查询若带分区过滤条件如WHERE year 2024可大幅减少扫描数据量。mode(append)增量追加适合日志、事件流等只增数据配合spark.microsoft.delta.optimizeWrite.enabledtrue可自动控制小文件数量。Delta 表 CRUD 操作Delta 格式提供 ACID 事务、版本化与时间旅行能力SKILL.md因此可以直接对表执行 UPDATE / DELETE / MERGE# UPDATE spark.sql( UPDATE silver_customers SET status active WHERE last_login 2024-01-01 -- Example date, adjust as needed ) # DELETE spark.sql( DELETE FROM silver_customers WHERE is_deleted true ) # MERGE (Upsert) spark.sql( MERGE INTO silver_customers AS target USING staging_customers AS source ON target.customer_id source.customer_id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * )UPDATE按条件更新满足条件的行示例中last_login 2024-01-01为占位条件实际需按业务口径调整。DELETE按条件物理逻辑删除行Delta 通过新的数据版本实现配合后续 VACUUM 才真正释放存储。MERGEUpsertON子句指定匹配键WHEN MATCHED THEN UPDATE SET *表示匹配时用源行更新目标全部字段WHEN NOT MATCHED THEN INSERT *表示未匹配时插入源行。这是数据仓库中增量更新、拉链表构建的核心语法。使用 StructType 定义显式 Schema当 CSV 等无 Schema 数据的列类型推断不可靠或需要强制约束如金额列必须为 Decimal、ID 列不允许为空时可显式定义 Schemafrom pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType, DecimalType schema StructType([ StructField(id, IntegerType(), False), StructField(name, StringType(), True), StructField(email, StringType(), True), StructField(amount, DecimalType(18, 2), True), StructField(created_at, TimestampType(), True) ]) df spark.read.format(csv) \ .schema(schema) \ .option(header, true) \ .load(Files/bronze/customers.csv)StructField(字段名, 数据类型, nullable)中的第三个参数为是否允许空值False表示非空约束读取时违反约束的行会报错起到质量把关作用。金额类字段使用DecimalType(18, 2)表示 18 位精度、2 位小数避免浮点误差是财务场景的标准做法。时间字段显式声明为TimestampType可避免inferSchema对日期时间格式的误判。在 Notebook 中使用 SQL MagicFabric Notebook 支持%%sqlMagic可直接对 Delta 表执行 SQL适合快速探查与聚合分析%%sql -- Query Delta table directly SELECT customer_id, COUNT(*) as order_count, SUM(amount) as total_amount FROM gold_orders GROUP BY customer_id ORDER BY total_amount DESC LIMIT 10该模式适合验证链路产物gold 层表、临时分析、生成报表数据集。注意%%sql单元格的结果会以 DataFrame 形式返回可继续在 Python 单元格中使用复杂的清洗逻辑仍建议用 PySpark API 实现SQL 专注聚合与探查。V-Order 优化与表维护V-Order面向语义模型读取的排序优化# Enable V-Order for read optimization spark.conf.set(spark.sql.parquet.vorder.enabled, true)如 SKILL.md 所述V-Order 会对 Delta 表数据做预排序以改善常见访问模式的查询性能。该优化对 Fabric 语义模型供 Power BI 使用的读取加速尤为关键因此凡是要接语义模型的表都应保持 V-Order 开启。OPTIMIZE / VACUUM小文件合并与存储清理随着数据持续摄入与更新表会积累大量小文件、历史版本需要定期维护SKILL.md%%sql -- Optimize table (compact small files) OPTIMIZE silver_transactions -- Optimize with Z-ordering on query columns OPTIMIZE silver_transactions ZORDER BY (customer_id, transaction_date) -- Vacuum old files (default 7 days retention) VACUUM silver_transactions -- Vacuum with custom retention VACUUM silver_transactions RETAIN 168 HOURSOPTIMIZE把小文件合并成大文件减少读取时的文件打开开销数据摄入和更新越频繁越需要定期执行。ZORDER BY (customer_id, transaction_date)在合并的同时按指定列做 Z 排序让相邻值在存储上聚簇可显著加速这些列上的过滤与连接。VACUUM清理过期文件、释放存储。默认保留 7 天即保留最近 7 天内的所有版本RETAIN 168 HOURS可自定义保留窗口168 小时 7 天。执行 VACUUM 会移除超出保留期的历史版本从而失去该时间点之前的时间旅行能力生产环境需评估合规与审计需求。在 Data Factory 编排中Fabric 还提供了专门的Lakehouse Maintenance管道活动用于自动执行 OPTIMIZE 与 VACUUM见 getdata.md可纳入日常 ETL 调度。增量加载模式Watermark增量加载只处理新增/变更数据是控制 ETL 成本与延迟的关键模式。参考 pyspark.md 中的经典实现from pyspark.sql.functions import col # Get last processed watermark last_watermark spark.sql( SELECT MAX(processed_timestamp) as watermark FROM silver_orders ).collect()[0][watermark] # Load only new records new_records spark.read.format(delta) \ .table(bronze_orders) \ .filter(col(created_at) last_watermark) # Merge new records new_records.createOrReplaceTempView(staging_orders) spark.sql( MERGE INTO silver_orders AS target USING staging_orders AS source ON target.order_id source.order_id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * )该模式的执行流程读取水位线从目标表silver_orders中取出上次处理的最大时间戳processed_timestamp作为本次增量窗口的下界。过滤新记录从源表bronze_orders中筛选created_at last_watermark的记录只对新增数据做后续处理。临时视图 MERGE把新记录注册为临时视图staging_orders再用 MERGE 以order_id为键做 Upsert实现“已存在则更新、不存在则插入”。该模式天然幂等重复执行只会重放上一次水位线之后的记录配合 MERGE 不会产生重复主键。对于完全流式的场景也可将new_records的读取替换为readStream并把水位线维护改为foreachBatch中的持久化逻辑。SCD Type 2 模式缓慢变化维对于需要保留历史版本的维度表如客户地址变更史使用 SCD Type 2每次变更不覆盖旧记录而是将旧记录关闭is_currentfalse并插入新版本。参考文档中的实现from pyspark.sql.functions import current_timestamp, lit # Close existing records spark.sql( UPDATE dim_customer SET is_current false, end_date current_timestamp() WHERE customer_id IN (SELECT customer_id FROM staging_customer) AND is_current true ) # Insert new versions spark.sql( INSERT INTO dim_customer SELECT customer_id, name, email, address, current_timestamp() as start_date, null as end_date, true as is_current FROM staging_customer )两步配合的语义关闭旧版本把staging_customer中出现的customer_id且当前仍为is_currenttrue的记录置为过期end_date打上当前时间戳。插入新版本将暂存表中的记录作为新版本插入start_datecurrent_timestamp()、end_datenull、is_currenttrue。最终每行通过is_current标识当前有效版本通过start_date/end_date保留历史区间业务侧查询“当前状态”过滤is_currenttrue审计与回溯查询直接读历史行。该模式与增量加载结合时staging_customer即为增量批次产生的临时视图形成完整的“增量 版本化”维度更新链路。端到端落地结合 Data Factory 的管道编排上述 PySpark 代码通常不是孤立运行的而是嵌入 Data Factory 管道详见 getdata.md。Fabric Data Factory 提供 180 数据源连接器、Copy Activity、Dataflow Gen2、Notebook Activity 与调度触发器常用管道活动包括活动说明Copy Data在数据源与 Lakehouse 之间搬运数据Notebook执行 Spark 笔记本即运行本文的 PySpark 代码Dataflow运行 Dataflow Gen2 转换Stored Procedure执行 SQL 存储过程ForEach遍历项If Condition条件分支Get Metadata获取文件/文件夹元数据Lakehouse Maintenance优化与 VACUUM Delta 表典型的日级 ETL 编排Bronze → Silver → GoldPipeline: Daily_ETL_Pipeline ├── Get Metadata (check for new files) ├── ForEach (process each file) │ ├── Copy Data (bronze layer) │ └── Notebook (silver transformation) ├── Notebook (gold aggregation) └── Lakehouse Maintenance (optimize tables)落地建议将本文的“Spark 会话配置”置于每个 Notebook 起始单元格将“增量加载 / SCD Type 2”封装为可复用 Notebook由管道按表传入参数将“OPTIMIZE / VACUUM”交给 Lakehouse Maintenance 活动在数据写入后执行形成“写入 → 增量合并 → 表维护”的闭环。小结本指南完整覆盖了 awesome-copilot 中 fabric-lakehouse 技能的 PySpark 参考 的全部核心模式会话级 V-Order 与 optimizeWrite 配置、多格式数据读取、Delta 托管表写入与 CRUD、显式 Schema 定义、Notebook SQL Magic、OPTIMIZE/ZORDER/VACUUM 表维护以及增量加载Watermark MERGE与 SCD Type 2 两个生产级数据模式并补充了 SKILL.md 中的 Lakehouse 存储模型、V-Order 原理与 getdata.md 的管道编排构成从代码片段到端到端调度的完整参考。在使用时请注意所有 SQL 条件如last_login阈值与表名均为示例需替换为实际表结构与业务口径VACUUM 保留期设置需结合时间旅行与合规需求综合权衡。【免费下载链接】awesome-copilotCommunity-contributed instructions, agents, skills, and configurations to help you make the most of GitHub Copilot.项目地址: https://gitcode.com/GitHub_Trending/aw/awesome-copilot创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表