【问题标题】:Better way to handle Uncaught Exceptions in ForkJoinPool Tasks/action在 ForkJoinPool 任务/操作中处理未捕获异常的更好方法
【发布时间】:2012-02-02 03:30:06
【问题描述】:

在使用ForkJoinPool 提交任务(RecursiveActionRecursiveTask)时,处理异常(未捕获)的更好方法是什么?

ForkJoinPool 接受 Thread.UncaughtExceptionHandler 在 WorkerThread 突然终止时处理异常(这无论如何都不受我们的控制),但是当 ForkJoinTask 抛出异常时,不使用此处理程序。我在实现中使用标准的submit/invokeAll 方式。

这是我的场景:

我有一个线程在无限循环中运行,从第 3 方系统读取数据。在这个线程中,我将任务提交给ForkJoinPool

new Thread() {
      public void run() {
         while (true) {
             ForkJoinTask<Void> uselessReturn = 
                   ForkJoinPool.submit(RecursiveActionTask);
         }
      }
 }

我正在使用 RecursiveAction,在少数情况下使用 RecursiveTask。这些任务使用submit() 方法提交给FJPool。 我想要一个类似于UncaughtExceptionHandler 的通用异常处理程序,如果任务抛出未经检查/未捕获的异常,我可以处理异常并在需要时重新提交任务。处理异常还可以确保如果一个/一些任务抛出异常,排队的任务不会被取消。

invokeAll() 方法返回一组 ForkJoinTasks 但这些 Tasks 在一个递归块中(每个任务调用compute() 方法并且可以进一步拆分[假设场景])

class RecursiveActionTask extends RecursiveAction {

    public void compute() {
       if <task.size() <= ACCEPTABLE_SIZE) {
          processTask() // this might throw an checked/unchecked exception
       } else {
          RecursiveActionTask[] splitTasks = splitTasks(tasks)
          RecursiveActionTasks returnedTasks = invokeAll(splitTasks);
          // the below code never executes as invokeAll submits the tasks to the pool 
          // and the flow never comes to the code below.
          // I am looking for some handling like this
          for (RecusiveActionTask task : returnedTasks) {
             if (task.isDone()) {
                task.getException() // handle this exception
             }
          }
       }
    }

}

我注意到当 3-4 个任务失败时,整个队列提交单元被丢弃。目前我在我个人不喜欢的 processTask 周围放置了一个try/catch。我正在寻找更通用的。

  1. 我还想知道所有失败的任务列表,以便我可以重新提交它们
  2. 当任务抛出异常时,线程是否会从池中被逐出(尽管我的分析发现它们没有 [但不确定])?
  3. 在 FutureTask 上调用 get() 方法更有可能让我的流程按顺序等待,直到任务完成。
  4. 我只想知道任务失败时的状态。我不在乎它什么时候完成(显然不想等一个小时后)

任何想法如何处理上述场景中的异常?

【问题讨论】:

    标签: java multithreading java-7 java.util.concurrent fork-join


    【解决方案1】:

    这就是我们在 Akka 中解决它的方法:

    /**
     * INTERNAL AKKA USAGE ONLY
     */
    final class MailboxExecutionTask(mailbox: Mailbox) extends ForkJoinTask[Unit] {
      final override def setRawResult(u: Unit): Unit = ()
      final override def getRawResult(): Unit = ()
      final override def exec(): Boolean = try { mailbox.run; true } catch {
        case anything ⇒
          val t = Thread.currentThread
          t.getUncaughtExceptionHandler match {
            case null ⇒
            case some ⇒ some.uncaughtException(t, anything)
          }
          throw anything
      }
     }
    

    【讨论】:

    • 我对 AKKA 很陌生,但如果我理解正确,您是否将代码抛出的异常设置为当前线程异常处理程序?如果是这样,线程是否会从池中删除,因为线程池执行器可能会检查此异常并将线程标记为坏线程。你能解释一下吗?
    【解决方案2】:

    @Rajendra,你做的一切都是对的,除了你应该使用 ForkJoinPool execute() 而不是 submit()。

    这样,如果 Runnable 失败,它将强制工作异常,并且它将被您的 UncaughtExceptionHandler 捕获。

    我不知道为什么会存在这种行为,但它会起作用!我已经很难学会了:(

    取自 Java 8 代码:
    提交正在使用 AdaptedRunnableAction()。
    执行正在使用 RunnableExecuteAction()(参见 rethrow(ex))。

     /**
     * Adaptor for Runnables without results
     */
    static final class AdaptedRunnableAction extends ForkJoinTask<Void>
        implements RunnableFuture<Void> {
        final Runnable runnable;
        AdaptedRunnableAction(Runnable runnable) {
            if (runnable == null) throw new NullPointerException();
            this.runnable = runnable;
        }
        public final Void getRawResult() { return null; }
        public final void setRawResult(Void v) { }
        public final boolean exec() { runnable.run(); return true; }
        public final void run() { invoke(); }
        private static final long serialVersionUID = 5232453952276885070L;
    }
    
    /**
     * Adaptor for Runnables in which failure forces worker exception
     */
    static final class RunnableExecuteAction extends ForkJoinTask<Void> {
        final Runnable runnable;
        RunnableExecuteAction(Runnable runnable) {
            if (runnable == null) throw new NullPointerException();
            this.runnable = runnable;
        }
        public final Void getRawResult() { return null; }
        public final void setRawResult(Void v) { }
        public final boolean exec() { runnable.run(); return true; }
        void internalPropagateException(Throwable ex) {
            rethrow(ex); // rethrow outside exec() catches.
        }
        private static final long serialVersionUID = 5232453952276885070L;
    }
    

    【讨论】:

      猜你喜欢
      • 2023-03-30
      • 2017-03-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-11-25
      • 2012-06-07
      • 1970-01-01
      相关资源
      最近更新 更多