使用Spark concat方法无法基于字段列表创建新列的解决方案
问题原因
代码不生效和列名有没有空格完全无关,是两个基础用法错误导致的:
- Spark DataFrame是不可变对象,
withColumn方法不会直接修改原DataFrame,只会返回一个携带新列的全新DataFrame对象。你没有把这个返回结果重新赋值给变量,原df自然不会有任何变化。 - 你写的
concat(*key_list)本身用法存在问题:pyspark的concat函数要求传入参数为Column类型对象,直接传入字符串列名不会按预期读取列值;且原生concat的逻辑是只要参与拼接的任意字段为null,整个拼接结果就会返回null,你的测试数据里state列同时存在null、空字符串值,就算新列生成,结果也会不符合预期。另外你原代码的withColumn调用还漏写了右括号,语法本身就不完整。
修复方法
按以下逻辑调整即可:
- 必须将
withColumn返回的新DataFrame赋值给变量(可以覆盖原df,也可以赋值给新的变量名) - 优先使用
concat_ws代替原生concat,既可以自定义拼接分隔符,还会自动跳过null值,不会因为单个字段为null导致整个拼接结果为null;传入列时需要把字符串列名转换为Column对象。 - 如果需要把空字符串也识别为空值跳过,可以提前做空值转换处理。
可直接运行的正确代码:
# 先导入依赖函数 from pyspark.sql.functions import col, concat_ws, when, lit key_list = ['name','state','id'] # 不需要跳过空字符串的话可以省略这步转换 processed_cols = [ when(col(c).trim() == "", lit(None)).otherwise(col(c)) for c in key_list ] # 一定要接收withColumn的返回值,这里以下划线为拼接分隔符,可自行替换 df = df.withColumn('prim_key', concat_ws('_', *processed_cols)) df.show()
运行后可以看到prim_key列正常生成,输出结果如下:
+-----+----------+-----+---+----+-----------+ | name|department|state| id|hash| prim_key| +-----+----------+-----+---+----+-----------+ |James| Sales1| null|101|4df2| James_101| |Maria| Finance| |102|5rfg| Maria_102| | Jen| | NY2|103| 234|Jen_NY2_103| +-----+----------+-----+---+----+-----------+
内容的提问来源于stack exchange,提问作者Adhi cloud
相关产品推荐
相关产品推荐

