标签: parallel-processing

Parallel.ForEach使用自定义TaskScheduler来防止OutOfMemoryException

我正在通过Parallel.ForEach处理各种大小的PDF(简单的2MB到几百MB的高DPI扫描)并且偶尔会遇到OutOfMemoryException - 可以理解的是由于进程是32位并且Parallel产生了线程. ForEach占用了大量未知的内存消耗工作.

限制MaxDegreeOfParallelism确实有效,尽管由于所述线程的内存占用量较小而导致有大量(10k +)批量的小型PDF需要处理时的吞吐量不足.这是一个CPU繁重的过程,在遇到偶尔的大型PDF组并获得OutOfMemoryException之前,Parallel.ForEach很容易达到100%的CPU.运行Performance Profiler会将其备份.

根据我的理解,为Parallel.ForEach设置分区器不会提高我的性能.

这导致我使用TaskScheduler传递给我的Parallel.ForEach 的自定义MemoryFailPoint检查.在它周围搜索似乎有关于创建自定义TaskScheduler对象的稀缺信息.

之间寻找在.NET专业任务调度4并行扩展附加功能,在C#中的自定义的TaskScheduler这里#2各种各样的回答,我已经建立了我自己的TaskScheduler,并有我的QueueTask方法,例如:

protected override void QueueTask(Task task)
{
    lock (tasks) tasks.AddLast(task);
    try
    {
        using (MemoryFailPoint memFailPoint = new MemoryFailPoint(600))
        {
            if (runningOrQueuedCount < maxDegreeOfParallelism)
            {
                runningOrQueuedCount++;
                RunTasks();
            }
        }
    }
    catch (InsufficientMemoryException e)
    {     
        // somehow return thread to pool?           
        Console.WriteLine("InsufficientMemoryException");
    }
}
Run Code Online (Sandbox Code Playgroud)

虽然try/catch有点贵,但我的目标是捕获600MB的可能最大大小PDF(+一点额外内存开销)将抛出OutOfMemoryException.当我捕获InsufficientMemoryException时,这个解决方案似乎杀掉了试图完成工作的线程.有了足够大的PDF,我的代码最终成为一个单一的线程Parallel.ForEach.

在Parallel.ForEach和OutOfMemoryExceptions上的Stackoverflow上发现的其他问题似乎不适合我在线程上使用动态内存的最大吞吐量的用例,并且通常只是MaxDegreeOfParallelism作为静态解决方案使用,例如:

因此,为了获得可变工作内存大小的最大吞吐量,可以:

  • 如果线程被拒绝通过MemoryFailPoint支票工作,我如何将线程返回到线程池中?
  • 当有空闲内存时,我如何/在哪里安全地生成新线程以重新开始工作? …

c# parallel-processing multithreading

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

调试并行Python程序(mpi4py)

我有一个mpi4py程序间歇性挂起。如何跟踪各个流程在做什么?

我可以在不同的终端上运行该程序,例如使用 pdb

mpiexec -n 4 xterm -e "python -m pdb my_program.py"
Run Code Online (Sandbox Code Playgroud)

但是,如果问题仅通过大量进程(在我的情况下为〜80)表现出来,则将变得很麻烦。另外,很容易捕获异常,pdb但是我需要查看跟踪以找出发生挂起的位置。

python debugging parallel-processing trace mpi4py

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

Future将并行运行部分代码

我有疑问future(),doFuture()用法.

我想N并行运行计算(使用foreach ... %dopar%) - N我的机器上有多少核心.为此,我使用future:

library(doFuture)
registerDoFuture()
plan(multiprocess)

foreach(i = seq_len(N)) %dopar% {
    foo <- rnorm(1e6)
}
Run Code Online (Sandbox Code Playgroud)

这就像一个魅力,因为我N并行运行计算.但是我需要实现另一个使用大量内核的分析步骤(例如,N).这是代码的样子:

