PySpark使用内置函数按同列item值拆分path列替代报错UDF
PySpark 按指定列值拆分字符串列实现方案
问题说明
现有结构如下的DataFrame:
+----+-------+ |item| path| +----+-------+ | -a-| a-b-c| | -b-| e-b-f| | -d-|e-b-d-h| | -c-| g-h-c| +----+-------+
需求为以同一行item列的取值为分隔符,拆分对应行的path列,过滤无法拆分的无效行,最终得到如下结果:
+----+--------+ |item| path| +----+--------+ | -b-| [a, c]| | -b-| [e, f]| | -d-|[e-b, h]| | -c-|[g-h, b]| +----+--------+
此前使用Python UDF实现拆分逻辑:
split_udf = udf(lambda a,b: a.split(b),T.ArrayType(T.StringType())) org = org.withColumn('crb_url', split_udf('path','item')[0])
小批量测试时运行正常,但在做DataFrame关联、写入Delta表操作时抛出错误:
AttributeError: 'NoneType' object has no attribute 'split'
报错原因是UDF未做空值兼容,且Python UDF本身性能较差,在分布式运行场景下稳定性不足。
实现方案
直接使用Spark内置split函数替代自定义UDF即可,天然兼容空值场景,运行性能远高于Python UDF,核心代码如下:
from pyspark.sql import functions as F result = org.filter( # 提前过滤item、path为null的行,避免无效计算 F.col("item").isNotNull() & F.col("path").isNotNull() ).withColumn( # 直接以每行item值为分隔符拆分path,完全等价于原UDF的a.split(b)逻辑 "path", F.split(F.col("path"), F.col("item")) ).filter( # 过滤拆分后数组长度不足2的行,即path中无对应item分隔串的无效行 F.size("path") == 2 ).select("item", "path")
注意事项
如果item列的取值包含正则特殊字符(如.、*、+等),可以将分隔符包裹为正则字面量,避免拆分逻辑不符合预期,拆分代码修改为:
"path", F.split( F.col("path"), F.concat(F.lit("\\Q"), F.col("item"), F.lit("\\E")) )
注:示例中
-c-行的输出存在笔误,按给出的原始数据path=g-h-c、item=-c-拆分后结果为["g-h", ""],如果需要得到[g-h, b]请核对原始数据的path取值,核心拆分逻辑不受影响。
逻辑说明
- 内置
split支持接收动态列作为分隔符参数,完全覆盖原UDF的拆分能力,不需要额外开发自定义逻辑 - 内置函数对null输入直接返回null,不会触发Python层面的NoneType方法调用异常
- 内置函数为Spark原生优化实现,大数据量场景下性能比Python UDF高3~10倍,关联、写入Delta表时稳定性更高
内容的提问来源于stack exchange,提问作者peer wild
相关产品推荐
相关产品推荐

