PySpark如何仅读取当日新增的Delta CSV文件?
用PySpark处理当日新增Delta相关文件的方案
一、仅读取当日上传的文件
生成当日日期后缀
先通过Python内置模块生成和文件后缀格式一致的当日日期字符串:from datetime import datetime # 生成YYYYMMDD格式的当日后缀 today_suffix = datetime.today().strftime("%Y%m%d")用通配符匹配当日文件
PySpark支持路径中使用通配符*过滤文件,直接拼接路径即可精准匹配当日新增的文件:# 替换为你的实际文件存储路径 target_path = f"/your/folder/path/*{today_suffix}.csv"该路径会自动匹配所有带当日日期后缀的CSV文件,刚好对应每日新增的2个文件。
二、加载文件(区分CSV和Delta Lake两种场景)
场景1:文件是带日期后缀的CSV(你描述的"Delta CSV文件"指增量CSV)
直接用Spark的CSV读取器加载匹配到的文件:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("DailyCSVProcess").getOrCreate() # 加载当日CSV,根据实际情况调整header、inferSchema参数 daily_df = spark.read.csv(target_path, header=True, inferSchema=True) # 在这里添加你的数据转换逻辑 # 示例:daily_df = daily_df.withColumn("load_date", lit(today_suffix))
场景2:文件是标准Delta Lake格式(后缀为delta目录)
如果是标准Delta Lake文件,用Delta专属格式加载:
# 匹配当日Delta目录 delta_target_path = f"/your/folder/path/*{today_suffix}.delta" daily_df = spark.read.format("delta").load(delta_target_path)
三、追加数据到现有表
处理完成后,用append模式将数据写入目标表:
写入普通Spark/Hive表
# 替换为你的目标表名 daily_df.write.mode("append").saveAsTable("your_target_table")
写入Delta Lake表
如果目标表是Delta格式,用以下方式:
# 写入Delta目录 daily_df.write.format("delta").mode("append").save("/path/to/your/delta_table") # 或写入Delta表 daily_df.write.format("delta").mode("append").saveAsTable("your_delta_target_table")
内容的提问来源于stack exchange,提问作者peace
相关产品推荐
相关产品推荐

