PySpark并行计算PI时如何去除calc_points_in_circ未使用的itr参数
解决方法
PySpark 的 filter 算子要求传入的处理函数必须接收 RDD 中的单个元素作为输入,所以原来的 itr 参数是算子调用时强制传入的,我们可以通过两种方式消除无用变量,且完全保留原有并行计算能力:
方案1:按Python惯例标记无用参数(修改成本最低)
Python 通用约定使用单下划线 _ 表示不需要用到的变量,我们只需要修改函数的参数名即可,其余逻辑完全不变:
from pyspark.context import SparkContext from pyspark.sql.session import SparkSession import random sc = SparkContext("local", "PI Approx") spark = SparkSession(sc) n = 1_000_000 # 仅将无用参数名改为 _,标识该参数无实际用途 def calc_points_in_circ(_): x = random.random() y = random.random() distance = (x*x + y*y) ** 0.5 return distance < 1 rdd = sc.parallelize(range(0, n)) pointsInCircle = rdd.filter(calc_points_in_circ).count() pi = 4 * pointsInCircle / n print(f"PI approximation: {pi:.3f}")
方案2:彻底移除函数的无用入参
如果不想在函数定义中保留无用参数,可以用 lambda 匿名函数承接 RDD 传入的元素并忽略,再调用无参的计算函数即可:
from pyspark.context import SparkContext from pyspark.sql.session import SparkSession import random sc = SparkContext("local", "PI Approx") spark = SparkSession(sc) n = 1_000_000 # 移除函数的无用入参 def calc_points_in_circ(): x = random.random() y = random.random() distance = (x*x + y*y) ** 0.5 return distance < 1 rdd = sc.parallelize(range(0, n)) # 用lambda接收RDD元素并忽略,调用无参的计算函数 pointsInCircle = rdd.filter(lambda _: calc_points_in_circ()).count() pi = 4 * pointsInCircle / n print(f"PI approximation: {pi:.3f}")
两种方案的执行效率和并行逻辑和原代码完全一致,没有额外性能损耗。
内容的提问来源于stack exchange,提问作者Afiq Rosli
相关产品推荐
相关产品推荐

