Spark懒执行是否会导致循环中quantile变量覆盖,计算结果错误?
问题解答
核心结论
不会出现你担心的「除最后一次迭代外,其他计算结果被覆盖」的问题,quantile_foo和quantile_bar能正确传递到对应循环迭代的计算中,但你的代码存在两处明显错误需要修正。
具体分析
approxQuantile是立即执行的Action
Spark的懒执行仅针对转换操作(比如withColumn、filter),但approxQuantile是**动作(Action)**方法——调用它时Spark会立刻触发Job计算分位数值,并将结果以Python列表的形式返回给本地变量。也就是说,每一轮循环里的quantile_foo和quantile_bar,都是当前迭代对应列(foo_i/bar_i)的真实分位数值,会被即时确定并保存,不会等到后续迭代才解析。懒执行不会导致变量引用被覆盖
虽然withColumn是懒执行的转换,但在调用withColumn生成新列时,quantile_foo[0]、quantile_foo[1]这些都是Python本地常量,会直接嵌入到Spark的执行计划里。比如i=0时,quantile_foo已经是foo_0的分位数,生成foo_quantile_0的表达式时用的就是这个具体数值,而非指向quantile_foo变量的引用。后续循环修改quantile_foo的值,完全不会影响之前已经生成的列逻辑。
代码中的错误修正
你的代码存在两处致命错误,和懒执行无关,但会导致计算完全错误:
- 笔误用错分位变量:在计算
bar_quantile_{i}时,你错误地使用了quantile_foo的数值(比如quantile_foo[0]),应该替换为quantile_bar,否则bar列的分位逻辑会用foo的分位数,结果完全不对。 when语法错误:F.col("foo_{}".format(str(i))>quantile_foo[0]的括号位置错误,正确写法是F.col(f"foo_{i}") > quantile_foo[0],否则会把整个比较表达式当成列名,触发语法异常。
修正后的代码片段示例:
import pyspark.sql.functions as F df = spark.sql("select * from some_table") for i in range(5): quantile_foo = df.approxQuantile(f"foo_{i}", [0.25, 0.5, 0.75], 0.05) quantile_bar = df.approxQuantile(f"bar_{i}", [0.25, 0.5, 0.75], 0.05) # 修正foo_quantile的语法和逻辑 df = df.withColumn(f"foo_quantile_{i}", F.when(F.col(f"foo_{i}") > quantile_foo[2], 100) # 补充75分位以上的逻辑 .when(F.col(f"foo_{i}") > quantile_foo[1], 75) .when(F.col(f"foo_{i}") > quantile_foo[0], 50) .otherwise(25) ) # 修正bar_quantile的变量引用和语法 df = df.withColumn(f"bar_quantile_{i}", F.when(F.col(f"bar_{i}") > quantile_bar[2], 100) .when(F.col(f"bar_{i}") > quantile_bar[1], 75) .when(F.col(f"bar_{i}") > quantile_bar[0], 50) .otherwise(25) )
内容的提问来源于stack exchange,提问作者Shengxin Huang
相关产品推荐
相关产品推荐

