如何在Spark中获取多文件中每行对应的源文件行号?
问题:为S3中的JSONL文件每行添加源文件名和行号
我有一个存储多个JSONL文件的S3存储桶,每个文件的每一行都是一个JSON字符串。目前我用以下代码读取文件:
Dataset<Row> dataset = spark.read().option("recursiveFileLookup", "true").json(path);
这样得到的Dataset包含JSON对象的所有字段,现在需要为每一行数据添加它在源文件中的对应行号,示例如下:
示例文件
file1.jsonl:
{"key1": "a"} {"key1": "b"} {"key1": "c"}
file2.jsonl:
{"key1": "x"} {"key1": "y"} {"key1": "z"}
期望结果
| key1 | file_name | line_num |
|---|---|---|
| a | file1 | 1 |
| b | file1 | 2 |
| c | file1 | 3 |
| x | file2 | 1 |
| y | file2 | 2 |
| z | file2 | 3 |
请问是否可以实现该需求?
解决方案
完全可以实现,核心做法是先按文本格式读取文件(保留每行内容和源文件信息),给每个文件内的行分配行号后,再解析JSON字段。
具体Java代码实现如下:
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.Window; import org.apache.spark.sql.types.StructType; import static org.apache.spark.sql.functions.*; // 1. 按文本读取所有文件,同时获取源文件名 Dataset<Row> textDataset = spark.read() .option("recursiveFileLookup", "true") .text(path) // 从完整路径里提取不带后缀的文件名,比如把"s3://bucket/path/file1.jsonl"转成"file1" .withColumn("file_name", regexp_replace(input_file_name(), ".*/(.*)\\.jsonl", "$1")); // 2. 按文件名分组,给每组内的行按顺序加行号 Dataset<Row> numberedDataset = textDataset .withColumn("line_num", row_number().over(Window.partitionBy("file_name").orderBy(lit(1)))); // 3. 解析JSON字符串,展开原有字段,再合并文件名和行号 // 先获取JSON的schema,你可以手动指定或者用推断的方式 StructType schema = spark.read().json(path).schema(); Dataset<Row> finalDataset = numberedDataset .withColumn("json_data", from_json(col("value"), schema)) .select("json_data.*", "file_name", "line_num"); finalDataset.show();
注意点
- 不能直接用
spark.read.json读取后加行号,因为JSON读取会打乱行的原始顺序关联,没法准确对应文件内的行号。 input_file_name()能拿到当前行的完整文件路径,用正则替换可以提取你需要的文件名格式,要是你的文件名规则不一样,调整正则就行。- 窗口函数里的
orderBy(lit(1))是因为文本读取的顺序就是文件内的行顺序,不需要额外排序;如果要确保绝对一致,也可以依赖Spark文本读取的默认顺序(通常和文件内行序一致)。 - 要是JSON字段固定,手动指定schema会比推断更高效,比如:
import org.apache.spark.sql.types.StringType; StructType schema = new StructType().add("key1", StringType);
内容的提问来源于stack exchange,提问作者shashank93rao
相关产品推荐
相关产品推荐

