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

PySpark如何仅读取当日新增的Delta CSV文件?

用PySpark处理当日新增Delta相关文件的方案

一、仅读取当日上传的文件

  1. 生成当日日期后缀
    先通过Python内置模块生成和文件后缀格式一致的当日日期字符串:

    from datetime import datetime
    
    # 生成YYYYMMDD格式的当日后缀
    today_suffix = datetime.today().strftime("%Y%m%d")
    
  2. 用通配符匹配当日文件
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:01:39