PySpark中DataFrame执行foreach操作内append列表外部为空问题
问题产生原因
- 第一个核心原因是Spark的分布式执行逻辑:
DataFrame.foreach()中传入的函数会被序列化后分发到各个Executor节点上执行,你在Driver端定义的lst变量只会被复制副本到Executor中,所有对lst的修改都是在Executor的副本上生效,不会同步回Driver端的原始lst变量,因此Driver端打印的lst始终是空的。 - 第二个原因是Python的变量作用域问题:你在
func2中写了lst = lst.append(x.firstname),Python会将函数内的lst识别为局部变量,而你在赋值前就尝试读取lst.append,哪怕是本地单进程执行这段代码也会抛出UnboundLocalError异常,另外list.append()方法返回值是None,赋值操作本身也没有意义。
可行解决方案
- 方案1:小数据量场景下直接将数据拉取到Driver端处理
如果你的DataFrame数据量不大,可以先通过collect()方法把所有行拉到Driver端,再遍历处理,代码示例:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName('SparkByExamples.com').getOrCreate() data = [('James','Smith','M',30),('Anna','Rose','F',41),('Robert','Williams','M',62), ] columns = ["firstname","lastname","gender","salary"] df = spark.createDataFrame(data=data, schema = columns) lst = [] # 直接在Driver端遍历 for row in df.collect(): lst.append(row.firstname) print(len(lst)) # 输出3
- 方案2:大数据量场景下使用Spark累加器收集结果
如果数据量大不能全量拉到Driver端,可以使用Spark提供的累加器组件,它支持Executor端分布式累加、Driver端读取最终结果,代码示例:
from pyspark.sql import SparkSession from pyspark import AccumulatorParam # 自定义列表类型累加器 class ListAccumulatorParam(AccumulatorParam): def zero(self, initialValue): return initialValue.copy() def addInPlace(self, v1, v2): v1.append(v2) return v1 spark = SparkSession.builder.appName('SparkByExamples.com').getOrCreate() data = [('James','Smith','M',30),('Anna','Rose','F',41),('Robert','Williams','M',62), ] columns = ["firstname","lastname","gender","salary"] df = spark.createDataFrame(data=data, schema = columns) # 初始化累加器 lst_acc = spark.sparkContext.accumulator([], ListAccumulatorParam()) def func2(x): global lst_acc lst_acc.add(x.firstname) df.foreach(func2) # Driver端读取累加器结果 lst = lst_acc.value print(len(lst)) # 输出3
- 方案3:仅提取单列值的场景直接用RDD算子简化处理
如果你的需求只是提取firstname列的所有值,可以直接用select加RDD算子实现,代码更简洁:
lst = df.select("firstname").rdd.map(lambda x: x[0]).collect() print(len(lst)) # 输出3
内容的提问来源于stack exchange,提问作者user3868051
相关产品推荐
相关产品推荐

