【问题标题】:Catch exception from parallel stream从并行流中捕获异常
【发布时间】:2017-03-26 18:33:27
【问题描述】:

我有一堆列作为 csv 文件中的字符串数组。现在我想解析它们。由于这种解析需要日期解析和其他不那么快的解析技术,我正在考虑并行性(我计时了,这需要一些时间)。我的简单方法:

Stream.of(columns).parallel().forEach(column -> 
    result[column.index] = parseColumn(valueCache[column.index], column.type));

Columns 包含ColumnDescriptor 元素,它只有两个属性,要解析的列索引和定义如何解析它的类型。没有其他的。 result 是一个 Object 数组,它采用结果数组。

现在的问题是 parse 函数抛出一个 ParseException,我在调用堆栈中处理得更远。由于我们在这里是并行的,它不能只是被抛出。处理这个问题的最佳方法是什么?

我有这个解决方案,但我有点畏缩阅读它。有什么更好的方法?

final CompletableFuture<ParseException> thrownException = new CompletableFuture<>();
Stream.of(columns).parallel().forEach(column -> {
    try {
        result[column.index] = parseColumn(valueCache[column.index], column.type);
    } catch (ParseException e) {
        thrownException.complete(e);
    }});

if(thrownException.isDone())
    //only can be done if there is a value set.
    throw thrownException.getNow(null);

注意:我不需要所有的例外。如果我按顺序解析它们,我也只会得到一个。所以没关系。

【问题讨论】:

  • 对我来说,其他可能性是可读的>
  • 我想这更像是一个实验,因为与从 CSV 文件读取相比,解析不太可能占用大量时间(即,这是过早的优化)。如果您在读取文件的同时进行解析,您可能会发现所有内容都在读取完成的同时被解析。
  • 读取数据是一次性操作。但是解析将使用不同的设置重复进行。这就是为什么我要优化那部分。
  • 为什么你的“简单方法”不起作用?除了声称你做不到之外,你没有解释这一点。

标签: java exception-handling parallel-processing java-8 java-stream


【解决方案1】:

问题是你的错误前提“因为我们在这里是并行的,所以不能随便抛出。”没有规范禁止在并行处理中抛出异常。您可以像在顺序流中一样在并行流中抛出该异常,如果它是已检查异常,则将其包装在未检查异常中。

如果线程中至少抛出一个异常,forEach 调用会将其(或其中之一)传播给调用者。

您可能遇到的唯一问题是,当前实现在遇到异常时不会等待所有线程完成。这可以解决使用

try {
    Arrays.stream(columns).parallel()
        .forEach(column -> 
            result[column.index] = parseColumn(valueCache[column.index], column.type));
} catch(Throwable t) {
    ForkJoinPool.commonPool().awaitQuiescence(1, TimeUnit.MINUTES);
    throw t;
}

但通常情况下,您不需要它,因为在特殊情况下您不会访问并发处理的结果。

【讨论】:

  • 首先感谢您花时间回答这个问题。这将导致与 john16384 答案中的相同点。我不喜欢在多线程环境中抛出异常。当我学习使用线程时,我学到了多个线程上的异常冒泡是不行的。所以即使它有效,我也不是很满意。
  • @findusl:这会很有趣,谁告诉你有什么理由,因为在你的余生中避免某事,只是因为某天有人说了一些坏话,听起来很教条。你试图避免它甚至不会改变任何东西,无论是语义上还是技术上。 parseColumn 仍然在多线程执行中抛出异常,有人会捕获它并将其交给作业启动线程。为什么手动做会更好,而不是让 Stream 框架做呢?
  • 我能看到的一个问题是流打破了检查异常气泡流,从这个角度来看,这可能有点麻烦。
  • @FrankHopkins 但这与顺序流没有什么不同。
  • @Holger 是的,完全不反对您的回答。这只是我能看到导致建立这样一个“规则”的一件事,或者通常会导致人们试图避免这两个概念的结合。不过,正如您所说,它更多地与流有关,而不是与并行化有关。
【解决方案2】:

我觉得问题比较多,串行解析的时候一般是怎么做的呢?

您是否在第一个异常处停止,并停止整个过程?在这种情况下,将异常包装在运行时异常中,让流中止并抛出它。捕获包装器异常,解包并处理。

您会跳过不良记录吗?然后要么 1. 在某个地方跟踪 List 中的错误,要么 2. 创建一个可以保存解析结果或错误的包装器对象(不要跟踪异常本身,只跟踪描述错误所需的最小值)。

之后检查第一个选项的列表中是否有错误,或者为第二个选项显示不同的错误记录。

【讨论】:

  • 运行时异常确实有效。我喜欢它胜过我的解决方案。仍然不太喜欢它,因为我通过这样的并行线程抛出了一个 RuntimeException 。我学会了在并行编程时要避免的事情。但是由于流为我处理它,我想如果这些天我没有得到更好的答案,我会标记你的正确:)
猜你喜欢
  • 2023-03-26
  • 2012-07-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-07-13
  • 2016-10-31
  • 1970-01-01
相关资源
最近更新 更多