如何用PySpark读取含Base64编码及换行的非标准JSON文件?
解决非标准多行JSON文件的Spark读取问题
你的问题核心是文件里的顶级JSON对象可能跨多行(因为Body字段嵌了带换行的JSON),没法直接用spark.read.json或按行分割读取。下面给两种靠谱的解决思路:
方法一:按大括号层级分割(无额外依赖)
每个顶级JSON都是{开头、}结尾,内嵌的大括号会增加嵌套层级,我们可以通过计数括号的方式精准分割出每个完整的顶级对象,同时处理字符串里的大括号避免误判:
代码示例(Python):
from pyspark.sql import SparkSession import json import base64 spark = SparkSession.builder.appName("FixNonStandardJson").getOrCreate() # 读取目标文件 file_rdd = spark.sparkContext.wholeTextFiles("/path/to/your/data.json") def split_top_level_jsons(text): jsons_list = [] current_chunk = [] bracket_count = 0 in_string = False escape_flag = False for char in text: current_chunk.append(char) # 处理转义字符,跳过后续判断 if escape_flag: escape_flag = False continue # 切换字符串状态:遇到双引号且不是转义的,进出字符串 if char == '"': in_string = not in_string # 标记转义符 elif char == '\\': escape_flag = True # 不在字符串里才统计括号 elif not in_string: if char == '{': bracket_count += 1 elif char == '}': bracket_count -= 1 # 括号计数归0,说明一个顶级JSON结束 if bracket_count == 0: jsons_list.append(''.join(current_chunk).strip()) current_chunk = [] return jsons_list # 分割出所有完整的顶级JSON字符串 json_strings_rdd = file_rdd.flatMap(lambda x: split_top_level_jsons(x[1])) # 解析JSON并处理Body字段 def parse_and_process_body(json_str): data = json.loads(json_str) body = data.get("Body") # 判断Body是Base64字符串还是内嵌JSON if isinstance(body, str): try: # 解码Base64 decoded = base64.b64decode(body).decode("utf-8") # 尝试把解码后的内容转成JSON(如果需要的话) try: data["Body"] = json.loads(decoded) except json.JSONDecodeError: # 不是JSON就保留解码后的字符串 data["Body"] = decoded except: # 解码失败就留原内容 pass # 内嵌JSON的情况直接保留 return data # 转成DataFrame result_df = json_strings_rdd.map(parse_and_process_body).toDF() # 查看结果 result_df.show(truncate=False)
方法二:用流式解析库处理(适合超大文件)
如果文件特别大,用wholeTextFiles加载到内存压力大,可以用ijson这个流式JSON解析库,它能逐块解析,不需要加载整个文件。
代码示例(Python):
先确保每个节点安装了ijson:pip install ijson
from pyspark.sql import SparkSession import ijson import base64 import json spark = SparkSession.builder.appName("StreamParseJson").getOrCreate() # 读取文件内容,按分区处理 file_rdd = spark.sparkContext.wholeTextFiles("/path/to/your/data.json").values() def parse_json_partition(iter_text): for text in iter_text: # 流式解析顶级JSON对象 parser = ijson.parse(text) current_obj = {} path_stack = [] for prefix, event, value in parser: if event == "start_map": if not path_stack: current_obj = {} path_stack.append("map") elif event == "end_map": path_stack.pop() if not path_stack: # 处理Body字段 body = current_obj.get("Body") if isinstance(body, str): try: decoded = base64.b64decode(body).decode("utf-8") try: current_obj["Body"] = json.loads(decoded) except json.JSONDecodeError: current_obj["Body"] = decoded except: pass yield current_obj elif event == "start_array": path_stack.append("array") elif event == "end_array": path_stack.pop() elif event in ["string", "number", "boolean", "null"]: # 解析属性路径,比如Properties.connectionDeviceId keys = prefix.split(".") temp = current_obj for key in keys[:-1]: temp = temp[key] temp[keys[-1]] = value # 解析并转成DataFrame result_df = file_rdd.mapPartitions(parse_json_partition).toDF() result_df.show(truncate=False)
两种方法对比:
- 方法一:无额外依赖,实现简单,适合中小文件
- 方法二:流式解析,内存占用低,适合超大文件,但需要安装第三方库
内容的提问来源于stack exchange,提问作者WilliamEllisWebb
相关产品推荐
相关产品推荐

