【问题标题】:Best way to implement graceful cancel for running async jobs in java在java中运行异步作业实现优雅取消的最佳方法
【发布时间】:2016-10-03 04:44:44
【问题描述】:

假设我有一个这样的接口,让我的应用程序中的组件运行作业 -

IJob {
    IResult execute();
    void cancel();
}

我想设置我的应用程序,以便异步运行这些作业。期望调用取消应该立即执行返回,结果表明它已被取消。

最好的设置方法是什么?我可以创建一个 Thread 对象来运行它,它有额外的取消方法,但我也在查看我不熟悉的 Future 接口。

FutureTask 的问题是取消不优雅,不允许我调用 job.cancel()。扩展 FutureTask 并实现我自己的处理是否是个好主意?

【问题讨论】:

  • Future 怎么“不优雅”?多线程很难,试图改变 Future 的实现很可能会引入 bug。
  • 好吧,“不优雅”可能不是最好的选择。我的意思是 cancel() 发送一个线程中断,它并没有真正给我时间来正确调用 job.cancel()。如果我创建一个调用 job.execute() (阻塞调用)的 FutureTask,那么我将如何处理取消请求以便我可以调用 job.cancel() ?
  • 什么?对不起,我完全不明白。你必须向我们展示一些代码来解释为什么你没有“时间”去做某事。 (响应您的编辑:让execute() 抛出InterruptedException。但Java ExecutorService 是一个更好的主意。)
  • 另外启动一个新线程只是为了阻止它的完成是一个 TERRIBLE 的想法。您还不如让原始线程自己完成任务,如果您所做的只是阻塞和等待,就没有任何收获。您已经停止了一个线程并在其位置启动了另一个线程。有什么意义?
  • 我将尝试添加更多上下文...我正在处理的环境是插件架构,其中插件提供 IJob 实现。我想避免插件必须了解有关线程或处理中断异常等的任何信息。它们只是实现执行和取消。

标签: java java.util.concurrent


【解决方案1】:

当您在任务中调用cancel 时,它会向运行任务的线程发送中断信号。您的任务需要定期检查该信号是否已发送,并在发送时做出相应的反应:

if (Thread.interrupted()) {
    performNecessaryCleanup();
    return;
}

