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

如何编写PySpark UDF解析混合竖线分隔的复杂列并拆分为多列

PySpark 非标准混合分隔符列拆分实现

问题特征

  • 原始列取值为非标准JSON格式,样例:{"1":"mediaMaaadadeftch||OAISAOID|true|ModsVersio|67900|clk|true|PPOOOS|20220501164113|34958|38177557..}
  • 值部分混用||、|两种竖线作为分隔符
  • 拆分后需按顺序生成col_1、col_2、col_3……列,第一列对应JSON的键值,后续列按顺序对应分隔后的值内容,第一行预期取值依次为1、mediaMaaadadeftch、OAISAOID、true……

核心处理逻辑

  • 先清洗原始字符串,去除首尾包裹的大括号,按第一个冒号切分,分离键部分和值部分,避免值内包含冒号时误拆分
  • 统一分隔符:将值部分所有连续双竖线||替换为单竖线|,再按单竖线切分值
  • 将键放在拆分结果列表的首位,和值拆分结果拼接为完整的字段序列
  • 增加异常兜底逻辑,格式错误、空值的行直接返回空数组,避免Spark任务中断
  • 基于拆分后的数组列,动态提取元素生成独立列

完整代码实现

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf, col, size, max as spark_max
from pyspark.sql.types import ArrayType, StringType

# 初始化Spark会话,已有运行环境可跳过
spark = SparkSession.builder.appName("split_raw_column").getOrCreate()

# 定义拆分UDF
@udf(returnType=ArrayType(StringType()))
def split_raw_col(input_str):
    if not input_str:
        return []
    try:
        # 去除首尾大括号
        content = input_str.strip("{}")
        # 按第一个冒号切分键和值
        key, raw_value = content.split(":", 1)
        # 去除键和值外层包裹的双引号
        key = key.strip('"')
        raw_value = raw_value.strip('"')
        # 统一分隔符后拆分
        value_list = raw_value.replace("||", "|").split("|")
        # 拼接键和值列表,过滤空内容
        full_list = [key] + value_list
        return [item.strip() for item in full_list if item.strip()]
    except:
        return []

# ---------------------- 调用示例 ----------------------
# 假设原始DataFrame为source_df,待拆分列名为raw_data
# 1. 先生成拆分后的数组列
df_split = source_df.withColumn("field_arr", split_raw_col(col("raw_data")))

# 2. 计算拆分后的最大列数,动态生成列
max_field_count = df_split.agg(size(spark_max("field_arr"))).first()[0]
final_df = df_split.select(
    *[col("field_arr")[idx].alias(f"col_{idx+1}") for idx in range(max_field_count)]
)

# 查看拆分结果
final_df.show(truncate=False)

使用说明

  • 如果提前明确拆分后的固定列数,可以跳过最大列数计算步骤,直接按固定索引提取列,运行性能更优
  • 如果原始值末尾存在无意义截断字符(比如样例末尾的..),可以在UDF返回结果前增加对应字符的过滤、替换逻辑
  • 异常捕获逻辑会将格式不符合预期的行所有拆分列置为null,不会导致批量任务失败

内容的提问来源于stack exchange,提问作者Sanjay Nirmal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 07:42:27