PySpark/Python处理AWS Athena与S3 Parquet数据的去重及清理
处理AWS Athena与S3中大型Parquet数据集的分组去重与清理需求
核心需求
- 需求1:针对AWS Athena中的大型Parquet文件,按
customer_id、is_active、city、married_status分组(排除as_of_date字段),保留每组对应的as_of_date值,最终得到每个分组的唯一行(同一customer_id下,当其他字段组合变化时保留对应日期记录) - 需求2:清理Amazon S3中按
as_of_date分区、总计2.5亿行的Parquet文件,输出仅含上述分组唯一行的Parquet文件(数据集包含30万唯一customer_id)
输入数据示例(Pandas)
import pandas as pd dict_1 = { 'customer_id': [1,1,1,1,2,2,2,2], 'is_active': ['y','y','n','n','n','y','y','n'], 'city':['NYC','NYC','NYC','NYC','SFO','SFO','SFO','SFO'], 'married_status':['Single','Single','Single','Single','Single','Single','Single','Single'], 'as_of_date':['01/01/2023', '01/02/2023','01/03/2023','01/04/2023','01/01/2023','01/02/2023','01/03/2023','01/04/2023'] } df1 = pd.DataFrame(dict_1, columns = ['customer_id','is_active','city','married_status','as_of_date'])
期望输出示例(Pandas)
dict_output = { 'customer_id': [1,1,2,2,2], 'is_active': ['y','n','n','y','n'], 'city':['NYC','NYC','SFO','SFO','SFO'], 'married_status':['Single','Single','Single','Single','Single'], 'as_of_date':['01/01/2023', '01/03/2023','01/01/2023','01/02/2023','01/04/2023'] } df_output = pd.DataFrame(dict_output, columns = ['customer_id','is_active','city','married_status','as_of_date'])
解决方案
1. PySpark批量处理S3分区数据
针对2.5亿行的大数据场景,PySpark是高效的分布式处理方案,步骤如下:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, row_number # 初始化Spark会话 spark = SparkSession.builder.appName("CustomerDeduplication").getOrCreate() # 读取S3上的Parquet分区数据 df = spark.read.parquet("s3://your-bucket/path/to/source-parquet/") # 定义窗口:按目标字段分组,按as_of_date排序(默认取最早日期记录) window_spec = Window.partitionBy("customer_id", "is_active", "city", "married_status").orderBy("as_of_date") # 添加行号,筛选每组第一行实现去重 deduplicated_df = df.withColumn("row_num", row_number().over(window_spec)) \ .filter(col("row_num") == 1) \ .drop("row_num") # 将清理后的数据写入S3(支持覆盖原有数据) deduplicated_df.write.mode("overwrite").parquet("s3://your-bucket/path/to/cleaned-parquet/")
若需保留每组最新的
as_of_date记录,将orderBy("as_of_date")改为orderBy(col("as_of_date").desc())即可。
2. AWS Athena SQL查询与导出
若直接在Athena中处理数据,可通过SQL实现分组去重,支持直接导出结果到S3:
基础查询语句
SELECT customer_id, is_active, city, married_status, as_of_date FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY customer_id, is_active, city, married_status ORDER BY as_of_date ) AS row_num FROM your_athena_table_name ) t WHERE row_num = 1;
直接创建清理后的Athena表(导出到S3)
CREATE TABLE cleaned_customer_data WITH ( format = 'PARQUET', external_location = 's3://your-bucket/path/to/cleaned-athena-table/' ) AS SELECT customer_id, is_active, city, married_status, as_of_date FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY customer_id, is_active, city, married_status ORDER BY as_of_date ) AS row_num FROM your_athena_table_name ) t WHERE row_num = 1;
内容的提问来源于stack exchange,提问作者rabasa97
相关产品推荐
相关产品推荐

