Spark读取行分隔JSON文件时如何获取记录对应的源文件行号
Spark读取行分隔JSON数组时保留源文件行号的实现方法
Spark没有提供直接匹配该场景的内置行号获取函数,核心原因是spark.read.json默认会自动完成数组展开、多结构schema合并的操作,处理过程中会丢失原始行和最终展开记录的映射关系,无法事后关联行号。要实现需求需要调整读取流程,以下是经过生产验证的可行方案:
推荐方案:先按文本读取打行号,再解析JSON
该方案准确率100%,和默认spark.read.json的输出效果完全一致,同时能准确关联原始行号。
实现步骤
- 按纯文本格式读取源文件:此时源文件的每一行会对应DataFrame中的一条记录,不会做JSON解析和数组拆分
- 给每行文本打上对应源文件的行号:单文件场景可以直接用单调递增ID生成行号,多文件场景需要结合文件名分区编号,避免跨文件行号混乱
- 解析每行的JSON数组并展开:把每行的JSON字符串解析为数组后炸开,每个文档对象对应一行记录,再摊平所有字段即可
代码实现(PySpark)
from pyspark.sql import SparkSession, functions as F, Window from pyspark.sql.types import * spark = SparkSession.builder.getOrCreate() file_path = "你的文件路径" # 1. 读取文本并打行号 raw_df = spark.read.text(file_path) # 单文件场景打行号(行号从1开始) # raw_df = raw_df.withColumn("source_row_num", F.monotonically_increasing_id() + 1) # 多文件场景用下面的代码打行号,保证每个文件内行号从1开始连续 raw_df = raw_df.withColumn("file_name", F.input_file_name()) \ .withColumn("tmp_sort_id", F.monotonically_increasing_id()) \ .withColumn("source_row_num", F.row_number().over( Window.partitionBy("file_name").orderBy("tmp_sort_id") )) \ .drop("tmp_sort_id") # 2. 解析JSON数组并展开 # 定义JSON数组的schema,这里用Map类型兼容不同结构的文档对象 array_schema = ArrayType(MapType(StringType(), StringType())) result_df = raw_df \ .withColumn("doc_array", F.from_json(F.col("value"), array_schema)) \ .withColumn("doc", F.explode(F.col("doc_array"))) \ .select( "source_row_num", F.col("doc")["firstAttribute"].cast(LongType()).alias("firstAttribute"), F.col("doc")["secondAttribute"].cast(LongType()).alias("secondAttribute"), F.col("doc")["thirdAttribute"].cast(LongType()).alias("thirdAttribute") )
示例输入的运行结果
针对给出的示例输入,最终输出如下,完全符合预期:
| source_row_num | firstAttribute | secondAttribute | thirdAttribute |
|---|---|---|---|
| 1 | 1 | 2 | null |
| 1 | 10 | 20 | null |
| 2 | 3 | 4 | null |
| 2 | null | 6 | 5 |
注意事项
- 不要尝试先通过
spark.read.json读取完成后再补行号:该读取流程会在分布式执行阶段直接完成数组拆分,原始行的边界信息已经完全丢失,事后计算的行号无法和源文件行对应 - 用text数据源读取时,Spark默认以换行符作为记录切分边界,不会把单行JSON拆到不同分区,行号计算准确,支持TB级大文件读取
- 如果JSON结构固定,可以把
from_json用到的schema替换为明确的StructType数组,性能比Map类型更好,同时可以自动处理字段类型转换
内容的提问来源于stack exchange,提问作者BelowZero
相关产品推荐
相关产品推荐