foreach(i = seq_len(N)) %dopar% {
    foo <- rnorm(1e6)
    write.table(foo, paste0("file_", i, ".txt"))
    # This step uses high number of cores 
    system(paste0("head ", "file_", i, ".txt", " > ", "file_head_", i, ".txt")
}
Run Code Online (Sandbox Code Playgroud)

我在运行多个rnormhead并行,但由于head使用了大量的内核(让我们假设这个)我的分析卡住.

题:

如何使用并行运行部分代码future?(如何仅rnorm并行运行然后head …

parallel-processing r

5
推荐指数
0
解决办法
186
查看次数

ParallelQuery.Aggregate不能并行运行的可能原因

我很感激PLYNQ专家的帮助!我会花时间回顾一下答案,我对math.SE有一个更为确定的概况.

我有一个类型的对象ParallelQuery<List<string>>,它有44个列表,我想并行处理(一次五个,比如说).我的流程有一个签名

private ProcessResult Process(List<string> input)
Run Code Online (Sandbox Code Playgroud)

处理将返回一个结果,这是一对布尔值,如下所示.

    private struct ProcessResult
    {
        public ProcessResult(bool initialised, bool successful)
        {
            ProcessInitialised = initialised;
            ProcessSuccessful = successful;
        }

        public bool ProcessInitialised { get; }
        public bool ProcessSuccessful { get; }
    }
Run Code Online (Sandbox Code Playgroud)

问题.给定一个IEnumerable<List<string>> processMe,我的PLYNQ查询尝试实现此方法:https://msdn.microsoft.com/en-us/library/dd384151(v = vs.110).aspx .它写成

processMe.AsParallel()
         .Aggregate<List<string>, ConcurrentStack<ProcessResult>, ProcessResult>
             (
                 new ConcurrentStack<ProcessResult>,   //aggregator seed
                 (agg, input) =>
                 {                         //updating the aggregate result
                     var res = Process(input);
                     agg.Push(res);
                     return agg;
                 },
                 agg => 
                 {                         //obtain the result …
Run Code Online (Sandbox Code Playgroud)

.net c# parallel-processing multithreading plinq

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

为什么并行任务在第一时间总是很慢?

我有一些分类器,我想对一个样本进行评估。由于它们彼此独立,因此可以并行运行此任务。这意味着我要并行化它。

我用python和bash脚本尝试过。问题是,当我第一次运行该程序时,大约需要30到40秒才能完成。当我连续多次运行该程序时,只需1s-3s即可完成。即使我用不同的输入来输入分类器,我也得到了不同的结果,因此似乎没有缓存。当我运行其他程序并随后重新运行该程序时,它又需要40秒钟才能完成。

我在htop中还观察到,第一次运行该程序时CPU利用率不高,但是当我一次又一次重新运行时,CPU利用率就很高。

有人可以向我解释这种奇怪的行为吗?我如何避免这种情况,这样即使程序第一次运行也会很快?

这是python代码:

import time
import os
from fastText import load_model
from joblib import delayed, Parallel, cpu_count
import json

os.system("taskset -p 0xff %d" % os.getpid())

def format_duration(start_time, end_time):
    m, s = divmod(end_time - start_time, 60)
    h, m = divmod(m, 60)
    return "%d:%02d:%02d" % (h, m, s)

def classify(x, classifier_name, path):
    f = load_model(path + os.path.sep + classifier_name)    
    labels, probabilities = f.predict(x, 2)
    if labels[0] == '__label__True':
        return classifier_name
    else:
        return None

if __name__ == '__main__':
    with open('classifier_names.json') as …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing bash joblib

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

在后台执行命令

我使用下面的代码,它是"npm install"的执行命令,现在在调试时我看到命令执行大约需要10..15秒(取决于我有多少模块).我想要的是这个命令将在后台执行,程序将继续.

cmd := exec.Command(name ,args...)
cmd.Dir = entryPath
Run Code Online (Sandbox Code Playgroud)

在调试中我看到要移动到下一行tass大约10..15秒...

我有两个问题:

  1. 我怎么能这样做?因为我想要并行做某事......
  2. 我怎么知道它什么时候结束?提供与此命令相关的其他逻辑,即在npm install完成之后,我需要做其他事情.

parallel-processing concurrency background process go

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

Python多处理模块中apply()和apply_async()之间的区别

我目前有一段代码产生多个进程,如下所示:

pool = Pool(processes=None)
results = [pool.apply(f, args=(arg1, arg2, arg3)) for arg3 in arg_list]
Run Code Online (Sandbox Code Playgroud)

我的想法是,这将使用之后可用的所有核心将工作划分为核心processes=None.但是,多处理模块docs中Pool.apply()方法的文档如下:

相当于apply()内置函数.它会一直阻塞,直到结果准备就绪,因此apply_async()更适合并行执行工作.此外,func仅在池中的一个工作程序中执行.

第一个问题: 我不清楚这一点.如何apply在工人之间分配工作,以及与什么方式不同apply_async?如果任务分布在工人之间,那怎么可能func只在其中一个工人中执行?

我的猜测:我的猜测是apply,在我当前的实现中,给一个具有一组参数的工作者提供一个任务,然后等待该工作完成,然后将下一组参数提供给另一个工作者.通过这种方式,我将工作发送到不同的流程,但没有发生并行性.这似乎是这种情况,因为apply事实上只是:

def apply(self, func, args=(), kwds={}):
    '''
    Equivalent of `func(*args, **kwds)`.
    Pool must be running.
    '''
    return self.apply_async(func, args, kwds).get()
Run Code Online (Sandbox Code Playgroud)

第二个问题:我也想更好地理解为什么在出台的文件,部分16.6.1.5.('使用工人池'),他们说,即使是apply_async像这样 的建筑[pool.apply_async(os.getpid, ()) for i in range(4)] 可能会使用更多的工艺,但它不确定它会.是什么决定是否使用多个流程?

python parallel-processing multiprocessing

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

不能腌制本地物体

我想在此代码中使用多进程包.我试图调用该函数create_new_population并将数据分发到8个处理器,但是当我这样做时,我得到了pickle错误.

通常函数会像这样运行: self.create_new_population(self.pop_size)

我尝试像这样分发工作:

f= self.create_new_population
pop = self.pop_size/8
self.current_generation = [pool.apply_async(f, pop) for _ in range(8)]
Run Code Online (Sandbox Code Playgroud)

我得到 或Can't pickle local object 'exhaust.__init__.<locals>.tour_select'
PermissionError: [WinError 5] Access is denied

我仔细阅读了这个帖子,并尝试使用Steven Bethard的方法绕过错误,允许通过copyreg进行方法酸洗/ 取消:

def _pickle_method(method)
def _unpickle_method(func_name, obj, cls)
Run Code Online (Sandbox Code Playgroud)

我还尝试使用pathos包没有任何运气.
我知道应该在
if __name__ == '__main__':块下调用代码,但我想知道是否可以在代码中尽可能少的更改来完成.

python parallel-processing multithreading python-multiprocessing pathos

5
推荐指数
0
解决办法
2353
查看次数

并行执行"git submodule foreach"

有没有办法git submodule foreach并行执行命令,类似于--jobs 8参数的工作方式git submodule update

例如,我们工作的项目之一涉及近200个子组件(子模块),我们大量使用该foreach命令对它们进行操作.我想加快它们的速度.

PS:如果解决方案涉及脚本,我在Windows上工作,大多数时候,使用git-bash.

git parallel-processing git-submodules

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

Julia中的并行循环 - 不希望在开始之前分工

我的机器有4个核心.当我使用@sync @parallel进行并行运行时,我注意到Julia在将作业发送到4个处理器之前将作业分成4 个:

# start of do_something.jl
function do_something(i, parts)
    procs = zeros(Int, parts)
    procs[i] = myid()
    total = 0.0
    for j = 1:i * 100000000
        total = total + 1e-6
    end
    return procs
end
# end of do_something.jl

# synctest3a.jl
addprocs(Sys.CPU_CORES)
@everywhere include("do_something.jl")
parts = 20
procs = @sync @parallel (+) for i = 1:parts
    do_something(i, parts)
end
@printf("procs=%s\n", procs)
Run Code Online (Sandbox Code Playgroud)

julia synctest3a.jl的结果,表示前5个被发送到处理器2,接下来的5个被发送到处理器3,依此类推:

procs=[2, 2, 2, 2, 2, 3, 3, 3, 3, 3, 4, 4, 4, 4, 4, …
Run Code Online (Sandbox Code Playgroud)

parallel-processing julia

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