You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 08:00:35