标签: parallel-processing

并行错误R:序列化错误(数据,节点$ con):写入连接时出错

我在这个主题上看过其他一些帖子,似乎没有一个与我遇到的问题完全相同.但是这里:

我正在使用并行运行一个函数

cores <- detectCores() cl <- makeCluster(8L,outfile="output.txt") registerDoParallel(cl) x <- foreach(i = 1:length(y), .combine='list',.packages=c('httr','jsonlite'), .multicombine=TRUE,.verbose=F,.inorder=F) %dopar% {function(y[i])}

这通常工作正常,但现在抛出错误:

序列化错误(数据,节点$ con):写入连接时出错

检查output.txt文件后,我看到:

starting worker pid=11112 on localhost:11828 at 12:38:32.867
starting worker pid=10468 on localhost:11828 at 12:38:33.389
starting worker pid=4996 on localhost:11828 at 12:38:33.912
starting worker pid=3300 on localhost:11828 at 12:38:34.422
starting worker pid=10808 on localhost:11828 at 12:38:34.937
starting worker pid=5840 on localhost:11828 at 12:38:35.435
starting worker pid=8764 on localhost:11828 at 12:38:35.940
starting worker pid=7384 on localhost:11828 at 12:38:36.448
Error in …
Run Code Online (Sandbox Code Playgroud)

parallel-processing foreach r

5
推荐指数
2
解决办法
3538
查看次数

为什么我的线程不能在后台运行?

在上市波纹管,我希望当我打电话t.detach()时创建线程行之后,该线程t将在后台运行,而printf("quit the main function now \n")会叫,然后main将退出.

#include <thread>
#include <iostream>

void hello3(int* i)
{

    for (int j = 0; j < 100; j++)
    {
        *i = *i + 1;
        printf("From new thread %d \n", *i);
        fflush(stdout);

    }

    char c = getchar();
 }

int main()
{
    int i;
    i = 0;
    std::thread t(hello3, &i);
    t.detach();
    printf("quit the main function now \n");
    fflush(stdout);
    return 0;
}
Run Code Online (Sandbox Code Playgroud)

然而,从它在屏幕上打印的内容来看并非如此.它打印

From new thread 1
From new thread …
Run Code Online (Sandbox Code Playgroud)

c++ parallel-processing multithreading c++11

5
推荐指数
1
解决办法
671
查看次数

有没有办法在Chapel中自定义整个数组语句的默认并行化行为?

根据Chapel的可用文档,(整个)数组语句如

A = B + alpha * C;   // with A, B, and C being arrays, and alpha some scalar
Run Code Online (Sandbox Code Playgroud)

在语言中实现如下的forall迭代:

forall (a,b,c) in zip(A,B,C) do
   a = b + alpha * c;
Run Code Online (Sandbox Code Playgroud)

因此,默认情况下,数组语句由并行线程团队执行.不幸的是,这似乎也完全排除了这些陈述的(部分或完全)矢量化.对于习惯于Fortran或Python/Numpy等语言的程序员而言,这可能会带来性能上的惊喜(默认行为通常是只对数组语句进行矢量化).

对于使用(整数)数组语句与小到中等大小的数组的代码,矢量化的丢失(由Linux硬件性能计数器确认)和并行线程固有的显着开销(不适合有效地利用细粒度数据) - 在这些问题中可用的并行性)可能导致性能的显着损失.作为一个例子,考虑以下版本的Jacobi迭代,它们都解决了300 x 300区域的相同问题:

Jacobi_1 使用数组语句,如下所示:

/*
 *  Jacobi_1
 *
 *  This program (adapted from the Chapel distribution) performs
 *  niter iterations of the Jacobi method for the Laplace equation
 *  using (whole-)array statements.
 *
 */

config var n = 300;                  // size of n …
Run Code Online (Sandbox Code Playgroud)

arrays parallel-processing chapel

5
推荐指数
1
解决办法
123
查看次数

简单地调用parallelStream是否并行运行任务?

我在互联网上看到了一些例子,为了使用流API来做并行的东西,只需调用这样的.parallelStream()方法:

mySet
    .parallelStream()
    ... // do my fancy stuff and collect
Run Code Online (Sandbox Code Playgroud)

但在其他情况下,我已经看到并行流在线程池子目录中使用,如下所示:

