如何在满足条件的情况下为所有工作人员返回的函数中编写并行for循环?
就是这样的:
function test(n)
@sync @parallel for i in 1:1000
{... statement ...}
if {condition}
return test(n+1)
end
end
end
Run Code Online (Sandbox Code Playgroud)
所有工人都停止在for循环上工作,只有主进程返回?(其他进程再次开始使用下一个for循环?)
如果输入大小太小,库会自动序列化流中地图的执行,但这种自动化不会,也不能考虑地图操作的重要程度.有没有办法强制parallelStream()实际并行化CPU 重图?
根据OCP的书,人们必须避免有状态的操作,否则称为有状态的lambda表达.本书中提供的定义是"有状态的lambda表达式,其结果取决于在执行管道期间可能发生变化的任何状态."
它们提供了一个示例,其中使用并行流将固定的数字集合添加到使用该.map()函数的同步ArrayList中.
arraylist中的顺序是完全随机的,这应该让人看到有状态的lambda表达式在运行时产生不可预测的结果.这就是为什么强烈建议在使用并行流时避免有状态操作以消除任何潜在的数据副作用.
它们没有显示无状态lambda表达式,它提供了解决同一问题的方法(向同步的arraylist添加数字),我仍然不知道使用map函数用数据填充空的同步arraylist的问题. ..在执行管道期间可能发生变化的状态究竟是什么?他们指的是Arraylist本身吗?就像当另一个线程决定在并行流仍处于添加数字并因此改变最终结果的过程中时将其他数据添加到ArrayList时?
也许有人可以为我提供一个更好的例子来说明有状态的lambda表达式是什么以及为什么要避免它.非常感谢.
谢谢
我在这里发现了一种"奇怪"的行为.我懂了:
{-# LANGUAGE BangPatterns #-}
import Data.List
import Control.Parallel
import Control.Parallel.Strategies
fib 0 = 1
fib 1 = 1
fib n = fib (n-1) + fib (n-2)
main =
let xs = [ fib (20 + n `mod` 2) | n <- [0..1000] ]
`using` test rseq
in print (sum xs)
test :: Strategy a -> Strategy [a]
test strat xs = do
parTraversable strat xs -- case #1
-- parList strat xs -- case #2
return xs
Run Code Online (Sandbox Code Playgroud)
case #1 …
我想用以下 C 代码进行块矩阵矩阵乘法。在这种方法中,大小为 BLOCK_SIZE 的块被加载到最快的缓存中,以减少计算过程中的内存流量。
void bMMikj(double **A , double **B , double ** C , int m, int n , int p , int BLOCK_SIZE){
int i, j , jj, k , kk ;
register double jjTempMin = 0.0 , kkTempMin = 0.0;
for (jj=0; jj<n; jj+= BLOCK_SIZE) {
jjTempMin = min(jj+ BLOCK_SIZE,n);
for (kk=0; kk<n; kk+= BLOCK_SIZE) {
kkTempMin = min(kk+ BLOCK_SIZE,n);
for (i=0; i<n; i++) {
for (k = kk ; k < kkTempMin ; k++) { …Run Code Online (Sandbox Code Playgroud) parallel-processing caching hpc matrix matrix-multiplication
我知道并行性和并发性之间的区别。我正在寻找如何在 Go 中实现并行性。我希望 goroutines 是并行的,但我发现的文档似乎相反。
设置 GOMAXPROCS 允许我们配置应用程序可用于并行运行的线程数。从 1.5 版开始,GOMAXPROCS 将核心数作为值。据我了解,从 1.5 版开始,goroutine 本质上是并行的。这个结论正确吗?
我在 StackOverflow 等网站上发现的每个问题似乎都已经过时,并且没有考虑到 1.5 版中的这一变化。请参阅:golang 中的并行处理
我的困惑源于试图在实践中测试这种并行性。我在 Go 1.10 中尝试了以下代码,但它没有并行运行。
package main
import (
"fmt"
"sync"
)
var wg sync.WaitGroup
func main() {
wg.Add(2)
go count()
go count()
wg.Wait()
}
func count() {
defer wg.Done()
for i := 0; i < 10; i++ {
fmt.Println(i)
}
}
Run Code Online (Sandbox Code Playgroud)
将 GOMAXPROCS 设置为 2 不会改变结果。我得到一个并发程序而不是并行程序。
我在 8 核系统上运行所有测试。
编辑:供将来参考,
我被这个博客迷住了:https : //www.ardanlabs.com/blog/2014/01/concurrency-goroutines-and-gomaxprocs.html,其中在一个小的 for 循环中实现了并行性而没有太多麻烦。@peterSO 的回答是完全有效的。出于某种原因,在 …
这是https://docs.python.org/2/library/multiprocessing.html#module-multiprocessing.pool关于并行处理的基本 Python 示例
from multiprocessing import Pool
def f(x):
return x*x
if __name__ == '__main__':
p = Pool(5)
print(p.map(f, [1, 2, 3]))
Run Code Online (Sandbox Code Playgroud)
由于某种原因,我无法在我的 PC 上运行。当我尝试执行第三个块时,程序冻结了。我的操作系统是 Windows 10。我在 Spyder IDE 上运行该程序,并且安装了 anaconda。可能是什么问题?
我正在尝试并行化 afor loop以加速我的代码,因为循环处理操作都是独立的。按照在线教程,multiprocessingPython 中的标准库似乎是一个好的开始,我已经将它用于基本示例。
但是,对于我的实际用例,我发现在 Windows 上运行时,并行处理(使用双核机器)实际上要慢一点(<5%)。然而,与串行执行相比,在 Linux 上运行相同的代码会使并行处理速度提高约 25%。
从文档中,我认为这可能与 Window 缺少 fork() 函数有关,这意味着该进程每次都需要重新初始化。但是,我不完全理解这一点,想知道是否有人可以确认这一点?
特别,
--> 这是否意味着调用 python 文件中的所有代码都会为 Windows 上的每个并行进程运行,甚至初始化类和导入包?
--> 如果是这样,是否可以通过将类的副本(例如使用 deepcopy)传递给新进程来避免这种情况?
--> 是否有任何提示/其他策略可以有效地并行化 unix 和 windows 的代码设计。
我的确切代码很长并且使用了很多文件,所以我创建了一个伪代码样式的示例结构,希望能显示这个问题。
# Imports
from my_package import MyClass
imports many other packages / functions
# Initialization (instantiate class and call slow functions that get it ready for processing)
my_class = Class()
my_class.set_up(input1=1, input2=2)
# Define main processing function to be used in loop
def calculation(_input_data):
# …Run Code Online (Sandbox Code Playgroud) python parallel-processing multithreading multiprocessing python-multiprocessing
我有一个带有 8 个逻辑处理器的 corei7 处理器。
我正在尝试使用 parallel.For 在 dotnet core 2.2 中运行并行任务。当我测量开始时间时,有 9 个任务并行启动。不是应该只有8个吗?
下面你可以看到:
i => [ThreadId],[ProcessorNumber] == starttime - endtime
并行任务结果
我正在尝试从大约 100 个文件中拆分和重新排序数据。这些文件包含有关大约的信息。40 万用户,按时间顺序排列。目的是遍历 100 个文件并为每个用户创建一个单独的文件。
该方法是使用 python 多处理库来创建多个进程: - 一个工作进程处理池,用于加载数据、对其进行排序并将批次附加到队列中。每个批次包含一个用户的数据。- 一个从队列中取出元素并将批次添加到队列中的进程
现在,下面的代码适用于装有 Windows 10、Python 3.7 的笔记本电脑。然而,我试图在带有 ubuntu 16.04 和 python 3.5 的服务器上运行它。虽然在 Windows 上代码运行没有问题,但服务器内存不足。现在笔记本电脑有 8GB 的 RAM,而服务器有 256GB。我在这里缺少什么?
编辑:我已将代码替换为您应该能够运行的简化版本。代码现在创建进程,将 numpy 数组添加到队列中。另一个进程使数组出列。在数组排队期间,打印队列的大小。由于元素出列,队列大小保持较低。但是,内存使用量不断增加。
import numpy as np
import sys
import multiprocessing
import time
def queue_remover(q):
while(1):
if not q.empty():
data = q.get()
def queue_adder(q, unused_number):
for _ in range(5):
q.put(np.zeros((1000,1000)))
print('q size ', q.qsize())
sys.stdout.flush()
time.sleep(1)
if __name__ == "__main__":
list_of_numbers = list(range(500))
m = multiprocessing.Manager()
queue = m.Queue(maxsize=40000)
writer = …Run Code Online (Sandbox Code Playgroud) python parallel-processing multiprocessing python-3.x python-multiprocessing