【讨论】:

    【解决方案2】:

    使用并发时,请使用语言提供的功能,而不是手动实现。

    据我了解,ExecutorService 应该是适合您的工具,因为您可以:

    • 为其提供将异步运行并可以返回结果的作业
    • 关闭执行器,以便取消所有正在运行的作业

    例子

    public static void main(String[] args) {
        ExecutorService executor = Executors.newSingleThreadExecutor();
        List<Future<IResult>> results = new ArrayList<>();
    
        for (int i = 0; i < 10; i++) {
            results.add(executor.submit(new Job(i))); //start jobs
        }
    
        executor.shutdownNow(); //attempts to stop all running jobs
    
        //program flow immediatly continues
    }
    

    就像@JoeC 在他的回答中解释的那样,保证所有作业停止的条件是 中断 在每个作业内部进行管理,因为每个线程都将被标记为 在调用shutdownNow()时被打断。

    if (Thread.interrupted()) {
        //return result cancelled
    }
    

    【讨论】:

    • 问题是 IJob 实现是由插件提供的,理想情况下不应该关心线程中断,或者与线程有关的任何事情。
    • 我现在的想法是扩展 FutureTask 并做一个非常简单的 cancel() 覆盖,它只调用 super.cancel() 和 job.cancel()。
    • @TR1096 好的,但最后,如果 IJob 不关心线程中断,则不保证在调用 cancel() 或 @987654326 时会取消 IJob @ 或任何类似的方法。
    【解决方案3】:

    调用IJob.execute()FutureTask.run()会阻塞当前线程,需要调度IJob未来任务

    调度 FutureTask 是最好的选择,消费者可以调用 FutureTask.get() 并等待结果(即使你调用 IJob.cancel( ))。

    我做了一个模拟 IJobIResult 的演示,它使用普通线程进行调度,在生产中你应该有一个像以前一样的 ExecutorService发布示例。

    如你所见,主线程可以检查调用FutureTask.isDone()的状态,基本上你是在检查结果是否已经设置。设置的结果意味着IJob的线程已经结束。

    您几乎可以随时调用 IJob.cancel() 来完成 FutureTask 中的包装 IJob厘米。

    模拟工作:

    public class MockJob implements IJob {
    
        private boolean cancelled;
    
        public MockJob() {
        }
    
        @Override
        public IResult execute() {
            int count = 0;
            while (!cancelled) {
                try {
                    count++;
                    System.out.println("Mock Job Thread: count = " + count);
                    if (count >= 10) {
                        break;
                    }
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    cancelled = true;
                }
            }
            return new MockResult(cancelled, count);
        }
    
        @Override
        public void cancel() {
            cancelled = true;
        }
    }
    

    模拟结果:

    public class MockResult implements IResult {
    
        private boolean cancelled;
        private int result;
    
        public MockResult(boolean cancelled, int result) {
            this.cancelled = cancelled;
            this.result = result;
        }
    
        public boolean isCancelled() {
            return cancelled;
        }
    
        public int getResult() {
            return result;
        }
    }
    

    主类:

    public class Main {
    
        public static void main(String[] args) throws InterruptedException {
            // Job
            IJob mockJob = new MockJob();
    
            // Async task
            FutureTask<IResult> asyncTask = new FutureTask<>(mockJob::execute);
            Thread mockJobThread = new Thread(asyncTask);
    
            // Show result
            Thread showResultThread = new Thread(() -> {
                try {
                    IResult result = asyncTask.get();
                    MockResult mockResult = (MockResult) result;
                    Thread thread = Thread.currentThread();
                    System.out.println(String.format("%s: isCancelled = %s, result = %d",
                            thread.getName(),
                            mockResult.isCancelled(),
                            mockResult.getResult()
                    ));
                } catch (InterruptedException | ExecutionException ex) {
                    // NO-OP
                }
            });
    
            // Check status
            Thread monitorThread = new Thread(() -> {
                try {
                    while (!asyncTask.isDone()) {
                        Thread thread = Thread.currentThread();
                        System.out.println(String.format("%s: asyncTask.isDone = %s",
                                thread.getName(),
                                asyncTask.isDone()
                        ));
                        Thread.sleep(1000);
                    }
                } catch (InterruptedException ex) {
                    // NO-OP
                }
                Thread thread = Thread.currentThread();
                System.out.println(String.format("%s: asyncTask.isDone = %s",
                        thread.getName(),
                        asyncTask.isDone()
                ));
            });
    
            // Async cancel
            Thread cancelThread = new Thread(() -> {
                try {
                    // Play with this Thread.sleep, set to 15000
                    Thread.sleep(5000);
                    if (!asyncTask.isDone()) {
                        Thread thread = Thread.currentThread();
                        System.out.println(String.format("%s: job.cancel()",
                                thread.getName()
                        ));
                        mockJob.cancel();
                    }
                } catch (InterruptedException ex) {
                    // NO-OP
                }
            });
    
            monitorThread.start();
            showResultThread.start();
            cancelThread.setDaemon(true);
            cancelThread.start();
            mockJobThread.start();
        }
    }
    

    输出(Thread.sleep(5000)):

    Thread-2: asyncTask.isDone = false
    Thread-0: count = 1
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 2
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 3
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 4
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 5
    Thread-3: job.cancel()
    Thread-2: asyncTask.isDone = false
    Thread-1: isCancelled = true, result = 5
    Thread-2: asyncTask.isDone = true
    

    输出(Thread.sleep(15000)):

    Thread-2: asyncTask.isDone = false
    Thread-0: count = 1
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 2
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 3
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 4
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 5
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 6
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 7
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 8
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 9
    Thread-2: asyncTask.isDone = false
    Thread-0: count = 10
    Thread-1: isCancelled = false, result = 10
    Thread-2: asyncTask.isDone = true
    

    【讨论】:

    • Ahhh 我没有考虑直接从主线程调用 job.cancel(),我当然可以,因为它有一个句柄。轻微的皱纹是,在插件作业运行后我需要处理结果(基本上,写输出),这部分作业也需要是可取消的。不过,我想我可以处理这个问题......
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-11-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-04-16
    相关资源
    最近更新 更多