Pyspark如何在RDD Map函数中过滤DataFrame并基于全表计算新增列
错误原因
你遇到的PicklingError是因为在RDD的map算子回调函数中直接引用了Driver端创建的DataFrame对象:Spark需要将算子闭包序列化后发送给Executor执行,而DataFrame内部包含不可序列化的_thread.RLock线程锁对象,因此触发序列化失败。
此外这种实现方式存在严重性能问题:每条数据处理都会触发一次全表扫描的Action操作,数据量稍大就会完全不可用。
正确实现方案
建议使用Spark原生的自连接+聚合算子实现需求,逻辑和你原代码的统计规则完全一致,且无需在Executor端引用Driver对象:
from pyspark.sql import functions as F df = SPARK_DATAFRAME # 第一步:提前提取quittierung的时间部分 df = df.withColumn("quittierung_time", F.substring("quittierung", 12, 8)) # 第二步:自连接匹配统计条件 df_left = df.alias("left") df_right = df.alias("right") joined_df = df_left.join( df_right, (F.col("left.datum") == F.col("right.datum")) & (F.col("left.start") == F.col("right.start")) & (F.col("right.uhrzeit") < F.col("left.uhrzeit")) & (F.col("right.quittierung_time") > F.col("left.uhrzeit")), how="left" # 左连接保留原始表所有行 ) # 第三步:分组统计符合条件的行数,得到最终结果 result_df = joined_df.groupBy( "left.datum", "left.start", "left.uhrzeit", "left.quittierung" ).agg(F.count("right.datum").alias("count")) result_df.show()
如果是同分组(datum+start)下数据量不大的场景,也可以用窗口函数实现,性能更高:
from pyspark.sql import functions as F from pyspark.sql.window import Window df = SPARK_DATAFRAME df = df.withColumn("quittierung_time", F.substring("quittierung", 12, 8)) # 按datum、start分区,统计同分组下所有符合条件的行数 w = Window.partitionBy("datum", "start") df = df.withColumn("count", F.sum( F.when( (F.col("uhrzeit") < F.col("uhrzeit").over(w)) & (F.col("quittierung_time").over(w) > F.col("uhrzeit")), 1 ).otherwise(0) ).over(w)) df.show()
内容的提问来源于stack exchange,提问作者Giuseppe
相关产品推荐
相关产品推荐

