标签: parallel-processing

当Java 8 Stream抛出RuntimeException时,预期的行为是什么?

遇到RuntimeException流处理时,流处理是否应该中止?应该先完成吗?是否应重新抛出异常Stream.close()?异常是否被重新抛出或被包裹?Java Streamjava.util.stream包的JavaDoc 无话可说.

我发现的所有关于Stackoverflow的问题似乎都集中在如何从功能界面中包装已检查的异常以便编译代码.事实上,互联网上的博客文章和类似文章都集中在同一个警告上.这不是我关心的问题.

我根据自己的经验知道,一旦抛出一个顺序流的处理就会中止,RuntimeException并且这个异常会按原样重新抛出.仅当客户端线程抛出异常时,对于并行流才是相同的.

但是,此处的示例代码演示了如果在并行流处理期间由"工作线程"(=与调用终端操作的线程不同的线程)抛出异常,则此异常将永远丢失并且流处理完成.

示例代码将首先IntStream并行运行.然后是"正常" Stream并行.

这个例子将表明,

1)IntStream如果遇到RuntimeException,则中止并行处理没有问题.异常被重新抛出,包装在另一个RuntimeException中.

2)Stream不好玩.实际上,客户端线程永远不会看到抛出的RuntimeException的痕迹.流不仅完成处理; 将处理更多元素而不是limit()指定的元素!

在示例代码中,IntStream使用IntStream.range()生成."正常" Stream没有"范围"的概念,而是由1:s组成,但调用Stream.limit()将流限制为10亿个元素.

这是另一个转折点.生成IntStream的示例代码执行如下操作:

IntStream.range(0, 1_000_000_000).parallel().forEach(..)
Run Code Online (Sandbox Code Playgroud)

将其更改为生成的流,就像代码中的第二个示例一样:

IntStream.generate(() -> 1).limit(1_000_000_000).parallel().forEach(..)
Run Code Online (Sandbox Code Playgroud)

此IntStream的结果是相同的:异常被包装并重新抛出,处理中止.但是,第二个流现在也将包装并重新抛出异常,而不是处理超出限制的元素!因此:更改第一个流的生成方式会对第二个流的行为产生副作用.对我来说,这很奇怪.

ForkJoinPool.invoke()的 JavaDoc 并ForkJoinTask表示异常被重新抛出,这就是我对并行流的期望.

背景

当处理并行流中的元素时,我遇到了这个"问题" Collection.stream().parallel()(我还没有验证它的行为,Collection.parallelStream()但它应该是相同的).发生的事情是"工作线程"崩溃然后静静地离开,而所有其他线程成功完成了流.我的应用程序使用默认的异常处理程序将异常写入日志文件.但是甚至没有创建这个日志文件.线程和他的例外根本就消失了.由于我需要在捕获运行时异常时立即中止,因此一种替代方法是编写将此异常泄漏给其他工作程序的代码,使其在任何其他线程抛出异常时不愿意继续.当然,这并不能保证流实现只是继续生成尝试完成流的新线程.所以我可能最终不会使用并行流,而是使用线程池/执行器进行"正常"并发编程.

这表明运行时异常丢失的问题不会与使用流Stream.generate()或流生成的流隔离Stream.limit().最重要的是,我很想知道...是预期的行为?

java parallel-processing multithreading runtimeexception java-stream

4
推荐指数
1
解决办法
871
查看次数

是否有可能将此for循环并行化?

