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

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);
    }
}

问题分析与解决方案

核心问题

  1. 逻辑方向错误:当前代码仅将最忙线程设为暂停,与「其余线程暂停」的需求完全相反。
  2. wait/notify使用不当:每个线程用自身对象锁等待,唤醒时仅调用notify()只能唤醒一个线程;且设置paused与检查等待的逻辑未同步,存在竞态条件。
  3. 状态分散管理:每个线程维护独立的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 22:34:55