我在这个主题上看过其他一些帖子,似乎没有一个与我遇到的问题完全相同.但是这里:
我正在使用并行运行一个函数
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) 在上市波纹管,我希望当我打电话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) 根据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) 我在互联网上看到了一些例子,为了使用流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网站,它不使用线程/异步。登录需要花费很多时间,因为它对不同的微服务进行了十二个调用,但一次都同步进行:每个调用都等待另一个完成,然后再进行下一个调用。但是,API调用中的任何一个都不取决于其他任何一个的结果。
但是我们确实需要获得所有结果并将其合并,然后再继续。看起来确实很明显,我们应该能够并行进行这十二个调用,但是要等到它们全部完成后,才能在接下来的步骤中使用它们的数据。
因此,在调用之前和之后,一切都是同步的。但是,最好将它们各自并行,异步(或只是并行)发送出去,然后我们仅受单个最慢的调用的限制,而不是所有调用的总顺序时间。
我读过Java 8围绕着这套很棒的新操作CompletableFuture。但是我还没有在任何地方解释我的用法。我们不希望结果有希望-我们很高兴等到它们全部完成然后继续。JS具有Promise.all(),但即使如此,它也会返回一个承诺。
我能想到的就是在进行异步调用后稍等一下,直到我们得到所有结果后才继续。显然是疯了。
我在这里想念什么吗?因为对我来说似乎很明显,但是似乎没有人对此有问题-否则这种方法是如此简单,没人问,我只是不明白。
java parallel-processing asynchronous java-8 completable-future
我创建了一个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) 假设yo = Yo()是一个带有方法的大对象double,它返回其参数乘以2.
如果我通过yo.double到imap的multiprocessing,那么它是非常缓慢的,因为每一个函数调用创建一个副本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 …
我希望澄清我对.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",但据我所知,同样的事情发生在背景中关于分配到多处理器/核心.
任何人都可以帮助澄清情况,如果可能的话,链接到文档/文章.
我知道这听起来像是要求与同事一起解决争论,但我真的想澄清我是否正确理解这一点.
我有一个实时的Linux桌面应用程序(用C编写),我们正在移植到ARM(4核Cortex v8-A72 CPU).在架构上,它结合了高优先级显式pthread(其中6个)和一对GCD(libdispatch)工作队列(一个并发和另一个串行).
我的担忧有两个方面:
select声明中)parallel-processing operating-system arm multiprocessing grand-central-dispatch
我有一个项目,我在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) java ×2
python ×2
arm ×1
arrays ×1
asynchronous ×1
c# ×1
c++ ×1
c++11 ×1
chapel ×1
concurrency ×1
foreach ×1
haskell ×1
java-8 ×1
java-stream ×1
python-2.7 ×1
r ×1
tree ×1