基于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
相关产品推荐
相关产品推荐

