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
相关产品推荐
相关产品推荐

