【问题标题】:How to stop all comsumers with BlockingQueue如何使用 BlockingQueue 停止所有消费者
【发布时间】:2018-03-21 17:06:06
【问题描述】:

我正在尝试开发一个生产者-消费者系统,生产者将文件插入阻塞队列,消费者获取文件并处理它们。

我想创建一个选项来停止和恢复系统(停止所有消费者线程)。

夏天:

  1. 用户单击“停止”按钮调用 MyProgram.java 中的 stop() 函数。
  2. 每个消费者线程都会检查正在运行的进程是否处于活动状态并终止该进程。将消费者 while 更改为 false 后。
  3. 我得到了以下所有异常 (java.lang.InterruptedException)

** 我附上了我的课程的伪代码。 我究竟做错了什么 ?

错误 - 来自每个消费者线程

java.lang.InterruptedException
Mar 21, 2018 6:43:15 PM Consumer run
SEVERE: null
java.lang.InterruptedException
    at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.reportInterruptAfterWait(AbstractQueuedSynchronizer.java:2014)
    at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2048)
    at java.util.concurrent.ArrayBlockingQueue.take(ArrayBlockingQueue.java:403)
    at Consumer.run(Consumer.java:75)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)

Mar 21, 2018 6:43:15 PM Consumer run
SEVERE: null
java.lang.InterruptedException
    at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.reportInterruptAfterWait(AbstractQueuedSynchronizer.java:2014)
    at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2048)
    at java.util.concurrent.ArrayBlockingQueue.take(ArrayBlockingQueue.java:403)
    at Consumer.run(Consumer.java:75)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)

消费者

public class Consumer implements Runnable {
    protected BlockingQueue<File> queue = null;
    private volatile Boolean threadRunning = true;
    private File originalFile;

    public void run() {
        while (threadRunning) {
            try {
                originalFile = queue.take();
                tempFile = new File(tempPath.toString() + "\\" + originalFile.getName());;
                try {
                    process = getProcessBuilder(tempFile).start();
                    process.waitFor();

                } catch (IOException e) {
                    System.err.println(e);

                }
            } catch (InterruptedException e) {
                System.err.println(e);
            }
        }
    }

    public void stop() {
        threadRunning = false;
        if (process != null && process.isAlive()) {
            process.destroy();
        }
    }
}

制片人

public class Producer implements Runnable {

    protected BlockingQueue<File> queue = null;
    private volatile boolean running;

    public Producer(BlockingQueue queue) {
        this.queue = queue;
        running=true;
    }


    @Override
    public void run() {
        System.out.println("Producer Started");
        try {
            while (running) {
            // Adding files to Queue ........
            }
        } catch (IOException ex) {
                    System.err.println(ex);
        } catch (InterruptedException ex) {
                    System.err.println(ex);
        }
    }

}

MyProgram 类

public class MyProgram {

    private ExecutorService service;
    private Vector<Consumer> listOfThreads;

    private TaskParameters task;

    private Producer producer;

    public MyProgram(String taskDirPath, String taskName) {
        start();
    }

    public void start() {
        createAndExecuteThreads();
    }

    private void createAndExecuteThreads() {

        BlockingQueue<File> queue = new ArrayBlockingQueue(100);
        producer = new Producer(queue);
        new Thread(producer).start();
        service = Executors.newCachedThreadPool();
        for (int i = 0; i < task.getNumOfThreads(); ++i) {
            listOfThreads.add(new Consumer(queue, "Consumer " + (i + 1)));
        }
        for (int i = 0; i < listOfThreads.size(); ++i) {
            service.submit(listOfThreads.get(i));
        }
        service.shutdown();
    }

    public void stop() {
        Running = false;
        for (int i = 0; i < listOfThreads.size(); i++) {
            listOfThreads.get(i).stop();
        }
        service.shutdownNow();
    }



}

【问题讨论】:

    标签: java multithreading thread-safety producer-consumer blockingqueue


    【解决方案1】:

    我认为您的所有交易都在等待queue.take() 呼叫,如果您确定队列总是被填满,您的方法可能会起作用......

    否则,阻止一切的更好方法是

    • 提供者清空队列并开始使用 X 个假文件对其进行归档,例如名为 stop.stop 的文件(X 是消费者线程的数量)
    • 消费者检测到虚假文件并将自己的 threadRunning 设置为 false。

    【讨论】:

    • 感谢您的回答,我看到了您所说的“毒丸”解决方案,非常好的解决方案。但我认为它不能解决我的问题,因为我想立即“杀死”正在运行的进程。据我了解,消费者将完成该过程的工作,然后他才会停止提取文件。没有?
    猜你喜欢
    • 1970-01-01
    • 2019-09-25
    • 2017-09-23
    • 1970-01-01
    • 1970-01-01
    • 2017-11-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多