You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于PySpark DataFrame将文件增量加载至Azure Delta表VAS

实现Azure存储中CSV文件全量+增量加载至Delta表(PySpark)

一、环境配置

首先配置Azure存储访问凭据(以ADLS Gen2的Service Principal认证为例),并启用Delta Lake支持:

# 替换为你的存储账户、租户ID、客户端ID和密钥
spark.conf.set("fs.azure.account.auth.type.<storage-account>.dfs.core.windows.net", "OAuth")
spark.conf.set("fs.azure.account.oauth.provider.type.<storage-account>.dfs.core.windows.net", "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider")
spark.conf.set("fs.azure.account.oauth2.client.id.<storage-account>.dfs.core.windows.net", "<client-id>")
spark.conf.set("fs.azure.account.oauth2.client.secret.<storage-account>.dfs.core.windows.net", "<client-secret>")
spark.conf.set("fs.azure.account.oauth2.client.endpoint.<storage-account>.dfs.core.windows.net", "https://login.microsoftonline.com/<tenant-id>/oauth2/token")

# 导入所需PySpark函数与Delta模块
from delta.tables import *
from pyspark.sql.functions import current_date, input_file_name

二、全量加载初始文件

加载VAS文件夹下的初始CSV文件,添加加载日期和文件名列,写入Delta主表,并创建文件跟踪表记录已加载文件:

# 源文件夹路径
source_path = "abfss://<container>@<storage-account>.dfs.core.windows.net/VAS/"

# 加载CSV并添加衍生列
full_load_df = spark.read.csv(source_path, header=True, inferSchema=True) \
    .withColumn("load_date", current_date()) \
    .withColumn("file_name", input_file_name())

# 写入主Delta表VAS
full_load_df.write.format("delta").mode("overwrite").saveAsTable("VAS")

# 创建已加载文件跟踪表,用于后续去重
loaded_files_df = full_load_df.select("file_name").distinct()
loaded_files_df.write.format("delta").mode("overwrite").saveAsTable("VAS_Loaded_Files")

三、增量加载新增文件

次日文件夹新增文件后,通过对比跟踪表筛选未加载文件,完成增量加载并更新跟踪表:

# 获取已加载的文件列表
loaded_files = spark.table("VAS_Loaded_Files").select("file_name").rdd.flatMap(lambda x: x).collect()
loaded_files_set = set(loaded_files)

# 列出源文件夹下所有CSV文件
all_csv_files = [file.path for file in dbutils.fs.ls(source_path) if file.path.endswith(".csv")]
# 筛选未加载的新增文件
new_files = [file for file in all_csv_files if file not in loaded_files_set]

if new_files:
    # 加载新增文件并添加衍生列
    incremental_df = spark.read.csv(new_files, header=True, inferSchema=True) \
        .withColumn("load_date", current_date()) \
        .withColumn("file_name", input_file_name())

    # 追加数据到主Delta表
    incremental_df.write.format("delta").mode("append").saveAsTable("VAS")

    # 更新已加载文件跟踪表
    new_files_df = incremental_df.select("file_name").distinct()
    new_files_df.write.format("delta").mode("append").saveAsTable("VAS_Loaded_Files")
else:
    print("无新增CSV文件需要加载")

核心要点

  • 文件去重逻辑:通过VAS_Loaded_Files表持久化已加载文件名,比依赖文件修改时间更稳定,适用于文件可能被重写的场景。
  • 衍生列说明:current_date()生成数据加载当日的日期,input_file_name()返回当前数据行对应的源CSV文件路径。
  • 认证方式适配:如果使用SAS令牌访问存储,替换为以下配置:
    spark.conf.set("fs.azure.account.key.<storage-account>.dfs.core.windows.net", "<your-sas-token>")
    

内容的提问来源于stack exchange,提问作者Priyam Singh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 04:10:13