使用Vert.x逐行读取大文件(AsyncFile与RecordParser及背压问题)
嘿,这个问题我之前也踩过坑!确实,RecordParser本身并不是WriteStream,所以没法直接和Pump搭配使用——毕竟Pump的核心作用就是在ReadStream和WriteStream之间流转数据,自动处理背压问题。不过咱们有几种靠谱的方案,既能实现大文件逐行处理,又能妥善处理背压:
方案一:手动控制流的暂停/恢复(简单直接)
既然AsyncFile是ReadStream,我们可以在逐行处理的时候手动控制它的读取节奏,避免数据堆积。这种方式不需要额外的适配器,逻辑清晰:
// 推荐用异步open,避免阻塞事件循环 vertx.fileSystem().open("path/to/your/large/file", new OpenOptions(), ar -> { if (ar.failed()) { ar.cause().printStackTrace(); return; } AsyncFile asyncFile = ar.result(); RecordParser recordParser = RecordParser.newDelimited("\n", bufferedLine -> { // 这里是你的逐行处理逻辑,假设是异步操作(比如写DB、调用外部API) processLineAsync(bufferedLine.toString()) .onComplete(res -> { // 处理完成后,恢复文件读取 asyncFile.resume(); if (res.failed()) { res.cause().printStackTrace(); } }); // 开始处理当前行时,暂停读取,避免数据积压 asyncFile.pause(); }); // 把文件流传给RecordParser asyncFile.handler(recordParser); // 监听文件读取结束 asyncFile.endHandler(v -> System.out.println("所有行处理完成")); // 处理读取异常 asyncFile.exceptionHandler(err -> err.printStackTrace()); }); // 模拟异步处理方法 private Future<Void> processLineAsync(String line) { Promise<Void> promise = Promise.promise(); // 这里替换成你的实际业务逻辑,比如延迟模拟处理耗时 vertx.setTimer(100, id -> promise.complete()); return promise.future(); }
这个方案的核心逻辑:
- 每次拿到一行数据就暂停文件读取,防止后续数据不断涌入内存
- 等当前行的异步处理完成后,再恢复读取下一行
- 完全手动实现背压控制,适合处理逻辑明确、异步操作清晰的场景
方案二:将RecordParser适配为WriteStream(复用Pump能力)
如果想保留Pump自动处理背压的能力,我们可以写一个简单的适配器类,把RecordParser包装成WriteStream,这样就能和Pump无缝配合了:
// 自定义WriteStream适配器,将数据转发给RecordParser class RecordParserWriteStream implements WriteStream<Buffer> { private final RecordParser recordParser; private final Promise<Void> endPromise = Promise.promise(); private Handler<Void> drainHandler; private int pendingTasks = 0; private final int MAX_PENDING = 100; // 自定义最大待处理任务数,控制背压 public RecordParserWriteStream(RecordParser recordParser) { this.recordParser = recordParser; } @Override public WriteStream<Buffer> write(Buffer data) { pendingTasks++; recordParser.handle(data); // 如果待处理任务数没到上限,直接触发drain if (pendingTasks < MAX_PENDING && drainHandler != null) { vertx.runOnContext(v -> drainHandler.handle(null)); } return this; } @Override public void end() { endPromise.complete(); } @Override public void end(Buffer data) { write(data); end(); } @Override public WriteStream<Buffer> setWriteQueueMaxSize(int maxSize) { return this; } @Override public boolean writeQueueFull() { // 当待处理任务数达到上限时,告诉Pump暂停读取 return pendingTasks >= MAX_PENDING; } @Override public WriteStream<Buffer> drainHandler(Handler<Void> handler) { this.drainHandler = handler; return this; } @Override public WriteStream<Buffer> exceptionHandler(Handler<Throwable> handler) { return this; } // 处理完一行后调用,减少待处理任务数并触发drain public void onLineProcessed() { pendingTasks--; if (pendingTasks < MAX_PENDING && drainHandler != null) { vertx.runOnContext(v -> drainHandler.handle(null)); } } }
然后使用这个适配器搭配Pump:
vertx.fileSystem().open("path/to/your/large/file", new OpenOptions(), ar -> { if (ar.failed()) { ar.cause().printStackTrace(); return; } AsyncFile asyncFile = ar.result(); RecordParserWriteStream writeStream = new RecordParserWriteStream( RecordParser.newDelimited("\n", bufferedLine -> { processLineAsync(bufferedLine.toString()) .onComplete(res -> { if (res.failed()) { res.cause().printStackTrace(); } // 通知适配器任务完成 writeStream.onLineProcessed(); }); }) ); // 创建Pump并启动,自动处理背压 Pump pump = Pump.pump(asyncFile, writeStream); pump.start(); asyncFile.endHandler(v -> { writeStream.end(); System.out.println("文件处理完成"); }); asyncFile.exceptionHandler(err -> err.printStackTrace()); });
这个方案的优势:
- 复用了
Pump的背压处理逻辑,不需要手动管理pause/resume - 通过适配器里的
pendingTasks和MAX_PENDING,可以灵活控制并发处理的任务数,平衡性能和内存占用
额外提醒:
尽量避免使用openBlocking,因为它会阻塞Vert.x的事件循环,影响整个应用的性能。优先用异步的open方法,符合Vert.x的非阻塞设计理念。
内容的提问来源于stack exchange,提问作者lebolo
相关产品推荐
相关产品推荐

