如何用wr.s3.read_parquet对Parquet分区列去重过滤?
解决S3 Parquet数据集分区字段去重问题
问题根源
你之前的代码存在两个关键问题:
- 构造的路径已经指向了单个具体分区(
fecha_dato={fecha_dato_formato}/),读取的自然是该分区的全量数据,filters条件完全无效。 - awswrangler的
filters参数仅支持==, !=, >, <, >=, <=这类比较操作符,不存在unique这种去重逻辑的操作。
解决方案
方案1:读取分区列后去重(简单直接)
直接指定只加载fecha_dato列,再通过drop_duplicates去重,不需要读取全量数据:
# 读取整个数据集的fecha_dato分区列 df_fecha = wr.s3.read_parquet( path=param.delivery["output_path"], dataset=True, columns=["fecha_dato"] # 仅加载分区列,大幅减少IO ) # 去重得到唯一值 unique_fechas = df_fecha.drop_duplicates() print(unique_fechas.head(5))
方案2:直接枚举S3分区前缀(性能最优)
如果只需要获取分区的唯一值,不需要读取Parquet文件内容,直接列举S3上的分区前缀即可,效率更高:
import boto3 import pandas as pd s3 = boto3.client('s3') # 从S3路径中拆分bucket和前缀 s3_path = param.delivery["output_path"] bucket = s3_path.split('/')[2] prefix = '/'.join(s3_path.split('/')[3:]) + 'fecha_dato=' # 列举所有fecha_dato分区的前缀 response = s3.list_objects_v2(Bucket=bucket, Prefix=prefix, Delimiter='/') unique_fechas = [] for prefix_item in response.get('CommonPrefixes', []): # 提取分区值 fecha_val = prefix_item['Prefix'].split('=')[-1].rstrip('/') unique_fechas.append(fecha_val) # 转换为DataFrame格式 df_unique_fechas = pd.DataFrame({'fecha_dato': unique_fechas}) print(df_unique_fechas.head(5))
内容的提问来源于stack exchange,提问作者Jeanpiere Alcocer
相关产品推荐
相关产品推荐

