基于异步队列的Java多线程文件写入线程安全问题排查
首先,咱们来拆解下你遇到的核心问题:虽然writeToFile方法加了synchronized,但执行emptyQueues写入文件时,其他线程依然能往队列里添加数据。这本质是锁的范围和对象不匹配导致的并发问题,具体原因和优化方案如下:
问题根源
锁对象不统一:
writeToFile的synchronized锁住的是WriterImpl实例本身,而run方法里执行emptyQueues时,用的是fastQQueue1和fastQQueue2作为锁对象。这两组锁完全独立,互相不干扰,所以emptyQueues执行期间,writeToFile依然能获取实例锁并修改队列。队列大小判断非原子操作:
run方法里的if(fastQQueue1.size() > MAX_QUEUE_SIZE)是单独的判断逻辑,从判断完成到进入同步块的间隙,队列大小可能已经被其他线程修改,导致判断结果失效。未利用
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

