读取大JSON文件至PySpark DataFrame前的验证与修复方法咨询
问题解答
先修复损坏行再读取到PySpark DataFrame是否可行?
完全可行,而且对于1500行的小文件来说,这种方案高效且易实现。具体流程可以是:
- 用Python本地文件操作逐行读取原文件,针对每行的损坏JSON做修复(比如补全缺失的括号、移除多余的闭合符号)
- 将修复后的完整JSON行写入临时文件
- 最后用PySpark的
spark.read.json()读取这个修复好的文件
1500行的文件量级很小,本地处理的开销几乎可以忽略,不会对后续PySpark的处理效率造成影响。
简单正则替换是否足以移除JSON对象中间的意外换行?
分两种场景判断:
- 如果意外换行是JSON字符串内部的无意义换行(比如
"place_name":"Chi\ncago"这类情况),可以用正则精准替换掉引号内的换行,比如通过匹配引号间的内容并移除其中的换行符来处理 - 如果换行导致JSON结构断裂(比如一个完整的JSON对象被拆成了多行),单纯正则替换就不够了。这种情况需要先通过统计括号的开闭数量来判断行的完整性,比如当某行的左括号数不等于右括号数时,继续读取下一行拼接,直到括号数量平衡,再做后续修复
参考代码片段
import re def fix_single_json_line(line): # 移除行尾多余的冗余闭合括号(如示例中的...}) cleaned = re.sub(r'}\s*$', '}', line.rstrip()) # 清除字符串内部的意外换行 cleaned = re.sub(r'(?<=")(.*?)\n(.*?)(?=")', r'\1\2', cleaned) # 补全缺失的闭合括号 open_count = cleaned.count('{') close_count = cleaned.count('}') if open_count > close_count: cleaned += '}' * (open_count - close_count) return cleaned # 批量修复文件 with open('damaged_data.json', 'r') as in_f, open('fixed_data.json', 'w') as out_f: for raw_line in in_f: fixed_line = fix_single_json_line(raw_line) out_f.write(fixed_line + '\n') # 读取到PySpark DataFrame from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ReadFixedJSON").getOrCreate() df = spark.read.json("fixed_data.json") df.show()
内容的提问来源于stack exchange,提问作者Boris
相关产品推荐
相关产品推荐

