Spark如何自动读取S3桶中每周一更新的周度Parquet文件
S3周度Parquet文件自动读取实现方案
不需要依赖特殊的内置日期函数,用Python标准库datetime就能实现符合需求的日期计算,完全匹配「整周读取固定周度上传文件、到点自动切换新文件」的逻辑。
核心计算逻辑
你的场景是文件固定每周一上传,整周运行都读当周周一的文件,核心就是先计算出运行当日对应的当周周一日期,格式化后拼接成完整S3路径即可:
- 用
datetime.date.today()获取当前日期 - 用
weekday()方法计算当前日期到本周一的天数差:该方法返回值规则为周一=0、周二=1……周日=6,返回值就是当前日期距离本周一的间隔天数 - 用当前日期减去间隔天数,就能得到当周周一的日期
- 按照文件命名的日期格式转成字符串,拼接成完整S3读取路径
可直接复用的代码
import datetime from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 计算当周周一日期 today = datetime.date.today() days_to_monday = today.weekday() current_week_monday = today - datetime.timedelta(days=days_to_monday) # 匹配文件名的日期格式,示例为XXXXXX_2022-06-15这种横杠分隔格式 file_date = current_week_monday.strftime("%Y-%m-%d") # 拼接完整S3路径,替换成你自己的桶名和文件前缀即可 s3_file_path = f"s3://你的桶名/文件前缀_{file_date}.parquet" # 读取文件,注意读取parquet时format要对应写parquet,不要写错成csv sdf = spark.read.format("parquet").load(s3_file_path)
特殊场景适配
- 如果周一任务运行时新文件还未上传完成,需要读取上一周的周一文件,只需要在计算间隔时多减7天即可:
# 计算上周一日期 last_week_monday = today - datetime.timedelta(days=days_to_monday + 7) - 如果固定上传日期不是周一,比如每周三上传,只需要调整偏移计算逻辑,就能自动取最近一个上传日的日期:
周三对应
weekday()返回值为2,偏移量计算为(today.weekday() - 2) % 7,用当前日期减去这个偏移量就能得到最近一个周三的日期,其他上传日同理替换对应weekday值即可。 - 如果文件名的日期格式不是
YYYY-MM-DD,比如是无横杠的YYYYMMDD格式,只需要修改strftime()的参数为"%Y%m%d"即可。
内容的提问来源于stack exchange,提问作者DOT GAMING
相关产品推荐
相关产品推荐

