Spark缓冲区大小限制问题的PySpark解决方案咨询
解决PySpark缓冲区大小限制问题的三种方案
原统计代码运行时遇到缓冲区大小限制问题,代码如下:
# Calculate the statistics stats = df.groupBy("EventType").agg( size(collect_set("Parameters")).alias("ParameterLength"), collect_list("Parameters").alias("Parameters"), (count("*") / df.count() * 100).alias("Frequency"), )
以下是针对需求的三种解决方案:
方案1:拆分处理大参数数据
先过滤掉EventType为"A"或"B"的数据生成主汇总表,再单独将这两类数据输出至独立文件:
from pyspark.sql.functions import col, size, collect_set, collect_list, count # 生成排除A、B类的主汇总表 main_stats = df.filter(~col("EventType").isin("A", "B")).groupBy("EventType").agg( size(collect_set("Parameters")).alias("ParameterLength"), collect_list("Parameters").alias("Parameters"), (count("*") / df.filter(~col("EventType").isin("A", "B")).count() * 100).alias("Frequency"), ) # 单独处理A、B类数据并输出到独立文件 ab_df = df.filter(col("EventType").isin("A", "B")) ab_stats = ab_df.groupBy("EventType").agg( size(collect_set("Parameters")).alias("ParameterLength"), collect_list("Parameters").alias("Parameters"), (count("*") / ab_df.count() * 100).alias("Frequency"), ) # 输出到独立文件,可根据需求调整存储格式和路径 ab_stats.write.mode("overwrite").parquet("path/to/ab_event_stats")
方案2:截断A、B类的参数列表
对EventType为"A"或"B"的Parameters列表截断至前100个元素,非A/B类保持完整:
from pyspark.sql.functions import when, slice, collect_list, size, count, col stats = df.groupBy("EventType").agg( # 针对A/B类取实际长度与100的最小值,其余正常计算 when(col("EventType").isin("A", "B"), size(slice(collect_list("Parameters"), 1, 100)) ).otherwise(size(collect_set("Parameters"))).alias("ParameterLength"), # 截断A/B类的参数列表至前100个 when(col("EventType").isin("A", "B"), slice(collect_list("Parameters"), 1, 100) ).otherwise(collect_list("Parameters")).alias("Parameters"), (count("*") / df.count() * 100).alias("Frequency"), )
注:
slice函数参数为slice(列名, 起始位置, 截取长度),这里从第1个元素开始截取100个,超过100的部分会被自动截断。
方案3:缩减Parameters数据体积
将Parameters中的{"Key":XXX, "Value":YYY}字典结构转为Spark StructType或元组(数组),减少冗余存储开销:
转为StructType
from pyspark.sql.functions import struct, collect_list, size, count, col # 将字典转为StructType,保留Key和Value字段 df_struct = df.withColumn("Parameters", struct(col("Parameters.Key").alias("Key"), col("Parameters.Value").alias("Value"))) # 执行统计逻辑 stats = df_struct.groupBy("EventType").agg( size(collect_set("Parameters")).alias("ParameterLength"), collect_list("Parameters").alias("Parameters"), (count("*") / df_struct.count() * 100).alias("Frequency"), )
转为元组(数组)
from pyspark.sql.functions import array, collect_list, size, count, col # 将字典转为(Key, Value)数组形式,Spark中元组以数组结构存储 df_tuple = df.withColumn("Parameters", array(col("Parameters.Key"), col("Parameters.Value"))) # 执行统计逻辑 stats = df_tuple.groupBy("EventType").agg( size(collect_set("Parameters")).alias("ParameterLength"), collect_list("Parameters").alias("Parameters"), (count("*") / df_tuple.count() * 100).alias("Frequency"), )
注:StructType和数组结构相比原始字典,省去了重复存储"Key"、"Value"键名的开销,能有效降低单条数据体积,缓解缓冲区压力。
内容的提问来源于stack exchange,提问作者Yiffany
相关产品推荐
相关产品推荐

