Spark中withColumn是否会将UDF结果拉取到Driver及相关操作疑问
PySpark UDF与数据操作疑问解答
执行代码
df.repartition(16).withColumn( 'data', explode(udf('col1','col2')) ).write.mode('overwrite').parquet('result.parquet')
问题解答
上述操作中是否会将UDF的结果传输到Driver节点?
不会。UDF的计算逻辑在Executor节点上分布式执行,计算后的结果直接在Executor端参与后续的explode和Parquet写入流程,全程不会将数据传输回Driver节点——只有collect()、take()这类主动拉取数据到Driver的操作才会触发这类传输,这段代码里没有这类动作。上述哪些操作涉及数据传输?
repartition(16):该操作会触发数据shuffle,需要在不同Executor节点之间传输数据,以此重新分配数据分区,是典型的跨节点数据传输场景。- Parquet写入操作:如果目标存储是分布式存储(如HDFS),Executor会将各自分区的数据写入到存储系统中;如果是本地存储,每个Executor也会将自己处理的分区数据写入本地文件,这也属于数据从Executor到存储介质的传输。其中核心的跨节点数据传输来自
repartition的shuffle过程。
- 执行器节点是否会各自保存独立的Parquet文件?
是的。repartition(16)将DataFrame划分为16个分区,每个分区对应一个独立的任务由Executor执行,每个任务执行完成后,会将对应分区的数据写入一个独立的Parquet文件(文件命名通常为part-xxxxxx开头的格式)。最终输出目录下会生成与分区数对应的多个Parquet文件。
内容的提问来源于stack exchange,提问作者Doofenshmirtz
相关产品推荐
相关产品推荐

