Java中文本文件随机访问:1.4TB WikiData JSON文件跳转指定行
处理超大WikiData JSON Dump的随机位置读取问题
直接跳转到第20000000行的方案不可行,原因有两个:一是1.4TB级别的文件无法直接计算指定行对应的字节偏移,必须预先扫描记录(这和你之前统计行数的耗时一样);二是WikiData的JSON结构不一定和行号严格对齐,跳转后大概率落在某个JSON对象的中间,导致解析失败。下面针对你提到的两种提速场景,给出可行的替代方案:
一、本机多SSD场景:预生成偏移量索引
先花一次时间扫描大文件,记录每个完整JSON对象的起始字节偏移,把这些偏移量存在一个小索引文件里,后续就能快速定位到目标对象:
- 生成索引(一次性操作)
用Jackson的流式解析器扫描文件,跳过对象内容只记录偏移:import com.fasterxml.jackson.core.JsonFactory; import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.JsonToken; import java.io.*; import java.util.ArrayList; import java.util.List; public class IndexGenerator { public static void main(String[] args) throws IOException { String inputFile = "path/to/your/wikidata.dump"; String indexFile = "path/to/index.dat"; RandomAccessFile raf = new RandomAccessFile(inputFile, "r"); JsonFactory factory = new JsonFactory(); JsonParser jp = factory.createParser(raf); List<Long> offsets = new ArrayList<>(); while (jp.nextToken() != JsonToken.END_OBJECT) { // 记录当前JSON对象的起始字节偏移 offsets.add(jp.getCurrentLocation().getByteOffset()); // 跳过当前对象的所有子节点,直接跳到下一个对象 jp.skipChildren(); } // 将偏移量写入索引文件 try (ObjectOutputStream oos = new ObjectOutputStream(new FileOutputStream(indexFile))) { oos.writeObject(offsets); } jp.close(); raf.close(); } } - 快速定位目标对象
读取索引文件找到目标偏移,用RandomAccessFile定位后再解析:
这种方式下,索引文件体积很小(每记录一个偏移只占8字节,2000万条也才160MB左右),后续定位几乎是瞬时的。import com.fasterxml.jackson.core.JsonFactory; import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.JsonToken; import java.io.*; import java.util.List; public class RandomReader { public static void main(String[] args) throws IOException, ClassNotFoundException { String inputFile = "path/to/your/wikidata.dump"; String indexFile = "path/to/index.dat"; int targetObjectIndex = 20000000 - 1; // 索引从0开始 // 读取索引文件 List<Long> offsets; try (ObjectInputStream ois = new ObjectInputStream(new FileInputStream(indexFile))) { offsets = (List<Long>) ois.readObject(); } // 定位到目标偏移 RandomAccessFile raf = new RandomAccessFile(inputFile, "r"); raf.seek(offsets.get(targetObjectIndex)); // 开始解析目标对象 JsonFactory factory = new JsonFactory(); JsonParser jp = factory.createParser(raf); jp.nextToken(); // 指向JSON对象起始 // 这里写你的对象处理逻辑 // ... jp.close(); raf.close(); } }
二、Apache Spark集群场景:分片流式解析
Spark的分布式特性天然适合处理超大文件,不需要手动拆分,只需利用其分片机制,在每个分片里用流式解析处理完整的JSON对象:
import com.fasterxml.jackson.core.{JsonFactory, JsonToken} import org.apache.spark.sql.SparkSession object WikiDataSparkParse { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("WikiDataParse") .getOrCreate() val inputFile = "path/to/your/wikidata.dump" val fileSize = new java.io.File(inputFile).length() val splitSize = 64 * 1024 * 1024 // 64MB分片,可根据集群配置调整 // 生成所有分片的字节范围 val splits = (0 to (fileSize / splitSize).toInt).map(i => { val start = i * splitSize val end = if (i == (fileSize / splitSize).toInt) fileSize else (i + 1) * splitSize (start, end) }) // 并行处理每个分片 spark.sparkContext.parallelize(splits) .mapPartitions { iter => val factory = new JsonFactory() iter.flatMap { case (start, end) => val raf = new java.io.RandomAccessFile(inputFile, "r") raf.seek(start) val jp = factory.createParser(raf) val results = collection.mutable.ListBuffer[String]() // 跳过分片起始位置的不完整对象,找到第一个合法的START_OBJECT var token = jp.nextToken() while (token != null && token != JsonToken.START_OBJECT) { token = jp.nextToken() } // 处理分片内的所有完整对象,直到接近分片末尾 while (token != null && jp.getCurrentLocation.getByteOffset < end) { // 读取完整对象并处理 val jsonNode = jp.readValueAsTree() results += jsonNode.toString() // 跳到下一个对象 token = jp.nextToken() } jp.close() raf.close() results } } // 后续输出或处理逻辑 .saveAsTextFile("path/to/output") spark.stop() } }
这种方案的核心是:每个分片只处理范围内的完整JSON对象,跳过开头的不完整部分,避免解析错误,同时Spark会自动把分片分配到集群节点上并行处理,最大化利用多机器资源。
内容的提问来源于stack exchange,提问作者Sogawa-sps
相关产品推荐
相关产品推荐