ForkJoinPool.commonPool().submit(() -> {
    mySet
        .parallelStream()
        ... // do my fancy stuff and collect
})
Run Code Online (Sandbox Code Playgroud)

只是调用parallelStream()执行多个并发线程中接下来的内容吗?就像在一些预配置的线程池或其他东西.或者我是否必须创建我的线程然后使用并行流?

java parallel-processing concurrency java-stream

5
推荐指数
1
解决办法
178
查看次数

您如何等待所有异步调用在Java中完成?

从概念上讲非常简单。我们有一个庞大的旧版Java网站,它不使用线程/异步。登录需要花费很多时间,因为它对不同的微服务进行了十二个调用,但一次都同步进行:每个调用都等待另一个完成,然后再进行下一个调用。但是,API调用中的任何一个都不取决于其他任何一个的结果。

但是我们确实需要获得所有结果并将其合并,然后再继续。看起来确实很明显,我们应该能够并行进行这十二个调用,但是要等到它们全部完成后,才能在接下来的步骤中使用它们的数据。

因此,在调用之前和之后,一切都是同步的。但是,最好将它们各自并行,异步(或只是并行)发送出去,然后我们仅受单个最慢的调用的限制,而不是所有调用的总顺序时间。

我读过Java 8围绕着这套很棒的新操作CompletableFuture。但是我还没有在任何地方解释我的用法。我们不希望结果有希望-我们很高兴等到它们全部完成然后继续。JS具有Promise.all(),但即使如此,它也会返回一个承诺。

我能想到的就是在进行异步调用后稍等一下,直到我们得到所有结果后才继续。显然是疯了。

我在这里想念什么吗?因为对我来说似乎很明显,但是似乎没有人对此有问题-否则这种方法是如此简单,没人问,我只是不明白。

java parallel-processing asynchronous java-8 completable-future

5
推荐指数
1
解决办法
2081
查看次数

顺序和并行编程之间的解决方案的差异

我创建了一个python代码,解决了一个组套索惩罚线性模型.对于那些不习惯使用这些模型的人来说,基本的想法是你输入数据集(x)和响应变量(y),以及参数(lambda1)的值,改变值此参数更改模型的解决方案.所以我决定使用多处理库并解决不同的模型(与不同的参数值相关联).我创建了一个名为"model.py"的python文件,其中包含以下函数:

# -*- coding: utf-8 -*-
from __future__ import division
import functools
import multiprocessing as mp
import numpy as np
from cvxpy import *

