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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 20:01:22