【问题标题】:How to thread-safe signal threads to pause in Java如何在 Java 中线程安全地向线程发出信号以暂停
【发布时间】:2013-01-13 13:44:27
【问题描述】:

我有一堆线程同时运行。有时一个线程需要通知其他线程等待它完成一个工作并再次发出信号让它们恢复。由于我对 Java 的同步有点陌生,我想知道做这种事情的正确方法是什么。我的代码是这样的:

private void Concurrent() {
    if (shouldRun()) {
        // notify threads to pause and wait for them
        DoJob();
        // resume threads
    }

    // Normal job...
}

更新:

请注意,我编写的代码在一个类中,每个线程都将执行该类。我无权访问这些线程或它们的运行方式。我只是在线程中。

更新 2:

我的代码来自爬虫类。爬虫类(crawler4j)知道如何处理并发。我唯一需要的是在运行一个函数之前暂停其他爬虫,然后再恢复它们。这段代码是我爬虫的基础:

   public class TestCrawler extends WebCrawler {
    private SingleThread()
    {
        //When this function is running, no other crawler should do anything
    }

    @Override
    public void visit(Page page) {
        if(SomeCriteria())
        {
            //make all other crawlers stop until I finish
            SingleThread();
            //let them resume
        }

        //Normal Stuff
    }
   }

【问题讨论】:

  • 您是否考虑过改用BlockingQueue?
  • 线程只能通过添加代码来停止。您无法安全地从外部停止/暂停线程。

标签: java multithreading synchronization


【解决方案1】:

下面是一个简短的例子,说明如何使用很酷的 Java 并发实现这一目标:

snip Pause 类不再影响旧代码。

编辑:

这是新的 Test 类:

package de.hotware.test;

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class Test {

    private Pause mPause;

    public Test() {
        this.mPause = new Pause();
    }

    public void concurrent() throws InterruptedException {
        while(true) {
            this.mPause.probe();
            System.out.println("concurrent");
            Thread.sleep(100);
        }
    }

    public void crucial() throws InterruptedException {
        int i = 0;
        while (true) {
            if (i++ % 2 == 0) {
                this.mPause.pause(true);
                System.out.println("crucial: exclusive execution");
                this.mPause.pause(false);
            } else {
                System.out.println("crucial: normal execution");
                Thread.sleep(1000);
            }
        }
    }

    public static void main(String[] args) {
        final Test test = new Test();
        Runnable run = new Runnable() {

            @Override
            public void run() {
                try {
                    test.concurrent();
                } catch (InterruptedException e) {
                    // TODO Auto-generated catch block
                    e.printStackTrace();
                }
            }

        };
        Runnable cruc = new Runnable() {

            @Override
            public void run() {
                try {
                    test.crucial();
                } catch (InterruptedException e) {
                    // TODO Auto-generated catch block
                    e.printStackTrace();
                }
            }

        };
        ExecutorService serv = Executors.newCachedThreadPool();
        serv.execute(run);
        serv.execute(run);
        serv.execute(cruc);
    }

}

还有实用程序 Pause 类:

package de.hotware.test;

import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;

/**
 * Utility class to pause and unpause threads
 * with Java Concurrency
 * @author Martin Braun
 */
public class Pause {

    private Lock mLock;
    private Condition mCondition;
    private AtomicBoolean mAwait;

    public Pause() {
        this.mLock = new ReentrantLock();
        this.mCondition = this.mLock.newCondition();
        this.mAwait = new AtomicBoolean(false);
    }

    /**
     * waits until the threads until this.mAwait is set to true
     * @throws InterruptedException
     */
    public void probe() throws InterruptedException {
        while(this.mAwait.get()) {
            this.mLock.lock();
            try {
                this.mCondition.await();
            } finally {
                this.mLock.unlock();
            }
        }
    }

    /**
     * pauses or unpauses
     */
    public void pause(boolean pValue) {
        if(!pValue){
            this.mLock.lock();
            try {
                this.mCondition.signalAll();
            } finally {
                this.mLock.unlock();
            }
        }
        this.mAwait.set(pValue);
    }

}

基本用法是在每次运行前调用probe()。如果在调用 pause(false) 之前暂停,这将阻塞。

