Spark Scala如何生成存储取值为true的列名的ArrayType类型字段
Spark 生成值为true的列名数组最优实现
这里优先推荐全内置函数的实现方式,无需自定义UDF,可充分利用Spark Catalyst优化器的能力,性能远高于UDF实现。
核心思路
遍历所有需要校验的布尔列,对每列做判断:如果列值为true则返回列名,否则返回null,最后移除数组中的所有null值,即可得到目标字段。
Scala 实现(Spark 2.4+)
import org.apache.spark.sql.functions.{col, array, array_remove, lit, when} // 定义需要校验的布尔列列表 val boolColumns = Seq("col1", "col2", "col3") val resultDF = inputDF.withColumn("colMap", array_remove( array(boolColumns.map(c => when(col(c) === true, lit(c)).otherwise(null)): _*), null ) )
如果是低于2.4版本的Spark,可使用filter函数替代array_remove:
import org.apache.spark.sql.functions.{col, array, filter, lit, when} val boolColumns = Seq("col1", "col2", "col3") val resultDF = inputDF.withColumn("colMap", filter( array(boolColumns.map(c => when(col(c) === true, lit(c)).otherwise(null)): _*), x => x.isNotNull ) )
PySpark 实现(Spark 2.4+)
from pyspark.sql.functions import col, array, array_remove, lit, when # 定义需要校验的布尔列列表 bool_columns = ["col1", "col2", "col3"] result_df = input_df.withColumn("colMap", array_remove( array(*[when(col(c) == True, lit(c)).otherwise(None) for c in bool_columns]), None ) )
效果验证
执行后得到的colMap字段完全符合预期:
- 当行的col1、col3为true时,返回
[col1, col3] - 当行所有列都为false时,返回空数组
[] - 当行仅col3为true时,返回
[col3]
内容的提问来源于stack exchange,提问作者Dariusz Krynicki
相关产品推荐
相关产品推荐