我得到了一些使用OpenMP进行并行化的代码,并且在各种函数调用中,我注意到这个for循环对计算时间有一些好的负罪感.

  double U[n][n];
  double L[n][n];
  double Aprime[n][n];
  for(i=0; i<n; i++) {
    for(j=0; j<n; j++) {
      if (j <= i) {
          double s;
          s=0;
          for(k=0; k<j; k++) {
            s += L[j][k] * U[k][i];
          } 
          U[j][i] = Aprime[j][i] - s;
      } else if (j >= i) {
          double s;
          s=0;
          for(k=0; k<i; k++) {
            s += L[j][k] * U[k][i];
          }
          L[j][i] = (Aprime[j][i] - s) / U[i][i];
      }
    }
Run Code Online (Sandbox Code Playgroud)

然而,在尝试并行化并在这里和那里应用一些信号量之后(没有运气),我开始认识到else if条件对早期的强烈依赖if(L[j][i]是一个处理过的数字U[i][i],可能在早期设置if …

c c++ parallel-processing openmp

4
推荐指数
1
解决办法
155
查看次数

并发与并行 - 特别是在C++中

我理解两者之间的基本区别,我经常在程序中使用std :: async,这给了我并发性.

是否有任何可靠/值得注意的库可以在C++中提供并行性?(我知道这可能是C++的一个特色17).如果是这样,您对它们的体验是什么?

谢谢!芭芭拉

c++ parallel-processing concurrency c++-standard-library

4
推荐指数
1
解决办法
1470
查看次数

akka-http查询不能并行运行

我对akka-http很陌生,在同一条路线上并行运行查询时遇到麻烦。

我有一条路线可能会非常快(如果已缓存)或没有(大量CPU多线程计算)返回结果。我想并行运行这些查询,以防一小段查询经过一长段计算之后又到达,我不希望第二个调用等待第一个调用完成。

但是,如果这些查询在同一路由上,则似乎不会并行运行(如果在不同的路由上,则并行运行)

我可以在一个基本项目中复制它:

并行调用服务器3次(使用http:// localhost:8080 / test上的3个Chrome选项卡),响应分别达到3.0秒,6.0秒和9.0秒。我想查询不会并行运行。

在具有jdk 8的Windows 10上的6核(带有HT)计算机上运行。

build.sbt

name := "akka-http-test"

version := "1.0"

scalaVersion := "2.11.8"

libraryDependencies += "com.typesafe.akka" %% "akka-http-experimental" % "2.4.11"
Run Code Online (Sandbox Code Playgroud)

* AkkaHttpTest.scala **

import java.util.concurrent.Executors

import akka.actor.ActorSystem
import akka.http.scaladsl.Http
import akka.http.scaladsl.server.Directives._
import akka.stream.ActorMaterializer

import scala.concurrent.{ExecutionContext, Future}

object AkkaHttpTest extends App {

  implicit val actorSystem = ActorSystem("system") // no application.conf here
  implicit val executionContext = 
                             ExecutionContext.fromExecutor(Executors.newFixedThreadPool(6))
  implicit val actorMaterializer = ActorMaterializer()

  val route = path("test") {
    onComplete(slowFunc()) { slowFuncResult =>
      complete(slowFuncResult) …
Run Code Online (Sandbox Code Playgroud)

parallel-processing scala akka akka-stream akka-http

4
推荐指数
1
解决办法
1316
查看次数

sizeof(MPI_INT)与sizeof(int)不同

我注意到int和double的大小与使用函数MPI_Type_size(MPI_INT,&MPI_INT_SIZE)计算的大小不同; 这是否意味着sizeof(MPI_INT)返回错误的值8 ?? 它通常应该是4谢谢你的回复

size parallel-processing types mpi

4
推荐指数
1
解决办法
1138
查看次数

Spark Scala - 如何对数据帧行进行分组并将复杂功能应用于组?

我正在努力解决这个超级简单的问题,我已经厌倦了,我希望有人可以帮我解决这个问题.我有一个像这样的形状的数据框:

---------------------------
|  Category  | Product_ID |
|------------+------------+
| a          | product 1  |
| a          | product 2  |
| a          | product 3  |
| a          | product 1  |
| a          | product 4  |
| b          | product 5  |
| b          | product 6  |
---------------------------

如何按类别对这些行进行分组并在Scala中应用复杂的功能?也许是这样的:

val result = df.groupBy("Category").apply(myComplexFunction)
Run Code Online (Sandbox Code Playgroud)

这个myComplexFunction应该为每个类别生成下表,并上传成Hive表中的成对相似性或将其保存到HDFS中:


+--------------------------------------------------+
|              | Product_1 | Product_2 | Product_3 |
+------------+------------+------------------------+
| Product_1    | 1.0       | 0.1       |    0.8    |
| Product_2    | 0.1 …

parallel-processing aggregate-functions dataframe apache-spark custom-function

4
推荐指数
1
解决办法
3303
查看次数

GridSearchCV的并行错误,适用于其他方法

我使用GridSearchCV遇到以下问题:它在使用时给我一个并行错误n_jobs > 1.同时n_jobs > 1与RadonmForestClassifier等单一模型一起工作正常.

下面是一个显示错误的简单工作示例:

train = np.random.rand(100,10)
targ = np.random.randint(0,2,100)

clf = ensemble.RandomForestClassifier(n_jobs = 2)
clf.fit(train,targ)
train = np.random.rand(100,10)
targ = np.random.randint(0,2,100)
?
clf = ensemble.RandomForestClassifier(n_jobs = 2)
clf.fit(train,targ)
Out[349]: RandomForestClassifier(bootstrap=True, class_weight=None,     criterion='gini',
            max_depth=None, max_features='auto', max_leaf_nodes=None,
            min_impurity_split=1e-07, min_samples_leaf=1,
            min_samples_split=2, min_weight_fraction_leaf=0.0,
            n_estimators=10, n_jobs=2, oob_score=False, random_state=None,
            verbose=0, warm_start=False)
Run Code Online (Sandbox Code Playgroud)

这个例子工作正常.

同时以下不起作用:

clf = ensemble.RandomForestClassifier()
param_grid = {'n_estimators': [10,20]}
grid_s= model_selection.GridSearchCV(clf, param_grid=param_grid_gb,n_jobs=-1,verbose=1)
grid_s.fit(train, targ)
Run Code Online (Sandbox Code Playgroud)

并给出以下错误:

Fitting 3 folds for each of 2 candidates, totalling 6 fits

ImportErrorTraceback (most …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing scikit-learn grid-search

4
推荐指数
1
解决办法
4950
查看次数

Swift 3 Parallel for/map Loop

这里有相当多的线程然而在Swift中使用Grand Central Dispatch来并行化并加速"for"循环?使用Swift <3.0代码而我无法在3中工作(参见代码).并行处理数组使用GCD使用指针并且它有点难看所以我要在这里断言我正在寻找好的Swift 3方法(当然尽可能高效).我也听说组很慢(?也许有人可以证实这一点.我也无法让小组工作.

这是我对跨步并行映射函数的实现(在Array的扩展中).它希望在全局队列上执行,以便不阻止UI.可能是并发位不需要在范围内,只需要余数循环.

extension Array {

    func parallelMap<R>(striding n: Int, f: @escaping (Element) -> R, completion: @escaping ([R]) -> ()) {
        let N = self.count
        var res = [R?](repeating: nil, count: self.count)
        DispatchQueue.global(qos: .userInitiated).async {
            DispatchQueue.concurrentPerform(iterations: N/n) {k in
                for i in (k * n)..<((k + 1) * n) {
                    res[i] = f(self[i]) //Error 1 here
                }
            }
            //remainder loop
            for i in (N - (N % n))..<N {
                res[i] = …
Run Code Online (Sandbox Code Playgroud)

parallel-processing grand-central-dispatch swift

4
推荐指数
1
解决办法
2119
查看次数

如何测量openmp中每个线程的执行时间?

我想衡量每个线程花费在执行代码块上的时间。我想看看我的负载平衡策略是否在工作人员之间平均分配块。通常,我的代码如下所示:

#pragma omp parallel for schedule(dynamic,chunk) private(i)
for(i=0;i<n;i++){
//loop code here
}
Run Code Online (Sandbox Code Playgroud)

更新我在gcc中使用openmp 3.1

c parallel-processing multithreading openmp

4
推荐指数
1
解决办法
2470
查看次数

使用CUDA流的优势

我试图了解流在哪里可以帮助我处理视频帧上的多个关注区域。如果使用支持流的NPP功能,是否会出现这样一种情况,即启动与ROI一样多的流?甚至可能为每个Stream创建一个CPU线程?还是使用一个流来处理所有ROI并可能使用来自CPU中多个线程的单个流的好处?

parallel-processing cuda emgucv managed-cuda opencv3.1

4
推荐指数
1
解决办法
2500
查看次数