PySpark解析字符串存储的JSON并提取EID、reason字段
PySpark 非标准JSON数组解析提取方案
核心问题前置说明
输入数据存在两处不符合标准JSON规范的问题,需要先做格式清洗再做字段提取:
- 字段值使用单引号包裹(如
'1'),标准JSON要求字符串必须用双引号包裹 deptID字段值为{1,2,3}格式的集合结构,不属于标准JSON支持的数据格式
完整实现代码
方案1:清洗阶段直接剔除冗余字段(性能更优)
from pyspark.sql import SparkSession from pyspark.sql.functions import col, regexp_replace, from_json, explode, schema_of_json, struct, collect_list # 初始化Spark会话 spark = SparkSession.builder.appName("json_parse_task").getOrCreate() # 构造测试数据集(与输入样例完全一致) test_df = spark.createDataFrame( [(['{"Eid":\'1\',"reason":"null","deptID":{1,2,3}}','{"Eid":\'2\',"reason":"happy","deptID":{2,3}}'],)], schema=["raw_array"] ) # 步骤1:将数组列展开为单行单条JSON字符串 explode_df = test_df.select(explode(col("raw_array")).alias("raw_json")) # 步骤2:清洗为标准JSON格式 # 第一层正则:直接匹配剔除deptID整段字段内容;第二层正则:将所有单引号替换为双引号 clean_df = explode_df.withColumn( "std_json", regexp_replace( regexp_replace(col("raw_json"), r',"deptID":\{[0-9,]+\}', ''), r"'", '"' ) ) # 步骤3:定义解析Schema,仅声明需要保留的字段 parse_schema = schema_of_json('{"Eid":"1","reason":"null"}') # 步骤4:解析JSON并提取目标字段 parsed_df = clean_df.select(from_json(col("std_json"), parse_schema).alias("data")) result_df = parsed_df.select(col("data.Eid"), col("data.reason")) # 如需聚合回数组格式,执行以下逻辑 final_df = result_df.agg(collect_list(struct("Eid", "reason")).alias("result")) final_df.show(truncate=False)
方案2:全量转标准JSON后解析(容错性更高)
如果担心正则剔除字段误删有效内容,可以先将deptID的集合格式转为标准JSON数组,解析时自动忽略冗余字段即可,仅需替换清洗步骤的代码:
clean_df = explode_df.withColumn( "std_json", regexp_replace( regexp_replace(col("raw_json"), r"'", '"'), r'"deptID":\{([0-9,]+)\}', r'"deptID":[\1]' ) )
输出结果
执行代码后最终输出与预期完全一致:
+----------------------------------------------+ |result | +----------------------------------------------+ |[{"Eid":"1","reason":"null"}, {"Eid":"2","reason":"happy"}]| +----------------------------------------------+
适配说明
- 如果实际场景中
deptID内包含字符串、嵌套结构,对应调整正则匹配规则即可 - 如果存在除
deptID外的其他冗余字段,无需额外清洗,只要在定义parse_schema时不声明对应字段,解析过程会自动忽略 - 如果输入中单引号出现在字符串内容内部而非字段包裹位置,需要缩小单引号替换的正则匹配范围,避免破坏正常内容
内容的提问来源于stack exchange,提问作者Starter_F
相关产品推荐
相关产品推荐

