在PySpark中为重复DUNS_NUMBER生成递增索引的dup_duns_number列
在PySpark中创建dup_duns_number列的实现方法
实现思路
要达成需求,需分三步操作:
- 统计每个
DUNS_NUMBER的出现频次,识别出重复值(频次>1) - 为重复的
DUNS_NUMBER分配从1开始的递增索引 - 将索引关联回原数据集,非重复值填充
NaN
代码实现
假设你的数据集是名为df的PySpark DataFrame,具体代码如下:
from pyspark.sql import Window import pyspark.sql.functions as F # 1. 计算每个DUNS_NUMBER的出现次数 count_df = df.groupBy("DUNS_NUMBER").agg(F.count("*").alias("count")) # 2. 筛选出重复的DUNS_NUMBER,并分配递增索引 dup_index_df = count_df.filter(F.col("count") > 1) \ .withColumn("dup_DUNS_NUMBER", F.row_number().over(Window.orderBy("DUNS_NUMBER"))) # 3. 将索引关联回原DataFrame,非重复值填充NaN result_df = df.join(dup_index_df, on="DUNS_NUMBER", how="left") \ .select("DUNS_NUMBER", F.col("dup_DUNS_NUMBER").cast("double")) \ .na.fill(value=None, subset=["dup_DUNS_NUMBER"])
代码说明
- 统计频次:通过
groupBy和agg(count)获取每个编号的出现次数 - 分配递增索引:利用
row_number()窗口函数,对重复编号按顺序生成从1开始的连续索引 - 关联并填充NaN:左连接原表与索引表,将索引列转为double类型(匹配示例中的1.0格式),最后用
na.fill将空值替换为NaN
结果验证
运行代码后,result_df的结构会与你提供的示例数据完全一致:重复的DUNS_NUMBER对应相同的递增索引,非重复值显示NaN。
内容的提问来源于stack exchange,提问作者G R
相关产品推荐
相关产品推荐

