PySpark字典查询执行节点及广播变量使用合理性疑问
PySpark字典查询节点问题解答
问题1:PySpark代码中查询字典时,操作在集群还是主节点执行?
要看字典定义和使用的位置:
- 如果是在主节点(Driver端)定义的字典,且直接在主节点代码里查询(没用到RDD/DataFrame的map、filter这类分布式算子),查询就在主节点执行。
- 如果是把字典放到集群节点(Executor端)的任务逻辑里用(比如
df.map(lambda x: my_dict[x])这种算子里查询),那查询是在集群的Executor节点执行。但要注意:直接把普通字典传到算子里,Spark会给每个任务序列化一份字典副本,效率很低。
问题2:collect()后创建字典并查询,操作在集群还是主节点执行?
df.collect()会把集群上的DataFrame数据拉到主节点(Driver端),之后用Python创建的字典是主节点本地对象,查询操作肯定在主节点执行。
关于广播变量意义的疑问解答
你的理解有偏差,广播变量的价值是解决Executor端需要重复使用大字典/大对象的场景:
- 比如你有个1GB的字典,要在分布式算子(比如join、map)里匹配数据。不用广播变量的话,每个任务都会传一份字典,100个任务就会传100GB数据,还占Executor大量内存。
- 用广播变量的话,Spark只会把字典传到每个Executor节点一次,该节点上的所有任务共享这一份副本,能大幅减少数据传输量和内存占用。
- 要是字典只是在主节点本地用,确实不需要广播变量——这和广播变量的适用场景完全不冲突。
内容的提问来源于stack exchange,提问作者Paawan Angra
相关产品推荐
相关产品推荐

