如何在PySpark中按需处理DataFrame条件并拆分重复/唯一数据
PySpark 实现方案
1. 初始化环境与创建示例DataFrame
首先创建Spark会话并加载示例数据:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, concat, lit, count # 初始化Spark会话 spark = SparkSession.builder.appName("ItemDataProcessing").getOrCreate() # 示例数据 data = [ ("20-10767-58V", 98003351), ("20-10087-58V", 87003872), ("20-10087-58V", 97098411), ("20-10i72-YTW", 99003351), ("27-1o121-YTW", 89659352), ("27-10991-YTW", 98678411), ("At81kk00", 98903458), ("Avp12225", 85903458), ("Akb12226", 99003458), ("Ahh12829", 98073458), ("Aff12230", 88803458), ("Ar412231", 92003458), ("Aju12244", 98773458) ] # 创建初始DataFrame df = spark.createDataFrame(data, ["item_cd", "item_nbr"])
2. 处理item_cd字段
使用条件判断对item_cd进行处理:含连字符-的保留原值,不含则追加4个尾随0:
processed_df = df.withColumn( "processed_item_cd", when(col("item_cd").contains("-"), col("item_cd")) .otherwise(concat(col("item_cd"), lit("0000"))) )
3. 拆分重复与唯一数据
先统计每个(processed_item_cd, item_nbr)组合的出现次数,再据此拆分数据:
# 统计组合出现次数 count_df = processed_df.groupBy("processed_item_cd", "item_nbr").agg(count("*").alias("occurrence_count")) # 提取唯一数据(仅出现1次的组合) unique_data_df = processed_df.join(count_df, on=["processed_item_cd", "item_nbr"], how="inner")\ .filter(col("occurrence_count") == 1)\ .drop("occurrence_count") # 提取重复数据(出现多次的组合) duplicate_data_df = processed_df.join(count_df, on=["processed_item_cd", "item_nbr"], how="inner")\ .filter(col("occurrence_count") > 1)\ .drop("occurrence_count")
验证结果
可以通过show()方法查看处理后的结果:
# 查看唯一数据 unique_data_df.show() # 查看重复数据 duplicate_data_df.show()
内容的提问来源于stack exchange,提问作者venkat
相关产品推荐
相关产品推荐

