Databricks:高效实现含Null值的DataFrame多列拼接Python UDF
问题描述
我有多个包含Null值的DataFrame,示例如下:
DataFrame1(df1)
EmpId FirstName LastName MiddleName 1 Anna Walter I 2 Jack Shaun 3 Andrew Hill
需要将FirstName、LastName、MiddleName列的值拼接为:
Anna|Walter|I Jack||Shaun Andrew|Hill|
DataFrame2(df2)
DeptId DeptName Category Contact Building 10 Arts 1 James 7 20 Science 2 30 Social Kim 4
需要将DeptName、Category、Contact、Building列的值拼接为:
Arts|1|James|7 Science|2|| Social||Kim|4
我需要一个可按如下方式调用的自定义函数(UDF):
fields1 = F.array('FirstName', 'LastName', 'MiddleName') df1.withColumn('Combination1', udf_xyz(fields1)) fields2 = F.array('DeptName', 'Category', 'Contact', 'Building') df2.withColumn('Combination2', udf_xyz(fields2))
我尝试过以下函数,但先将空列填充为''再还原的方式在处理百万级数据时开销极大:
def concatdf_ws(df, colName, fields): df = df.fillna('') df = df.withColumn(colName, concat_ws('|', fields)) df=df.select(*[when(trim(df[x])=='',None).otherwise(df[x]).alias(x) for x in df.columns]) df.display() return df;
求更优的实现方式。
最优实现方案
不需要全局修改原DataFrame的Null值,直接通过Spark原生函数或轻量UDF即可高效完成需求:
方案1:纯原生API实现(推荐,性能最优)
用transform将数组内的Null值转为空字符串,再用concat_ws拼接,全程分布式优化执行,无需修改原表数据:
from pyspark.sql import functions as F def concat_with_nulls(array_col, sep='|'): # 将数组中的Null转换为空字符串 transformed_arr = F.transform(array_col, lambda x: F.coalesce(x, F.lit(''))) # 按指定分隔符拼接 return F.concat_ws(sep, transformed_arr) # 调用方式完全符合你的要求 fields1 = F.array('FirstName', 'LastName', 'MiddleName') df1 = df1.withColumn('Combination1', concat_with_nulls(fields1)) fields2 = F.array('DeptName', 'Category', 'Contact', 'Building') df2 = df2.withColumn('Combination2', concat_with_nulls(fields2))
方案2:轻量UDF实现(如果必须用UDF形式)
直接在UDF内部处理Null值,仅针对目标列操作,避免全局扫描修改:
from pyspark.sql import functions as F from pyspark.sql.types import StringType @F.udf(StringType()) def udf_xyz(arr): # 遍历数组,将Null转为空字符串后拼接 return '|'.join([str(x) if x is not None else '' for x in arr]) # 调用方式与你的要求完全一致 fields1 = F.array('FirstName', 'LastName', 'MiddleName') df1 = df1.withColumn('Combination1', udf_xyz(fields1)) fields2 = F.array('DeptName', 'Category', 'Contact', 'Building') df2 = df2.withColumn('Combination2', udf_xyz(fields2))
性能优势说明
- 原方法的全局
fillna+还原操作会扫描并修改全表所有列,大数据量下IO和计算开销极大; - 方案1基于Spark原生优化算子,仅处理目标数组列,无需修改原表其他数据,分布式执行效率最高;
- 方案2的UDF仅在拼接阶段处理目标列的Null值,不会全局修改原数据,开销远低于原方法。
内容的提问来源于stack exchange,提问作者Nanda
相关产品推荐
相关产品推荐

