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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 20:51:12