标签: parallel-processing

使用Akka进行分叉和连接

问题陈述:我有一系列需要以并行方式处理的证券.在Java中,我使用线程池来处理每个安全性,并使用锁存器来倒计时.完成后我会做一些合并等.

所以我给我的SecurityProcessor(这是一个演员)发了消息,等待所有的期货完成.最后,我使用MergeHelper进行后处理.SecurityProcessor采用安全性,进行一些I/O处理并回复安全性

  val listOfFutures = new ListBuffer[Future[Security]]()
  var portfolioResponse: Portfolio = _
  for (security <- portfolio.getSecurities.toList) {
    val securityProcessor = actorOf[SecurityProcessor].start()
    listOfFutures += (securityProcessor ? security) map {
      _.asInstanceOf[Security]
    }
  }
  val futures = Future.sequence(listOfFutures.toList)
  futures.map {
    listOfSecurities =>
      portfolioResponse = MergeHelper.merge(portfolio, listOfSecurities)
  }.get
Run Code Online (Sandbox Code Playgroud)

这个设计是否正确,是否有更好/更酷的方式来使用akka实现这个常见问题?

parallel-processing scala future actor akka

7
推荐指数
1
解决办法
1302
查看次数

使用foreach包将行附加到dataframe

我在使用并行处理将值附加到数据框时遇到问题.

我有一个函数会做一些计算并返回一个数据帧,包括这些计算是一个随机抽样.

所以我做的是:

randomizex <- function(testdf)
{
    foreach(ind=1:1000)%dopar%
    {
      testdf$X = sample(testdf$X,nrow(testdf), replace=FALSE)
      fit = lm(X ~ Y, testdf)
      newdf <- rbind(newdf, data.frame(pc=ind, err=sum(residuals(fit)^2) ))

    }

return(newdf)
}
resdf = randomizex(mydf)
Run Code Online (Sandbox Code Playgroud)

当我查看结果时resdf,它是空的

如果我更换%dopar%%do%结果被正确地计算,但它太慢了..

反正有没有提高这一点?

parallel-processing foreach r

7
推荐指数
2
解决办法
8186
查看次数

C#任务是否在一个核心上运行?

C#任务是否在一个核心上运行?

我有一个项目,我需要决定要创建多少个任务.我需要创建尽可能多的计算机.这是处理器,核心或逻辑处理器的数量,我在三个选项之间感到困惑.

c# parallel-processing task-parallel-library

7
推荐指数
1
解决办法
2586
查看次数

从R中的并行进程写入文件时锁定文件

我用parSapply()parallel包河,我需要对大量的数据进行计算.即使并行执行也需要数小时,因此我决定定期将结果写入集群中的文件write.table(),因为当内存不足或其他一些随机原因导致进程崩溃时,我想继续计算把它停下来.我注意到我得到的一些csv文件行只是在中间切割,可能是由于多个进程同时写入文件.有没有办法在write.table()执行时暂时锁定文件,因此其他集群无法访问它,或者唯一的出路是从每个集群写入单独的文件然后合并结果?

parallel-processing r file-locking filelock

7
推荐指数
1
解决办法
1202
查看次数

并行I/O - 为什么它可以工作?

我有一个python函数,它从文本文件中读取一行并将其写入另一个文本文件.它会对文件中的每一行重复此操作.实质上:

Read line 1 -> Write line 1 -> Read line 2 -> Write line 2...
Run Code Online (Sandbox Code Playgroud)

等等.

我可以使用队列来传递数据来并行化这个过程,所以它更像是:

Read line 1 -> Read line 2 -> Read line 3...

              Write line 1 -> Write line 2....
Run Code Online (Sandbox Code Playgroud)

我的问题是 - 为什么这样做(因为我为什么加快速度?).听起来像是一个愚蠢的问题,但我在想 - 当然我的硬盘一次只能做一件事吗?那么为什么没有一个过程被搁置直到另一个过程完成?

当用高级语言写作时,这样的事情对用户是隐藏的.我想知道什么是低级别的?

python io parallel-processing

7
推荐指数
1
解决办法
1000
查看次数

SIMD XOR操作不如Integer XOR有效吗?

我有一个任务来计算数组中的xor-sum字节:

X = char1 XOR char2 XOR char3 ... charN;
Run Code Online (Sandbox Code Playgroud)

我正在尝试并行化它,而是使用__m128.这应该加速因子4.另外,要重新检查算法,我使用int.这应该加速因子4.测试程序是100行,我不能让它更短,但它很简单:

#include "xmmintrin.h" // simulation of the SSE instruction
#include <ctime>

#include <iostream>
using namespace std;

