You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 13:48:19