你的班级应该是这样的:

public class TestCrawler extends WebCrawler {

private Pause mPause;

public TestCrawler(Pause pPause) {
    this.mPause = pPause;
}

private SingleThread()
{
        //When this function is running, no other crawler should do anything
}

@Override
public void visit(Page page) {
    if(SomeCriteria())
    {
        //only enter the crucial part once if it has to be exclusive
        this.mPause.probe();
        //make all other crawlers stop until I finish
        this.mPause.pause(true);
        SingleThread();
        //let them resume
        this.mPause.pause(false);
    }
    this.mPause.probe();
    //Normal Stuff
}
}

【讨论】:

  • 好的,这里有很多东西。你确定我的案子需要这一切吗?能否请您阅读我的最新更新并告诉我如何修改我的代码。
  • concurrent(...) 在您的情况下与访问相同。 critical(...) 与 SingleThread(...) 相同。我不知道这个框架,但如果你有多个类的实例,你可能需要将一个对象传递给他们可以锁定的所有实例。不过,让我一起摆弄一些东西吧。
  • 我使用了你的代码。几次测试并运行良好。但是,当我在实际情况下进行测试时,似乎有一些错误。就像:Crawler1 -> Probe Crawler2 -> Probe Crawler3 -> Probe ... 你能帮忙吗?谢谢。
  • 您是否使用了我修改过的类中的代码?而且您可能必须使您的 SomeCriteria() 方法线程安全。
  • 我认为您的一般问题是多个线程成为独占线程。您必须使这部分线程安全,并且不会出现死锁。基本上,您需要做的是创建一个“原子”块,您可以在其中获得许可,并且一旦获得许可,就让其他人得不到许可。这可以通过 AtomicBoolean 和 getAndSet(value) 来完成。这将检索您是否被允许获得成为独占线程的权限,然后将其设置为 false(因此已被采用)。但是有了这个,你必须在将它设置为 true 时关心竞争条件。
【解决方案2】:
public class StockMonitor extends Thread {

    private boolean suspend = false; 
    private volatile Thread thread;

    public StockMonitor() {
        thread = this;
    }

    // Use name with underscore, in order to avoid naming crashing with
    // Thread's.
    private synchronized void _wait() throws InterruptedException {
        while (suspend) {
            wait();
        }
    }

    // Use name with underscore, in order to avoid naming crashing with
    // Thread's.
    public synchronized void _resume() {
        suspend = false;
        notify();
    }

    // Use name with underscore, in order to avoid naming crashing with
    // Thread's.
    public synchronized void _suspend() {
        suspend = true;
    }  

     public void _stop() { 
        thread = null;
        // Wake up from sleep.
        interrupt();     
     }

     @Override
     public void run() {
        final Thread thisThread = Thread.currentThread();
        while (thisThread == thread) {
            _wait();
            // Do whatever you want right here.
        }
     }
}

调用_resume 和_suspend 将使您能够恢复和暂停Thread。 _stop 会让你优雅地停止线程。请注意,一旦停止Thread,就无法再次恢复它。 Thread 不再可用。

代码是从一个真实世界的开源项目中挑选出来的:http://jstock.hg.sourceforge.net/hgweb/jstock/jstock/file/b17c0fbfe37c/src/org/yccheok/jstock/engine/RealTimeStockMonitor.java#l247

【讨论】:

  • 谢谢。我注意到您正在自己创建线程。我不能那样做。我已经更新了我的问题。那我该怎么做类似的事情呢?
  • @AlirezaNoori 我们可以这样做吗- 创建一个daemon thread 并要求其他线程Join 它?
  • 我认为我没有很好地解释我的问题。请阅读第二次更新。我认为它会澄清问题。
【解决方案3】:

你可以使用wait()和notify()

线程等待:

  // define mutex as field
  Object mutex = new Object();

   // later:
   synchronized(mutex) {
        wait();
   }

通知线程继续

   synchronized (mutex) {
        notify();
   }

【讨论】:

  • 你能给我一个简单的例子吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-07-27
  • 1970-01-01
  • 1970-01-01
  • 2016-03-21
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多