我有一个函数需要从一个数组中查看大约20K行,并为每个行应用一个外部脚本.这是一个缓慢的过程,因为PHP在继续下一行之前等待脚本执行.
为了使这个过程更快,我正在考虑同时在不同的部分运行该功能.因此,例如,行0到2000作为一个函数,2001到4000在另一个函数上,依此类推.我怎样才能以一种干净的方式做到这一点?我可以创建不同的cron作业,每个函数一个用不同的参数:myFunction(0, 2000),然后另一个cron作业myFunction(2001, 4000),等等,但这似乎不太干净.这样做的好方法是什么?
到目前为止,每当我需要使用时,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) 我有一个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) 我想并行执行我的rspec测试
如果我需要在6核上执行所有测试,我会使用 bundle exec rake parallel:spec[6]
如何为spec rake指定标签选项,例如 --tag ~chat
我无法解释这一点,但我在其他人的代码中发现了这种现象:
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 setDaemon()方法的文档,当我读到JVM退出而不等待守护程序线程完成时,我感到很困惑.
但是,由于本质上守护程序线程是Java Thread的,可能依赖于在JVM上运行来实现其功能,如果JVM在守护程序线程完成之前退出,守护程序线程如何能够存活?
我正在研究一个Torch/Lua项目,我在其中实现了一个人工神经网络模型.一切正常,但现在我想以下列方式修改我的代码.由于我的输入数据集非常大,我想将它除以N = 20跨度.
然后我想仅在第一个数据集范围内训练我的神经网络,然后并行测试其他N-1 = 19个跨度.
要运行所有这些并行作业,我需要将神经网络模型详细信息保存到文件中,然后为每19个作业加载它.
火炬有没有办法正确"写"人工神经网络模型到文件?
类包含应该只创建一次的属性.创建过程是通过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) 我花了几个小时来尝试并行化我的数字运算代码,但是当我这样做时它只会变慢.不幸的是,当我尝试将其减少到下面的示例时,问题就消失了,我真的不想在这里发布整个程序.所以问题是:在这类程序中我应该避免哪些陷阱?
(注意:Unutbu的答案在底部后跟进.)
以下是情况:
BigData包含大量内部数据的类.在该示例中,存在一个ff插值函数列表; 在实际的程序,还有更多,例如ffA[k],ffB[k],ffC[k].do_chunk().do_single()将在5秒内do_multi()运行并且将在55秒内运行.xi和yi数组切割成连续的块并迭代k每个块中的所有值来分解工作.这工作得更好一点.现在,无论是使用1,2,3或4个线程,总执行时间都没有差别.但当然,我希望看到实际的加速!def do_chunk(array1, array2, array3)并对该数组进行仅限于numpy的计算.在那里,有显着的速度提升.#!/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) 我试图将数据持久化到数据库中.我的持久化方法是异步的.
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方法中处理异步调用的最佳方法是什么?