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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:44:55