基于DataFrame现有列生成errorColumn的高效实现方案问询
问题描述
输入DataFrame包含sequence、registerNumber、first_name、middle_name、surname列,数据如下:
+--------+--------------+----------+-----------+-------+ |sequence|registerNumber|first_name|middle_name|surname| +--------+--------------+----------+-----------+-------+ | 1| XXXXXXXX| | | | | 1| XXXXXXXX| | | | | 2| YYYYYYYY| CAR弟LINE| ELIZABET弟| | | 3| YYYYYYYZ| | | GAL弟| | 4| YYYYYYYM| 弟AROLINE| | | | 5| YYYYYYYL| | | | +--------+--------------+----------+-----------+-------+
需新增errorColumn列,规则为:记录first_name、middle_name、surname中非空(含有效内容)的列名,多列时用-分隔;三列均为空则errorColumn留空。此前用concat实现性能较差,寻求更优方案。
期望输出:
+--------+--------------+----------+-----------+-------+-------------------+ |sequence|registerNumber|first_name|middle_name|surname|errorColumn | +--------+--------------+----------+-----------+-------+-------------------+ | 1| XXXXXXXX| | | | | | 1| XXXXXXXX| | | | | | 2| YYYYYYYY| CAR弟LINE| ELIZABET弟| |first_name-middle_name| | 3| YYYYYYYZ| | | GAL弟|surname | | 4| YYYYYYYM| 弟AROLINE| | |first_name | | 5| YYYYYYYL| | | | | +--------+--------------+----------+-----------+-------+-------------------+
优化方案
避免多次concat拼接(会生成大量中间对象,拖累性能),改用数组构造+过滤+单次字符串拼接的组合,Spark底层对数组操作的优化更充分,能显著提升处理效率。
方法1:DataFrame API(Python)
from pyspark.sql import functions as F # 定义需要检查的目标列 check_cols = ["first_name", "middle_name", "surname"] # 构造数组:列非空时保留列名,否则设为null col_array = F.array( *[F.when(F.trim(F.col(col)) != "", col).otherwise(F.lit(None)) for col in check_cols] ) # 过滤数组中的null值,再用"-"拼接剩余元素 df_result = df.withColumn( "errorColumn", F.concat_ws("-", F.filter(col_array, lambda x: x.isNotNull())) ) df_result.show(truncate=False)
方法2:DataFrame API(Scala)
import org.apache.spark.sql.functions._ val checkCols = Seq("first_name", "middle_name", "surname") val colArray = array( checkCols.map(col => when(trim(col) =!= "", lit(col)).otherwise(null)): _* ) val dfResult = df.withColumn( "errorColumn", concat_ws("-", filter(colArray, x => x.isNotNull)) ) dfResult.show(false)
方法3:Spark SQL
若偏好SQL语法,可注册临时视图后执行:
CREATE OR REPLACE TEMP VIEW temp_data AS SELECT * FROM your_input_df; SELECT sequence, registerNumber, first_name, middle_name, surname, concat_ws( "-", CASE WHEN trim(first_name) != "" THEN "first_name" END, CASE WHEN trim(middle_name) != "" THEN "middle_name" END, CASE WHEN trim(surname) != "" THEN "surname" END ) AS errorColumn FROM temp_data;
性能说明
- 上述方案通过
filter+concat_ws替代多次concat,减少了中间字符串对象的生成,Spark Catalyst优化器可对数组操作做深度优化,在大数据量场景下性能提升明显。 - 代码中用
trim是为了排除仅含空格的无效空值,若业务中空字符串(无空格)才算空,可直接去掉trim。
内容的提问来源于stack exchange,提问作者Muskan Makhija
相关产品推荐
相关产品推荐

