标签: parallel-processing

GNU并行组合,多次使用参数列表

我想使用以下来生成唯一的作业,其中{1}和{2}是唯一的元组:

parallel echo {1} {2} ::: A B C D ::: A B C D
Run Code Online (Sandbox Code Playgroud)

例如在python(itertools)中提供了这样一个组合生成器:

permutations('ABCD', 2)
Run Code Online (Sandbox Code Playgroud)

AB AC AD BA BC BD CA CB CD DA DB DC


有没有办法直接通过bash实现它?还是GNU并行本身?也许以某种方式跳过冗余工作?但是,如何检查已使用的参数组合.

parallel echo {= 'if($_==3) { skip() }' =} ::: {1..5}
Run Code Online (Sandbox Code Playgroud)

linux parallel-processing bash combinations gnu

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

Nifi-使用ExecuteStreamCommand进行并行和并发执行

目前,我在具有4个核心的边缘节点上运行Nifi.假设我有20个传入的流文件,并且我为ExecuteStreamCommand处理器提供10个并发任务,这是否意味着我只获得并发执行或并发和并行执行?

parallel-processing apache-nifi

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

为什么以下简单的并行化代码比Python中的简单循环慢得多?

一个简单的程序,用于计算数字平方并存储结果:

    import time
    from joblib import Parallel, delayed
    import multiprocessing

    array1 = [ 0 for i in range(100000) ]

    def myfun(i):
        return i**2

    #### Simple loop ####
    start_time = time.time()

    for i in range(100000):
        array1[i]=i**2

    print( "Time for simple loop         --- %s seconds ---" % (  time.time()
                                                               - start_time
                                                                 )
            )
    #### Parallelized loop ####
    start_time = time.time()
    results = Parallel( n_jobs  = -1,
                        verbose =  0,
                        backend = "threading"
                        )(
                        map( delayed( myfun ),
                             range( 100000 )
                             )
                        ) …
Run Code Online (Sandbox Code Playgroud)

python arrays parallel-processing function

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

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
查看次数

有没有办法在Chapel中自定义整个数组语句的默认并行化行为?

根据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)

arrays parallel-processing chapel

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

存在哪些用于并行化解析器的概念或算法?

对于已经以拆分格式给出的大量输入数据(例如大量的单个数据库条目),解析器的并行化似乎很容易,或者很容易通过快速的预处理步骤进行拆分(例如,将句子的语法结构解析为大型)文本。

似乎很难进行并行解析,这已经需要付出很多努力才能在给定输入中定位子结构。通用编程语言代码看起来像一个很好的例子。在Haskell之类的使用布局/缩进分隔单个定义的语言中,您可能会在找到新定义的开始后检查每行的前导空格数,跳过所有行,直到找到另一个定义并传递每个定义跳过块到另一个线程进行完全解析。

对于使用平衡括号定义范围的C,JavaScript等语言,进行预处理的工作量会更高。您需要遍历整个输入,从而计算大括号,注意字符串文字中的文本,等等。对于XML之类的语言而言,情况更糟,您还需要在打开/关闭标签中跟踪标签名称。

我发现CYK解析算法的并行版本似乎适用于所有无上下文语法。但是我很好奇还有什么其他通用概念/算法可以使解析器并行化,包括上述大括号计数这样的事情,它们仅适用于有限的一组语言。这个问题不是关于特定的实现,而是关于这些实现的思想。

algorithm parallel-processing parsing

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

简单地调用parallelStream是否并行运行任务?

我在互联网上看到了一些例子,为了使用流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 parallel-processing concurrency java-stream

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