Spark DataFrame多withColumn操作如何设置全局行过滤条件
方案1:抽取公共条件 + 批量生成列
把重复的判断逻辑提取成公共变量,再通过批量操作新增列,避免每次手写重复条件:
import org.apache.spark.sql.functions._ // 抽取重复的判断条件,仅需定义一次 val matchCondition = $"chaine" === "TF1" && $"cond" === 1 // 定义新列与对应取值的映射关系,新增列仅需修改该序列即可 val columnMapping = Seq( ("New", "YES!"), ("New2", "YES2!"), ("New3", "YES3!"), ("New4", "YES4!") ) // 批量新增所有列 val resultDF = columnMapping.foldLeft(df) { case (currentDF, (colName, colVal)) => currentDF.withColumn(colName, when(matchCondition, colVal)) }
该方案改动最小,适合新增列的命名和取值有统一规律的场景,后续新增列仅需要更新columnMapping序列即可。
方案2:拆分数据集再合并
将数据集按cond拆分后分别处理,再合并结果,完全不用在列判断中重复加条件:
import org.apache.spark.sql.functions._ // 处理cond=1的行,直接按原有逻辑新增列 val dfCond1 = df.filter($"cond" === 1) .withColumn("New", when($"chaine" === "TF1", "YES!")) .withColumn("New2", when($"chaine" === "TF1", "YES2!")) .withColumn("New3", when($"chaine" === "TF1", "YES3!")) .withColumn("New4", when($"chaine" === "TF1", "YES4!")) // 处理cond≠1的行,直接补全对应新列即可 val dfCondOther = df.filter($"cond" =!= 1) .withColumn("New", lit(null).cast("string")) .withColumn("New2", lit(null).cast("string")) .withColumn("New3", lit(null).cast("string")) .withColumn("New4", lit(null).cast("string")) // 合并两份数据得到全量结果 val resultDF = dfCond1.unionByName(dfCondOther)
该方案适合cond=1的行处理逻辑非常复杂的场景,逻辑拆分更清晰,Spark优化器会自动优化执行计划,不会有额外性能开销。
内容的提问来源于stack exchange,提问作者Haha
相关产品推荐
相关产品推荐

