如何编写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
相关产品推荐
相关产品推荐

