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

PySpark转换后丢失DataFrame记录,如何实现与Pandas一致的特征频率列?

PySpark 为全量行新增特征组合统计频率列

我是PySpark新手,希望基于合成数据为PySpark中每条记录/行/事件新增一列,展示对应特征的归一化统计频率(其中Type与Encoding_type列为类别型)。初始我有包含5列、名为sdf的Spark DataFrame,示例如下:

from pyspark.sql.types import StructType,StructField, StringType, IntegerType
data2 = [("Sentence",92,6,"False",49),
         ("Sentence",17,3,"False",15),
         ("Sentence",17,3,"False",15),
         (0         , 0,0,"False", 0),
         (0         , 0,0,"False", 0),
         (0         , 0,0,"False", 0)
  ]

schema = StructType([ 
    StructField("Type",              StringType(),  True), 
    StructField("Length",            IntegerType(), True), 
    StructField("Token_number",      IntegerType(), True), 
    StructField("Encoding_type",     StringType(), True), 
    StructField("Character_feature", IntegerType(), True) 
    
  ])
 
sdf = spark.createDataFrame(data=data2,schema=schema)
sdf.printSchema()
sdf.show(truncate=False)

# 输出
# root
# |-- Type: string (nullable = true)
# |-- Length: integer (nullable = true)
# |-- Token_number: integer (nullable = true)
# |-- Encoding_type: string (nullable = true)
# |-- Character_feature: integer (nullable = true)

# +--------+------+------------+-------------+-----------------+
# |Type    |Length|Token_number|Encoding_type|Character_feature|
# +--------+------+------------+-------------+-----------------+
# |Sentence|92    |6           |False        |49               |
# |Sentence|17    |3           |False        |15               |
# |Sentence|17    |3           |False        |15               |
# |0       |0     |0           |False        |0                |
# |0       |0     |0           |False        |0                |
# |0       |0     |0           |False        |0                |
# +--------+------+------------+-------------+-----------------+

现在我希望计算特征的统计频率并映射回主表sdf,当前实现代码如下:

# Statistical Preprocessing
def add_freq_to_features_(df):
  sdf_pltc = sdf.select('Type', 'Length', 'Token_number', 'Encoding_type', 'Character_feature')
  #sdf_pltc.show(truncate=0)

  sdf2 = (
      sdf_pltc
      .groupBy(sdf_pltc.columns)
      .agg(F.count('*').alias('Freq'))
    # .withColumn('Freq' , (col('Freq')  / col('Freq').sum()))    # Normalzing between 0 & 1
      .withColumn('Encoding_type', F.col('Encoding_type').cast('string'))
  )

  sdf2.show()
  
  return new_df

# Apply frequency allocation and merge with extracted features df
features_df = add_freq_to_features_(df)
features_df

# 当前输出
# +--------+------+------------+-------------+-----------------+----+
# |    Type|Length|Token_number|Encoding_type|Character_feature|Freq|
# +--------+------+------------+-------------+-----------------+----+
# |Sentence|    92|           6|        False|               49|   1|
# |       0|     0|           0|        False|                0|   3|
# |Sentence|    17|           3|        False|               15|   2|
# +--------+------+------------+-------------+-----------------+----+

受groupby()机制影响,当前PySpark输出和Pandas版本的结果不同,Pandas实现如下:

import pandas as pd

data = {'Type':              ["Sentence" , "Sentence" ,"Sentence" , "-", "-", "-"], 
        'Length':            [92, 17,17,0,0,0],
        'Token_number':      [6, 3,3,0,0,0],
        'Encoding_type':     ["False" , "False" ,"False" , "False", "False", "False"], 
        'Character_feature': [49, 15,15,0,0,0],
       }

# pass column names in the columns parameter 
df = pd.DataFrame(data)
#df

#Statistical Preprocessing
def add_freq_to_features(df):
  frequencies_df = df.groupby(list(df.columns)).size().to_frame().rename(columns={0: "Freq"})
 # frequencies_df["Freq"] = frequencies_df["Freq"] / frequencies_df["Freq"].sum() # Normalzing 0 & 1
  new_df = pd.merge(df, frequencies_df, how='left', on=list(df.columns))
  
  return new_df

features_df = add_freq_to_features(df)
features_df

# 期望输出
# +----------+------+------------+-------------+-----------------+-----+
# |Type      |Length|Token_number|Encoding_type|Character_feature|Freq |
# +----------+------+------------+-------------+-----------------+-----+
# |Sentence  |92    |6           |False        |49               |1    |
# |Sentence  |17    |3           |False        |15               |2    |
# |Sentence  |17    |3           |False        |15               |2    |
# |-         |0     |0           |False        |0                |3    |
# |-         |0     |0           |False        |0                |3    |
# |-         |0     |0           |False        |0                |3    |
# +----------+-----+-------------+-------------+-----------------+-----+

我暂未找到在PySpark中得到和Pandas版本一致输出的方法,可能是函数使用有误。


解决方案

你当前输出和预期不一致的核心原因是:分组聚合后的结果只有去重后的特征组合,没有将统计到的频率映射回原表的所有行。可以通过两种方法实现预期效果:

方法1:窗口函数(性能更优,仅需一次shuffle)

直接使用开窗统计,无需分组再join,一步生成所有行的频率列:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

def add_freq_to_features_(df):
    # 定义窗口:按所有列分区,相同特征组合的行属于同一个分区
    window_spec = Window.partitionBy(df.columns)
    # 统计每个分区的行数,即当前特征组合的频率
    new_df = df.withColumn("Freq", F.count("*").over(window_spec))
    # 若需要归一化到0-1区间,打开下面两行注释即可
    # total_cnt = df.count()
    # new_df = new_df.withColumn("Freq", F.col("Freq") / total_cnt)
    return new_df

features_df = add_freq_to_features_(sdf)
features_df.show()

方法2:分组后左连接(和Pandas逻辑完全对齐)

在你原有代码基础上补全左连接步骤即可:

from pyspark.sql import functions as F

def add_freq_to_features_(df):
    # 分组统计每个特征组合的频率
    freq_df = df.groupBy(df.columns).agg(F.count("*").alias("Freq"))
    # 左连接回原表,将频率映射到所有行
    new_df = df.join(freq_df, on=df.columns, how="left")
    # 若需要归一化到0-1区间,打开下面两行注释即可
    # total_cnt = df.count()
    # new_df = new_df.withColumn("Freq", F.col("Freq") / total_cnt)
    return new_df

features_df = add_freq_to_features_(sdf)
features_df.show()

两种方法的输出均与你给出的Pandas版本结果完全一致。


内容的提问来源于stack exchange,提问作者Mario

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 08:54:02