Java中如何在一个线程内暂停其他线程?
问题描述
我在多线程程序中需要实现以下逻辑:当某个线程的特定方法被调用时,其余所有线程暂停,待该方法执行完毕后再恢复其余线程。
我尝试用volatile boolean字段实现:需要暂停时将字段设为true,其他线程在while循环中检查该字段,若为true则调用this.wait();完成操作后将布尔值设为false并调用notify(),但该方案未生效。
Worker线程代码
import java.io.IOException; import java.io.RandomAccessFile; public class Worker extends Thread { RandomAccessFile file; long start; long interval; long end; long cursor; SourceProvider sp; volatile MultiThreadCopier copier; boolean isDone = false; volatile boolean paused = false; public Worker(String dest, long start, long interval, SourceProvider sp, MultiThreadCopier copier) { try { file = new RandomAccessFile(dest, "rws"); file.setLength(sp.size()); } catch (IOException e) { e.printStackTrace(); } this.start = start; this.interval = interval; this.end = start + interval - 1; this.cursor = start; this.sp = sp; this.copier = copier; } @Override public void run() { write(); while (isDone) { Pair<Worker, Long> busiestWorkerPair = copier.busiestWorker(); Worker busiestWorker = busiestWorkerPair.getKey(); if (busiestWorker == null) { break; } long workRemaining = busiestWorkerPair.getValue(); this.end = busiestWorker.end; this.start = busiestWorker.end - workRemaining / 2 + 1; busiestWorker.end = this.start - 1; isDone = false; busiestWorker.paused = false; synchronized (busiestWorker) { busiestWorker.notify(); } write(); } try { file.close(); } catch (IOException e) { e.printStackTrace(); } } public void write() { SourceReader sr = sp.connect(start); cursor = start; try { file.seek(cursor); } catch (IOException e) { e.printStackTrace(); } while (cursor <= end) { while (paused) { System.out.println(getName() + " paused"); try { synchronized (this) { wait(); } } catch (InterruptedException e) { e.printStackTrace(); } } try { file.write(sr.read()); } catch (IOException e) { e.printStackTrace(); } cursor++; } isDone = true; } }
MultiThreadCopier线程代码
public class MultiThreadCopier { public static final long SAFE_MARGIN = 6; Worker[] workers; public MultiThreadCopier(SourceProvider sourceProvider, String dest, int workerCount) { workers = new Worker[workerCount]; long interval = sourceProvider.size() / workerCount; for (int i = 0; i < workerCount; i++) { long start = i * interval; workers[i] = new Worker(dest, start, interval + ((i + 1 == workerCount) ? sourceProvider.size() % workerCount : 0), sourceProvider, this); } } public void start() { for (Worker worker : workers) { worker.start(); } } public Pair<Worker, Long> busiestWorker() { long mostWorkRemaining = 0; Worker busiestWorker = null; for (Worker worker : workers) { if (worker.isDone) continue; long remainingWork = worker.end - worker.cursor + 1; if (remainingWork < SAFE_MARGIN) continue; if (remainingWork <= mostWorkRemaining) continue; mostWorkRemaining = remainingWork; busiestWorker = worker; } if (busiestWorker == null) return new Pair<>(null, 0L); busiestWorker.paused = true; return new Pair<>(busiestWorker, mostWorkRemaining); } }
问题分析与解决方案
核心问题
- 逻辑方向错误:当前代码仅将最忙线程设为暂停,与「其余线程暂停」的需求完全相反。
- wait/notify使用不当:每个线程用自身对象锁等待,唤醒时仅调用
notify()只能唤醒一个线程;且设置paused与检查等待的逻辑未同步,存在竞态条件。 - 状态分散管理:每个线程维护独立的
paused字段,难以实现全局统一的暂停/恢复控制。
修正方案
1. 全局统一暂停控制(修改MultiThreadCopier)
添加全局暂停开关和锁对象,统一控制所有线程的暂停与恢复:
public class MultiThreadCopier { public static final long SAFE_MARGIN = 6; Worker[] workers; private volatile boolean globalPause = false; private final Object pauseLock = new Object(); public MultiThreadCopier(SourceProvider sourceProvider, String dest, int workerCount) { // 原有构造逻辑不变 workers = new Worker[workerCount]; long interval = sourceProvider.size() / workerCount; for (int i = 0; i < workerCount; i++) { long start = i * interval; workers[i] = new Worker(dest, start, interval + ((i + 1 == workerCount) ? sourceProvider.size() % workerCount : 0), sourceProvider, this); } } public void start() { for (Worker worker : workers) { worker.start(); } } // 暂停所有线程 public void pauseAllWorkers() { globalPause = true; } // 恢复所有线程 public void resumeAllWorkers() { synchronized (pauseLock) { globalPause = false; pauseLock.notifyAll(); } } // 供Worker线程检查是否需要暂停 public void checkPause() throws InterruptedException { synchronized (pauseLock) { while (globalPause) { pauseLock.wait(); } } } public Pair<Worker, Long> busiestWorker() { // 先暂停所有线程 pauseAllWorkers(); // 原有寻找最忙线程逻辑 long mostWorkRemaining = 0; Worker busiestWorker = null; for (Worker worker : workers) { if (worker.isDone) continue; long remainingWork = worker.end - worker.cursor + 1; if (remainingWork < SAFE_MARGIN) continue; if (remainingWork <= mostWorkRemaining) continue; mostWorkRemaining = remainingWork; busiestWorker = worker; } if (busiestWorker == null) { resumeAllWorkers(); return new Pair<>(null, 0L); } return new Pair<>(busiestWorker, mostWorkRemaining); } }
2. 修改Worker类适配全局控制
移除独立的paused字段,改用全局暂停检查:
import java.io.IOException; import java.io.RandomAccessFile; public class Worker extends Thread { RandomAccessFile file; long start; long interval; long end; long cursor; SourceProvider sp; final MultiThreadCopier copier; boolean isDone = false; public Worker(String dest, long start, long interval, SourceProvider sp, MultiThreadCopier copier) { try { file = new RandomAccessFile(dest, "rws"); file.setLength(sp.size()); } catch (IOException e) { e.printStackTrace(); } this.start = start; this.interval = interval; this.end = start + interval - 1; this.cursor = start; this.sp = sp; this.copier = copier; } @Override public void run() { write(); while (isDone) { Pair<Worker, Long> busiestWorkerPair = copier.busiestWorker(); Worker busiestWorker = busiestWorkerPair.getKey(); if (busiestWorker == null) { break; } long workRemaining = busiestWorkerPair.getValue(); this.end = busiestWorker.end; this.start = busiestWorker.end - workRemaining / 2 + 1; busiestWorker.end = this.start - 1; isDone = false; // 任务调整完成后恢复所有线程 copier.resumeAllWorkers(); write(); } try { file.close(); } catch (IOException e) { e.printStackTrace(); } } public void write() { SourceReader sr = sp.connect(start); cursor = start; try { file.seek(cursor); } catch (IOException e) { e.printStackTrace(); } while (cursor <= end) { try { // 每次写入前检查是否需要暂停 copier.checkPause(); file.write(sr.read()); } catch (IOException | InterruptedException e) { e.printStackTrace(); } cursor++; } isDone = true; } }
关键修正说明
- 全局统一控制:用
globalPause变量统一管理所有线程的暂停状态,避免分散维护的混乱。 - 正确使用wait/notify:使用全局锁对象
pauseLock,恢复时调用notifyAll()确保所有等待线程被唤醒;同步块包裹检查逻辑,避免竞态条件。 - 流程顺序优化:先暂停所有线程再调整任务,完成后恢复,确保任务调整过程中无线程执行写入操作,避免数据不一致。
内容的提问来源于stack exchange,提问作者Harry Green
相关产品推荐
相关产品推荐

