如何使用PySpark将无引号JSON转换为CSV或标准JSON格式
PySpark转换无引号类JSON字符串的实现方案
方法1:正则替换转标准JSON
适用场景:字段值中不包含=、,、{、}等特殊字符,执行效率高于自定义UDF。
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 假设原始无引号JSON存储在raw_str列 df = df.withColumn("standard_json", # 链式正则替换修正格式 F.regexp_replace( F.regexp_replace( F.regexp_replace( F.regexp_replace("raw_str", r"\{", '{"'), r"=", '":"' ), r",", '","' ), r"\}", '"}' ) ) # 定义字段schema解析JSON为结构化数据 json_schema = StructType([ StructField("PlatformVersion", IntegerType(), nullable=True), StructField("PlatformClient", StringType(), nullable=True), StructField("namespace", StringType(), nullable=True) ]) df = df.withColumn("parsed_data", F.from_json("standard_json", json_schema))
方法2:自定义UDF解析(兼容性更强)
适用场景:字段值可能包含=等特殊字符,可灵活适配不同格式的异常字符串。
from pyspark.sql import functions as F from pyspark.sql.functions import udf from pyspark.sql.types import MapType, StringType # 自定义解析函数 def parse_unquoted_str(raw_str): raw_str = raw_str.strip("{}") kv_list = raw_str.split(",") res = {} for kv in kv_list: # 仅分割第一个等号,避免值中包含等号导致解析错误 k, v = kv.split("=", maxsplit=1) res[k.strip()] = v.strip() return res parse_udf = udf(parse_unquoted_str, MapType(StringType(), StringType())) df = df.withColumn("parsed_map", parse_udf("raw_str"))
导出为CSV或直接写入关系库
解析得到结构化数据后,可以直接导出为CSV,也可以跳过文件导出步骤直接写入关系型数据库:
# 展开字段为独立列 df_final = df.select( F.col("parsed_map.PlatformVersion").alias("platform_version"), F.col("parsed_map.PlatformClient").alias("platform_client"), F.col("parsed_map.namespace").alias("namespace") ) # 导出为CSV df_final.write.csv("./output_csv", header=True, mode="overwrite") # 直接写入MySQL等关系库(无需中转CSV/JSON文件) df_final.write.format("jdbc") \ .option("url", "jdbc:mysql://数据库地址:端口/库名") \ .option("dbtable", "表名") \ .option("user", "账号") \ .option("password", "密码") \ .mode("append") \ .save()
注意:如果原始字符串的值中包含逗号、大括号等特殊分割符,需要根据实际日志格式调整解析逻辑,避免字段错位。
内容的提问来源于stack exchange,提问作者Aniruddh Kulkarni
相关产品推荐
相关产品推荐