def lm_gl_preprocessing(x, y, index, lambda1=None):
    lambda_vector = [lambda1]
    m = x.shape[1]
    n = x.shape[0]
    lambda_param = Parameter(sign="positive")
    m = m+1
    index = np.append(0, index)
    x = np.c_[np.ones(n), x]
    group_sizes = []
    beta_var = []
    unique_index = np.unique(index)
    for idx in unique_index:
        group_sizes.append(len(np.where(index == idx)[0]))
        beta_var.append(Variable(len(np.where(index == idx)[0])))
    num_groups = len(group_sizes)
    group_lasso_penalization = 0
    model_prediction = x[:, …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing multiprocessing python-2.7

5
推荐指数
1
解决办法
253
查看次数

将一个大对象的方法传递给imap:通过包装方法加速1000倍

假设yo = Yo()是一个带有方法的大对象double,它返回其参数乘以2.

如果我通过yo.doubleimapmultiprocessing,那么它是非常缓慢的,因为每一个函数调用创建一个副本yo,我认为.

即,这很慢:

from tqdm import tqdm
from multiprocessing import Pool
import numpy as np


class Yo:
    def __init__(self):
        self.a = np.random.random((10000000, 10))

    def double(self, x):
        return 2 * x

yo = Yo()    

with Pool(4) as p:
    for _ in tqdm(p.imap(yo.double, np.arange(1000))):
        pass
Run Code Online (Sandbox Code Playgroud)

输出:

0it [00:00, ?it/s]
1it [00:06,  6.54s/it]
2it [00:11,  6.17s/it]
3it [00:16,  5.60s/it]
4it [00:20,  5.13s/it]
Run Code Online (Sandbox Code Playgroud)

...

但是,如果我yo.double用函数包装double_wrap …

python parallel-processing multiprocessing

5
推荐指数
1
解决办法
75
查看次数

线程可以在Task.Factory.StartNew和Parallel.Invoke的不同处理器或内核上运行

我希望澄清我对.NET多线程的理解,特别是哪些.NET方法创建的线程可能会在多处理器/核心系统中的不同处理器或内核上同时执行.

在.NET TPL框架中,您可以使用Parallel.Invoke或Task.Factory.StartNew方法来实现某种并行性.

我的理解是,在这两种情况下.NET都会创建新的任务(在Parallel.Invoke的幕后),.NET环境然后在后台分配给托管线程,然后将其分配给线程,CPU可以分配给不同的线程核心或处理器取决于工作负载.这两种方法的主要区别在于语义 - Parallel.Invoke执行多个任务并等待它们完成; Task.Factory.StartNew在后台启动一个新任务.在这两种情况下,实际工作可以在不同的核心或处理器上完成.根据任务并行库(TPL).

我有一位同事确信只有Parallel.Invoke方法允许线程在不同的核心/处理器上执行,而Task.Factory.StartNew启动一个新线程但该线程只能在一个核心/处理器上调度 - 所以实际上并没有给出并行性.

我找不到任何明确说明是否属于这种情况的文件或文章.我的同事向我介绍了我正在查看的相同文章,例如基于任务的异步编程,我认为这可以验证我的理解,但我的同事认为验证了他的.

文档有时使用术语"并行处理"参考Parallel.Invoke和"异步任务"参考"Task.Factory.StartNew",但据我所知,同样的事情发生在背景中关于分配到多处理器/核心.

任何人都可以帮助澄清情况,如果可能的话,链接到文档/文章.

我知道这听起来像是要求与同事一起解决争论,但我真的想澄清我是否正确理解这一点.

c# parallel-processing multithreading task-parallel-library

5
推荐指数
1
解决办法
328
查看次数

重型线程消耗对ARM(4核A72)与x86(2核i5)的影响

我有一个实时的Linux桌面应用程序(用C编写),我们正在移植到ARM(4核Cortex v8-A72 CPU).在架构上,它结合了高优先级显式pthread(其中6个)和一对GCD(libdispatch)工作队列(一个并发和另一个串行).

我的担忧有两个方面:

  • 我听说ARM没有超越x86的方式,因此我的4核已经是上下文切换以跟上我的6 pthreads(和后台进程).我应该从中得到什么样的性能损失?
    • 我听说我应该期望这些ARM上下文切换效率低于x86.真的吗?
    • 一些pthreads是针对相当罕见的事件的高优先级处理程序,这会改变前景吗?(即他们坐在select声明中)
  • 我更大的担忧来自GCD在这个应用程序中的影响.我对GCD内部工作原理的理解是,它是一个与调度程序交互的动态扩展线程池,并将尝试添加更多线程以适应负载.在我看来,这对我的情景中的性能几乎完全是负面影响.(核心被完全消耗的系统中的IE)正确吗?

parallel-processing operating-system arm multiprocessing grand-central-dispatch

5
推荐指数
1
解决办法
342
查看次数

在Haskell中并行构造树的策略

我有一个项目,我在Haskell中构建一个决策树.生成的树将具有多个彼此独立的分支,因此我认为它们可以并行构建.

DecisionTree数据类型被限定如下所示:

data DecisionTree =
    Question Filter DecisionTree DecisionTree |    
    Answer DecisionTreeResult

instance NFData DecisionTree where
    rnf (Answer dtr)            = rnf dtr
    rnf (Question fil dt1 dt2)  = rnf fil `seq` rnf dt1 `seq` rnf dt2
Run Code Online (Sandbox Code Playgroud)

这是构造树的算法的一部分

constructTree :: TrainingParameters -> [Map String Value] -> Filter -> Either String DecisionTree    
constructTree trainingParameters trainingData fil =    
    if informationGain trainingData (parseFilter fil) < entropyLimit trainingParameters    
    then constructAnswer (targetVariable trainingParameters) trainingData    
    else
        Question fil <$> affirmativeTree <*> negativeTree `using` …
Run Code Online (Sandbox Code Playgroud)

parallel-processing tree haskell

5
推荐指数
1
解决办法
163
查看次数