#include <stdlib.h> // rand

const int NIter = 100;

const int N = 40000000; // matrix size. Has to be dividable by 4.
unsigned char str[N] __attribute__ ((aligned(16)));

template< typename T >
T Sum(const T* data, const int N)
{
    T sum = 0;
    for ( int i = 0; i < N; ++i …
Run Code Online (Sandbox Code Playgroud)

c++ parallel-processing performance simd seeding

7
推荐指数
2
解决办法
1182
查看次数

我可以在ConcurrentBag上使用普通的foreach吗?

在代码的并行部分中,我将每个线程的结果保存到ConcurrentBag中。但是,完成此操作后,我需要遍历每个结果并通过我的评估算法运行它们。普通的foreach实际上会遍历所有成员,还是我需要特殊的代码?我还考虑过使用队列之类的东西来代替行李,但我不知道哪种方法最好。在并行代码的末尾,袋子通常只包含20个左右的项目。

即,实际上将为ConcurrentBag的所有成员访问并运行foreach吗?

ConcurrentBag futures = new ConcurrentBag();
foreach(move in futures)
{
 // stuff
}
Run Code Online (Sandbox Code Playgroud)

c# parallel-processing foreach parallel.foreach

7
推荐指数
1
解决办法
4739
查看次数

多处理错误的管道错误.Queue

在python2.7中,multiprocessing.Queue在从函数内部初始化时抛出一个破坏的错误.我提供了一个重现问题的最小例子.

#!/usr/bin/python
# -*- coding: utf-8 -*-

import multiprocessing

def main():
    q = multiprocessing.Queue()
    for i in range(10):
        q.put(i)

if __name__ == "__main__":
    main()
Run Code Online (Sandbox Code Playgroud)

抛出下面破裂的管道错误

Traceback (most recent call last):
File "/usr/lib64/python2.7/multiprocessing/queues.py", line 268, in _feed
send(obj)
IOError: [Errno 32] Broken pipe

Process finished with exit code 0
Run Code Online (Sandbox Code Playgroud)

我无法破译原因.我们无法从函数内部填充Queue对象,这当然很奇怪.

python parallel-processing multiprocessing python-2.7

7
推荐指数
2
解决办法
1万
查看次数

Python多处理:为什么较大的chunksize较慢?

我一直在使用Python的多处理模块分析一些代码('job'函数只是对数字进行平方).

data = range(100000000)
n=4
time1 = time.time()
processes = multiprocessing.Pool(processes=n)
results_list = processes.map(func=job, iterable=data, chunksize=10000)
processes.close()
time2 = time.time()
print(time2-time1)
print(results_list[0:10])
Run Code Online (Sandbox Code Playgroud)

我发现奇怪的一件事是最佳的chunksize似乎是大约10k元素 - 这在我的计算机上花了16秒.如果我将chunksize增加到100k或200k,那么它会减慢到20秒.

这种差异可能是由于长时间列表中酸洗所需的时间更长吗?100个元素的块大小需要62秒,我假设是由于在不同进程之间来回传递块所需的额外时间.

python parallel-processing multiprocessing python-multiprocessing

7
推荐指数
1
解决办法
1390
查看次数

同时运行的多个Python实例限制为35个

我在并行计算集群的不同处理器上运行Python 3.6脚本作为多个单独的进程.最多35个进程同时运行没有问题,但第36行(以及更多)在第二行崩溃并出现分段错误import pandas as pd.有趣的是,第一行import os不会引起问题.完整的错误消息是:

OpenBLAS blas_thread_init: pthread_create: Resource temporarily unavailable
OpenBLAS blas_thread_init: RLIMIT_NPROC 1024 current, 2067021 max
OpenBLAS blas_thread_init: pthread_create: Resource temporarily unavailable
OpenBLAS blas_thread_init: RLIMIT_NPROC 1024 current, 2067021 max
OpenBLAS blas_thread_init: pthread_create: Resource temporarily unavailable
OpenBLAS blas_thread_init: RLIMIT_NPROC 1024 current, 2067021 max
OpenBLAS blas_thread_init: pthread_create: Resource temporarily unavailable
OpenBLAS blas_thread_init: RLIMIT_NPROC 1024 current, 2067021 max
OpenBLAS blas_thread_init: pthread_create: Resource temporarily unavailable
OpenBLAS blas_thread_init: RLIMIT_NPROC 1024 current, 2067021 max
OpenBLAS blas_thread_init: pthread_create: Resource temporarily …
Run Code Online (Sandbox Code Playgroud)

python linux parallel-processing multiprocessing openblas

7
推荐指数
1
解决办法
1512
查看次数