如何使用PySpark读取S3中特定时间点之后创建的JSON文件
PySpark 读取S3指定时间后创建的JSON文件实现方案
方案1:预过滤S3文件路径再读取(推荐,适合数十万文件场景)
该方案先通过S3 SDK直接遍历目标路径下的文件元数据,过滤出符合时间要求的文件路径后再提交给Spark读取,避免Spark全量扫描文件的额外开销,性能远高于后处理过滤。
- 前置依赖:安装boto3
pip install boto3 - 代码示例:
import boto3 from datetime import datetime # 初始化S3客户端,提前配置好IAM权限或AK/SK s3 = boto3.client('s3') bucket_name = "替换为你的S3桶名" prefix = "替换为JSON文件所在的路径前缀(无需带桶名)" # 时间阈值需指定UTC时区,和S3的文件修改时间时区对齐 cutoff_time = datetime(2024, 1, 1, tzinfo=datetime.timezone.utc) valid_paths = [] continuation_token = None # 分页遍历S3对象,避免单次请求返回数据超限 while True: kwargs = {'Bucket': bucket_name, 'Prefix': prefix} if continuation_token: kwargs['ContinuationToken'] = continuation_token resp = s3.list_objects_v2(**kwargs) for obj in resp.get('Contents', []): # 过滤文件夹、非JSON文件、时间早于阈值的文件 if not obj['Key'].endswith('/') and obj['Key'].endswith('.json') and obj['LastModified'] > cutoff_time: valid_paths.append(f"s3a://{bucket_name}/{obj['Key']}") if not resp.get('IsTruncated'): break continuation_token = resp.get('NextContinuationToken') # 仅读取过滤后的合规文件 df = spark.read.json(valid_paths)
方案2:通过Spark内置元数据列过滤(写法更简洁)
Spark 2.1及以上版本支持直接读取文件的系统元数据,包括文件修改时间、路径等信息,过滤逻辑直接通过Spark SQL算子实现,无需额外依赖。缺点是Spark会先扫描全量文件的元数据,文件量级特别大时性能弱于方案1。
- 代码示例:
from pyspark.sql.functions import col # 读取所有JSON文件,同时携带文件元数据 df = spark.read.json("s3a://替换为你的桶名/替换为路径前缀/*") \ .withColumn("file_modify_time", col("_metadata.file_modification_time")) # 过滤指定时间之后的文件,时间后缀Z代表UTC时区 filtered_df = df.where(col("file_modify_time") > "2024-01-01T00:00:00Z") # 可选:删除不需要的元数据列 final_df = filtered_df.drop("file_modify_time")
注意事项
- 如果你的S3路径已经按时间做了分区命名,比如
s3://bucket/path/year=2024/month=05/day=10/格式,优先使用Spark分区裁剪,直接添加where(year>=2024 and month>=1 and day>=1)过滤条件,性能比以上两个方案都高 - 所有时间过滤逻辑都要注意时区对齐,S3返回的文件修改时间默认是UTC时区,避免时区差导致漏读或多读文件
- 提前配置好S3访问权限,不管是boto3还是Spark的S3A客户端,都需要有对应路径的读权限
内容的提问来源于stack exchange,提问作者Sree
相关产品推荐
相关产品推荐

