标签: parallel-processing

并行执行功能

我有一个函数需要从一个数组中查看大约20K行,并为每个行应用一个外部脚本.这是一个缓慢的过程,因为PHP在继续下一行之前等待脚本执行.

为了使这个过程更快,我正在考虑同时在不同的部分运行该功能.因此,例如,行0到2000作为一个函数,2001到4000在另一个函数上,依此类推.我怎样才能以一种干净的方式做到这一点?我可以创建不同的cron作业,每个函数一个用不同的参数:myFunction(0, 2000),然后另一个cron作业myFunction(2001, 4000),等等,但这似乎不太干净.这样做的好方法是什么?

php parallel-processing

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

如何获得Python多处理池剩余的"工作量"?

到目前为止,每当我需要使用时,multiprocessing我都是通过手动创建"进程池"并与所有子进程共享工作队列来完成的.

例如:

from multiprocessing import Process, Queue


class MyClass:

    def __init__(self, num_processes):
        self._log         = logging.getLogger()
        self.process_list = []
        self.work_queue   = Queue()
        for i in range(num_processes):
            p_name = 'CPU_%02d' % (i+1)
            self._log.info('Initializing process %s', p_name)
            p = Process(target = do_stuff,
                        args   = (self.work_queue, 'arg1'),
                        name   = p_name)
Run Code Online (Sandbox Code Playgroud)

这样我就可以在队列中添加东西,这些东西将由子进程使用.然后,我可以通过检查以下内容来监控处理的进度Queue.qsize():

    while True:
        qsize = self.work_queue.qsize()
        if qsize == 0:
            self._log.info('Processing finished')
            break
        else:
            self._log.info('%d simulations still need to be calculated', qsize)
Run Code Online (Sandbox Code Playgroud)

现在我认为这multiprocessing.Pool可以简化很多代码.

我无法找到的是如何监控仍有待完成的"工作量".

请看以下示例:

from multiprocessing import …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing pool process multiprocessing

9
推荐指数
1
解决办法
8082
查看次数

如何并行scipy稀疏矩阵乘法

我有一个scipy.sparse.csr_matrix格式的大型稀疏矩阵X,我想用一个利用并行性的numpy数组W来乘以它.经过一些研究后,我发现我需要在多处理中使用Array,以避免在进程之间复制X和W(例如:如何在Python多处理中将Pool.map与Array(共享内存)结合起来?并将共享的只读数据复制到Python多处理的不同过程?).这是我最近的尝试

import multiprocessing 
import numpy 
import scipy.sparse 
import time 

def initProcess(data, indices, indptr, shape, Warr, Wshp):
    global XData 
    global XIndices 
    global XIntptr 
    global Xshape 

    XData = data 
    XIndices = indices 
    XIntptr = indptr 
    Xshape = shape 

    global WArray
    global WShape 

    WArray = Warr     
    WShape = Wshp 

