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

基于异步队列的Java多线程文件写入线程安全问题排查

问题分析与解决方案

首先,咱们来拆解下你遇到的核心问题:虽然writeToFile方法加了synchronized,但执行emptyQueues写入文件时,其他线程依然能往队列里添加数据。这本质是锁的范围和对象不匹配导致的并发问题,具体原因和优化方案如下:

问题根源

  1. 锁对象不统一:
    writeToFile的synchronized锁住的是WriterImpl实例本身,而run方法里执行emptyQueues时,用的是fastQQueue1和fastQQueue2作为锁对象。这两组锁完全独立,互相不干扰,所以emptyQueues执行期间,writeToFile依然能获取实例锁并修改队列。

  2. 队列大小判断非原子操作:
    run方法里的if(fastQQueue1.size() > MAX_QUEUE_SIZE)是单独的判断逻辑,从判断完成到进入同步块的间隙,队列大小可能已经被其他线程修改,导致判断结果失效。

  3. 未利用LinkedBlockingQueue原生特性:
    LinkedBlockingQueue本身就是线程安全的有界队列,自带阻塞等待的put方法,你完全不用手动实现队列满时的阻塞逻辑,重复造轮子反而容易出错。

优化方案

咱们可以利用LinkedBlockingQueue的原生特性简化代码,同时彻底解决并发问题,具体重构步骤如下:

1. 初始化有界队列

把队列初始化为指定容量的有界队列,这样队列满时put方法会自动阻塞,直到有空闲位置:

@Component("writer")
public class WriterImpl implements Writer {
    private volatile boolean isRunning; // 用volatile保证多线程可见性
    private PrintWriter fastQWriter1, fastQWriter2;
    private final int MAX_QUEUE_SIZE = 5000;
    // 初始化指定容量的有界LinkedBlockingQueue
    private final Queue<FastQRecord> fastQQueue1 = new LinkedBlockingQueue<>(MAX_QUEUE_SIZE);
    private final Queue<FastQRecord> fastQQueue2 = new LinkedBlockingQueue<>(MAX_QUEUE_SIZE);
    
    // ... 其他方法
}

2. 简化writeToFile方法

去掉synchronized,改用LinkedBlockingQueue的put方法,它会自动处理队列满时的阻塞,且本身是线程安全的:

@Override
public void writeToFile(FastQRecord one, FastQRecord two) throws InterruptedException {
    // 队列满时自动阻塞,直到有空闲位置
    fastQQueue1.put(one);
    fastQQueue2.put(two);
}

注:put方法会抛出InterruptedException,你可以选择在方法上抛出,或者内部捕获并处理(比如中断当前线程)。

3. 优化emptyQueues方法

用drainTo批量取出队列元素,比循环poll更高效,且是原子操作:

private void emptyQueues() {
    List<FastQRecord> batch1 = new ArrayList<>();
    fastQQueue1.drainTo(batch1); // 一次性取出队列所有元素
    for (FastQRecord record : batch1) {
        writeRecord(fastQWriter1, record);
    }

    List<FastQRecord> batch2 = new ArrayList<>();
    fastQQueue2.drainTo(batch2);
    for (FastQRecord record : batch2) {
        writeRecord(fastQWriter2, record);
    }

    // 批量写入后刷新,减少IO次数
    fastQWriter1.flush();
    fastQWriter2.flush();
}

// 提取重复写入逻辑,让代码更简洁
private void writeRecord(PrintWriter writer, FastQRecord record) {
    writer.println(record.getId());
    writer.println(record.getRawSequence());
    writer.println(record.getPlus());
    writer.println(record.getQualityString());
}

4. 重构run方法

确保线程退出前写完队列剩余数据,同时避免无数据时频繁空转消耗CPU:

@Override
public void run() {
    isRunning = true;
    // 循环条件:线程运行中 或 队列还有未写入的数据
    while (isRunning || !fastQQueue1.isEmpty() || !fastQQueue2.isEmpty()) {
        emptyQueues();
        try {
            // 短暂休眠,避免空转
            Thread.sleep(10);
        } catch (InterruptedException e) {
            // 捕获中断信号,终止线程
            Thread.currentThread().interrupt();
            break;
        }
    }
}

5. 完善close方法

确保关闭前写完所有剩余数据,安全释放资源:

@Override
public void close() {
    isRunning = false;
    // 等待队列数据写完(可添加超时逻辑,避免无限等待)
    while (!fastQQueue1.isEmpty() || !fastQQueue2.isEmpty()) {
        try {
            Thread.sleep(50);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            break;
        }
    }
    // 安全关闭流,避免空指针或异常
    safeClose(fastQWriter1);
    safeClose(fastQWriter2);
}

private void safeClose(PrintWriter writer) {
    if (writer != null) {
        writer.flush();
        writer.close();
    }
}

额外性能优化建议

  • 对于20-50GB的大文件写入,建议用BufferedWriter包装输出流,大幅提升写入性能:
    fastQWriter1 = new PrintWriter(new BufferedWriter(new FileWriter(fastQ1)));
    
  • 可以考虑用ExecutorService管理线程,比手动启动线程更灵活可控;
  • 单个文件不建议多线程写入(会导致数据混乱),如果要提升吞吐量,可以拆分数据到多个文件,用多个写入线程分别处理。

内容的提问来源于stack exchange,提问作者Sam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:54:10