PySpark DataFrame调用foreach的lambda可打印但无法追加元素至列表
解决PySpark中
foreach无法向驱动端列表添加元素的问题 问题原因
PySpark是分布式计算框架,foreach方法中的lambda逻辑是在Executor节点上执行的:
- 驱动端(Driver)定义的
obj列表仅存在于Driver进程的内存中 - 每个Executor会复制一份
obj的副本,lambda里的append操作只修改了Executor本地的副本,这些修改不会同步回Driver端的原始列表,所以最终打印的obj还是空列表。
解决方案
如果需要将DataFrame中的数据收集到驱动端的列表,直接使用collect()方法即可——它会把所有Executor上的数据拉取到Driver端,此时在Driver端操作列表就能得到预期结果:
修改后的代码
from pyspark.sql.types import Row from pyspark.sql import SparkSession spark = SparkSession.builder\ .master("local[*]")\ .appName("ETL")\ .config("spark.executor.logs.rolling.time.interval", "daily")\ .getOrCreate() sample_streamed_data = spark.createDataFrame([ {"value": Row(apple='12', banana='14', carrot='0')}, {"value": Row(apple='10', banana='10', carrot='2')} ]) print(f"{sample_streamed_data} --> {type(sample_streamed_data)}") # 从DataFrame拉取数据到驱动端并转换为目标列表 obj = [row.asDict().get('value') for row in sample_streamed_data.collect()] print(obj) # 验证输出 for item in obj: print(item)
执行输出
DataFrame[value: struct<apple:string,banana:string,carrot:string>] --> <class 'pyspark.sql.dataframe.DataFrame'> [Row(apple='12', banana='14', carrot='0'), Row(apple='10', banana='10', carrot='2')] Row(apple='12', banana='14', carrot='0') Row(apple='10', banana='10', carrot='2')
注意事项
如果DataFrame数据量极大,collect()会将所有数据加载到Driver端内存,可能引发内存溢出。这种场景下应该避免将数据拉取到Driver端,而是在分布式环境中完成数据处理逻辑。
内容的提问来源于stack exchange,提问作者Manoj
相关产品推荐
相关产品推荐

