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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 08:14:59