问题陈述:我有一系列需要以并行方式处理的证券.在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实现这个常见问题?
我在使用并行处理将值附加到数据框时遇到问题.
我有一个函数会做一些计算并返回一个数据帧,包括这些计算是一个随机抽样.
所以我做的是:
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%结果被正确地计算,但它太慢了..
反正有没有提高这一点?
C#任务是否在一个核心上运行?
我有一个项目,我需要决定要创建多少个任务.我需要创建尽可能多的计算机.这是处理器,核心或逻辑处理器的数量,我在三个选项之间感到困惑.
我用parSapply()从parallel包河,我需要对大量的数据进行计算.即使并行执行也需要数小时,因此我决定定期将结果写入集群中的文件write.table(),因为当内存不足或其他一些随机原因导致进程崩溃时,我想继续计算把它停下来.我注意到我得到的一些csv文件行只是在中间切割,可能是由于多个进程同时写入文件.有没有办法在write.table()执行时暂时锁定文件,因此其他集群无法访问它,或者唯一的出路是从每个集群写入单独的文件然后合并结果?
我有一个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)
我的问题是 - 为什么这样做(因为我为什么加快速度?).听起来像是一个愚蠢的问题,但我在想 - 当然我的硬盘一次只能做一件事吗?那么为什么没有一个过程被搁置直到另一个过程完成?
当用高级语言写作时,这样的事情对用户是隐藏的.我想知道什么是低级别的?
我有一个任务来计算数组中的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) 在代码的并行部分中,我将每个线程的结果保存到ConcurrentBag中。但是,完成此操作后,我需要遍历每个结果并通过我的评估算法运行它们。普通的foreach实际上会遍历所有成员,还是我需要特殊的代码?我还考虑过使用队列之类的东西来代替行李,但我不知道哪种方法最好。在并行代码的末尾,袋子通常只包含20个左右的项目。
即,实际上将为ConcurrentBag的所有成员访问并运行foreach吗?
ConcurrentBag futures = new ConcurrentBag();
foreach(move in futures)
{
// stuff
}
Run Code Online (Sandbox Code Playgroud) 在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的多处理模块分析一些代码('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
我在并行计算集群的不同处理器上运行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)