如何不使用pandas API将pandas loc逻辑改写为PySpark DataFrame实现
PySpark 实现pandas loc条件赋值方案
PySpark 本身不支持pandas 原地逐行修改的loc语法,且计算是按列操作顺序逐层迭代的,之前写不出预期结果,核心原因大概率是没有对齐pandas代码的顺序执行逻辑——pandas里每执行一行loc赋值,Total列的值就会立刻更新,后续的条件判断全部基于更新后的值,不能把所有条件平行放在同一层计算。
具体实现逻辑完全对齐你给出的pandas代码执行顺序:
- 第一步:计算初始
Total值,即四个level列的和除以Number列 - 第二步:判断如果
level4 > 0,当前行Total加4,否则保持原值 - 第三步:基于第二步更新后的
Total,判断如果level3 > 0且Total < 1,当前行Total加3,否则保持原值 - 第四步:基于第三步更新后的
Total,判断如果level2 > 0且Total < 1,当前行Total加2,否则保持原值 - 第五步:基于第四步更新后的
Total,判断如果level1 > 0且Total < 1,当前行Total加1,否则保持原值
完整可直接运行的代码如下:
# 先导入PySpark内置函数 from pyspark.sql import functions as F # 把df替换成你自己的DataFrame变量名即可 df = df.withColumn( "Total", (F.col("level1") + F.col("level2") + F.col("level3") + F.col("level4")) / F.col("Number") ).withColumn( "Total", F.when(F.col("level4") > 0, F.col("Total") + 4).otherwise(F.col("Total")) ).withColumn( "Total", F.when((F.col("level3") > 0) & (F.col("Total") < 1), F.col("Total") + 3).otherwise(F.col("Total")) ).withColumn( "Total", F.when((F.col("level2") > 0) & (F.col("Total") < 1), F.col("Total") + 2).otherwise(F.col("Total")) ).withColumn( "Total", F.when((F.col("level1") > 0) & (F.col("Total") < 1), F.col("Total") + 1).otherwise(F.col("Total")) )
关键注意点:链式调用withColumn时,后续步骤可以直接引用前序步骤刚更新的列值,刚好对齐pandas逐行顺序执行的效果。不要把所有条件判断塞到同一个withColumn里,这种写法里所有条件取到的都是最原始的初始Total值,结果会和预期偏差很大。
可以拿同一份测试数据分别跑pandas原版逻辑和上述PySpark代码,输出的Total列值完全一致。
内容的提问来源于stack exchange,提问作者IknewIt
相关产品推荐
相关产品推荐

