使用PySpark识别同名同邮编客户并标记首项生成新表
PySpark 实现客户重复识别与标记方案
需求说明
识别姓名与邮政编码相同的不同客户:
- 每组中首次出现的记录:
FIRST_CUSTOMER_ID填充为该组首个客户ID,NEW_OLD_FLAG标记为true - 同组其余记录:
FIRST_CUSTOMER_ID填充为组内首个客户ID,NEW_OLD_FLAG标记为false
最终将处理结果存储至新表。
示例数据
输入表(customer_details)
| customer_id | name | zip_code | |
|---|---|---|---|
| 101 | Alice | 90210 | alice@example.com |
| 102 | Alice | 90210 | alice.a@example.com |
| 103 | Bob | 10001 | bob@example.com |
| 104 | Bob | 10001 | bob.b@example.com |
| 105 | Charlie | 60601 | charlie@example.com |
输出表(processed_customers)
| customer_id | name | zip_code | FIRST_CUSTOMER_ID | NEW_OLD_FLAG | |
|---|---|---|---|---|---|
| 101 | Alice | 90210 | alice@example.com | 101 | true |
| 102 | Alice | 90210 | alice.a@example.com | 101 | false |
| 103 | Bob | 10001 | bob@example.com | 103 | true |
| 104 | Bob | 10001 | bob.b@example.com | 103 | false |
| 105 | Charlie | 60601 | charlie@example.com | 105 | true |
PySpark 实现代码
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, first, row_number # 初始化SparkSession spark = SparkSession.builder.appName("CustomerDuplicateHandle").getOrCreate() # 读取原始客户表(可根据实际存储格式替换为csv/jdbc等) raw_df = spark.read.parquet("path/to/customer_details") # 定义窗口规则:按姓名+邮编分组,按客户ID排序(若有创建时间字段,建议用时间排序更贴合"首次出现"逻辑) window_spec = Window.partitionBy("name", "zip_code").orderBy("customer_id") # 计算分组行号、填充首客户ID、标记首次记录 processed_df = raw_df.withColumn( "row_num", row_number().over(window_spec) ).withColumn( "FIRST_CUSTOMER_ID", first("customer_id").over(window_spec) ).withColumn( "NEW_OLD_FLAG", col("row_num") == 1 ).drop("row_num") # 移除临时行号字段 # 写入结果表(支持overwrite追加模式,存储格式可按需调整) processed_df.write.mode("overwrite").parquet("path/to/processed_customers") # 若需写入Hive表,可替换为: # processed_df.write.mode("overwrite").saveAsTable("your_database.processed_customers")
关键细节说明
- 窗口分组逻辑:
partitionBy("name", "zip_code")确保同姓名同邮编的客户被归为一组;orderBy字段决定"首次出现"的判定依据,优先用业务上的时间字段(如create_time)而非ID。 first()函数:在窗口内提取分组的首个客户ID,保证同组所有记录的FIRST_CUSTOMER_ID一致。- 标记字段生成:通过行号判断是否为组内第一条记录,行号为1时标记为
true,其余为false。
内容的提问来源于stack exchange,提问作者Blogger22
相关产品推荐
相关产品推荐

