如何在PySpark版Glue Job中获取带前缀的S3存储桶元数据?
在Glue Job(PySpark)中获取S3 Parquet文件元数据的简便方法
方案一:利用Spark/Hadoop FileSystem API(推荐,无需额外依赖)
既然你已经在Glue中用PySpark处理Parquet文件,直接通过Spark底层的Hadoop文件系统API获取元数据是最贴合环境的方式,不需要额外引入boto3依赖,逻辑更简洁。
代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import input_file_name import datetime # Glue环境中可直接使用已初始化的spark对象,无需重复创建 spark = SparkSession.builder.getOrCreate() # 替换为你的Parquet文件存储路径 s3_parquet_path = "s3://your-bucket/parquet-output/" # 方式1:直接遍历路径下所有文件的元数据 hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration() fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) target_path = spark._jvm.org.apache.hadoop.fs.Path(s3_parquet_path) # 过滤掉S3的虚目录,只处理实际文件 for status in fs.listStatus(target_path): if not status.isDirectory(): full_path = status.getPath().toString() file_size = status.getLen() # 文件大小(字节) # 将毫秒时间戳转换为可读格式 modify_time = datetime.datetime.fromtimestamp(status.getModificationTime() / 1000) print(f"文件路径: {full_path}, 大小: {file_size} bytes, 修改时间: {modify_time}") # 方式2:结合DataFrame获取元数据(适合需与业务数据关联的场景) parquet_df = spark.read.parquet(s3_parquet_path) # 给每条数据关联其所属的文件路径 df_with_file = parquet_df.withColumn("file_path", input_file_name()) # 提取唯一文件路径并获取对应元数据 unique_files = df_with_file.select("file_path").distinct().collect() for row in unique_files: file_path = row["file_path"] file_status = fs.getFileStatus(spark._jvm.org.apache.hadoop.fs.Path(file_path)) print(f"文件路径: {file_path}, 大小: {file_status.getLen()} bytes, 修改时间: {datetime.datetime.fromtimestamp(file_status.getModificationTime()/1000)}")
方案二:修正boto3 Paginator的使用方式
如果坚持使用boto3,大概率是之前的分页逻辑或过滤规则有问题,以下是正确的实现:
代码实现
import boto3 from botocore.exceptions import ClientError # Glue默认会使用Job角色的权限,无需手动配置密钥 s3_client = boto3.client('s3') bucket_name = "your-bucket" # 注意前缀格式:如果要指定某目录下的文件,结尾可以不带斜杠,避免遗漏子目录文件 target_prefix = "path/to/parquet-output/" try: paginator = s3_client.get_paginator('list_objects_v2') # 分页遍历,Delimiter用于区分虚目录和实际文件 page_iterator = paginator.paginate(Bucket=bucket_name, Prefix=target_prefix, Delimiter='/') for page in page_iterator: # 跳过虚目录,只处理实际文件对象 if 'Contents' in page: for obj in page['Contents']: # 过滤掉可能的虚目录条目 if not obj['Key'].endswith('/'): file_key = obj['Key'] file_size = obj['Size'] create_time = obj['LastModified'] print(f"文件Key: {file_key}, 大小: {file_size} bytes, 创建时间: {create_time}") except ClientError as e: print(f"获取元数据失败: {e.response['Error']['Message']}")
注意事项
- 确保Glue Job的IAM角色拥有
s3:ListBucket权限(针对目标存储桶),否则两种方案都会失败。 - Spark返回的
getModificationTime是毫秒级时间戳,需要转换为可读格式;boto3返回的LastModified直接是datetime对象,可直接使用。 - 如果Parquet是分区存储的,两种方案都会自动遍历所有分区下的文件,无需额外处理。
内容的提问来源于stack exchange,提问作者Bee
相关产品推荐
相关产品推荐