def dot2(args):
    rowInds, i = args     

    global XData 
    global XIndices
    global XIntptr 
    global Xshape 

    data = numpy.frombuffer(XData, dtype=numpy.float)
    indices = numpy.frombuffer(XIndices, dtype=numpy.int32)
    indptr = numpy.frombuffer(XIntptr, dtype=numpy.int32)
    Xr = scipy.sparse.csr_matrix((data, indices, …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing scipy sparse-matrix

9
推荐指数
1
解决办法
4253
查看次数

如何使用parallel_tests运行标记的rspec测试?

我想并行执行我的rspec测试

如果我需要在6核上执行所有测试,我会使用 bundle exec rake parallel:spec[6]

如何为spec rake指定标签选项,例如 --tag ~chat

parallel-processing rspec ruby-on-rails-3

9
推荐指数
1
解决办法
1743
查看次数

SecurityException来自并行流中的I/O代码

我无法解释这一点,但我在其他人的代码中发现了这种现象:

import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.util.stream.Stream;

import org.junit.Test;

public class TestDidWeBreakJavaAgain
{
    @Test
    public void testIoInSerialStream()
    {
        doTest(false);
    }

    @Test
    public void testIoInParallelStream()
    {
        doTest(true);
    }

    private void doTest(boolean parallel)
    {
        Stream<String> stream = Stream.of("1", "2", "3");
        if (parallel)
        {
            stream = stream.parallel();
        }
        stream.forEach(name -> {
            try
            {
                Files.createTempFile(name, ".dat");
            }
            catch (IOException e)
            {
                throw new UncheckedIOException("Failed to create temp file", e);
            }
        });
    }
}
Run Code Online (Sandbox Code Playgroud)

在启用安全管理器的情况下运行时,仅仅调用parallel()流,或者parallelStream()从集合中获取流时,似乎可以保证所有执行I/O的尝试都会抛出SecurityException.(最有可能的,它调用任何方法可以抛出 …

java parallel-processing securitymanager java-8 java-stream

9
推荐指数
1
解决办法
992
查看次数

JVM退出后,守护程序线程如何生存?

我正在阅读有关Java setDaemon()方法的文档,当我读到JVM退出而不等待守护程序线程完成时,我感到很困惑.

但是,由于本质上守护程序线程是Java Thread的,可能依赖于在JVM上运行来实现其功能,如果JVM在守护程序线程完成之前退出,守护程序线程如何能够存活?

java parallel-processing multithreading jvm daemon

9
推荐指数
1
解决办法
5765
查看次数

Torch/Lua,如何将经过训练的神经网络模型保存到文件中?

我正在研究一个Torch/Lua项目,我在其中实现了一个人工神经网络模型.一切正常,但现在我想以下列方式修改我的代码.由于我的输入数据集非常大,我想将它除以N = 20跨度.

然后我想仅在第一个数据集范围内训练我的神经网络,然后并行测试其他N-1 = 19个跨度.

要运行所有这些并行作业,我需要将神经网络模型详细信息保存到文件中,然后为每19个作业加载它.

火炬有没有办法正确"写"人工神经网络模型到文件?

parallel-processing lua torch

9
推荐指数
1
解决办法
5105
查看次数

如何在MultiThread场景中对仅应执行一次的代码进行单元测试?

类包含应该只创建一次的属性.创建过程是通过Func<T>参数传递的.这是缓存方案的一部分.

测试时要注意,无论有多少线程尝试访问该元素,创建只会发生一次.

单元测试的机制是在访问器周围启动大量线程,并计算调用创建函数的次数.

这根本不是确定性的,没有任何保证可以有效地测试多线程访问.也许一次只有一个线程可以锁定.(实际上,getFunctionExecuteCount如果lock不存在,则在7到9之间......在我的机器上,没有任何保证在CI服务器上它将是相同的)

如何以确定的方式重写单元测试?如何确保lock多个线程多次触发?

using Microsoft.VisualStudio.TestTools.UnitTesting;
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

namespace Example.Test
{
    public class MyObject<T> where T : class
    {
        private readonly object _lock = new object();
        private T _value = null;
        public T Get(Func<T> creator)
        {
            if (_value == null)
            {
                lock (_lock)
                {
                    if (_value == null)
                    {
                        _value = creator();
                    }
                }
            }
            return _value;
        }
    }

    [TestClass]
    public class UnitTest1
    {
        [TestMethod] …
Run Code Online (Sandbox Code Playgroud)

.net c# parallel-processing multithreading unit-testing

9
推荐指数
1
解决办法
582
查看次数

使用numpy/scipy最大限度地减少Python multiprocessing.Pool的开销

我花了几个小时来尝试并行化我的数字运算代码,但是当我这样做时它只会变慢.不幸的是,当我尝试将其减少到下面的示例时,问题就消失了,我真的不想在这里发布整个程序.所以问题是:在这类程序中我应该避免哪些陷阱?

(注意:Unutbu的答案在底部后跟进.)

以下是情况:

  • 它是关于一个模块,它定义了一个BigData包含大量内部数据的类.在该示例中,存在一个ff插值函数列表; 在实际的程序,还有更多,例如ffA[k],ffB[k],ffC[k].
  • 计算将被归类为"令人尴尬的并行":可以一次在较小的数据块上完成工作.在这个例子中,那是do_chunk().
  • 在我的实际程序中,示例中显示的方法将导致最差的性能:每个块大约1秒(在单个线程中完成的实际计算时间的0.1秒左右).因此,对于n = 50,do_single()将在5秒内do_multi()运行并且将在55秒内运行.
  • 我还尝试通过将xiyi数组切割成连续的块并迭代k每个块中的所有值来分解工作.这工作得更好一点.现在,无论是使用1,2,3或4个线程,总执行时间都没有差别.但当然,我希望看到实际的加速!
  • 这可能是相关的:Multiprocessing.Pool使Numpy矩阵乘法更慢.但是,在程序的其他地方,我使用了一个多处理池进行更加孤立的计算:一个看起来类似的函数(没有绑定到类),def do_chunk(array1, array2, array3)并对该数组进行仅限于numpy的计算.在那里,有显着的速度提升.
  • CPU使用率随着预期的并行进程数量而变化(三个线程的CPU使用率为300%).
#!/usr/bin/python2.7

import numpy as np, time, sys
from multiprocessing import Pool
from scipy.interpolate import RectBivariateSpline

_tm=0
def stopwatch(msg=''):
    tm = time.time()
    global _tm
    if _tm==0: _tm = tm; return
    print("%s: %.2f seconds" % (msg, tm-_tm))
    _tm = tm

class …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing numpy pool multiprocessing

9
推荐指数
1
解决办法
4735
查看次数

在Akka actor的receive方法中处理异步调用的最佳方法

我试图将数据持久化到数据库中.我的持久化方法是异步的.

class MyActor(persistenceFactory:PersistenceFactory) extends Actor {
  def receive: Receive = {
    case record: Record =>
      // this method is asynchronous, immediate return Future[Int]
      persistenceFactory.persist(record) 
  }
}
Run Code Online (Sandbox Code Playgroud)

当应用程序在增加的负载下运行时,瓶颈就是我们内存不足或没有线程可用.

那么在Akka actor的receive方法中处理异步调用的最佳方法是什么?

parallel-processing multithreading asynchronous scala akka

9
推荐指数
1
解决办法
555
查看次数