标签: parallel-processing

Joblib 和内存限制并行化

在 Python 中,joblib提供了一个非常好的工具来执行令人尴尬的并行执行。我对此有点陌生,我正在尝试找出如何处理可能受内存限制的工作。

例如,请考虑以下情况:

from joblib import Parallel, delayed
import numpy as np

num_workers=8
num_jobs=100
memory_range=1000

def my_function(x):
    np.random.seed(x)
    array_length=(np.random.randint(memory_range)+1)**3
    big_list_of_numbers=np.random.uniform(array_length)
    return(np.sum(big_list_of_numbers))

results=Parallel(n_jobs=num_workers)(delayed(my_function)(i) for i in range(num_jobs))
Run Code Online (Sandbox Code Playgroud)

每个作业的内存使用量可能从几百字节到大约 0.8 GB 不等。

在完美的世界中,我想给 joblib 一些总内存预算(例如“我希望 99% 确定所有工作进程的总内存使用量小于 8GB”),让 joblib 估计各个工作进程的内存使用量分布到目前为止它已经看到的作业,并使用它来动态调整工作池的大小。

不幸的是,似乎(?)当您调用 Parallel() 时,活动工作线程的数量是静态确定的。如果是这种情况,我可以使用 resource.getrusage() 从一些示例作业中计算出内存使用情况,自己估计分布,并使用固定上限。如果内存使用情况变化很大,那么这要么是危险的,要么是低效的。

如果我要并行化一个特定的函数,我只需对其进行分析即可完成。但是,对于我的应用程序,我不会提前知道该功能。

我不想重新发明这个轮子。所以,我的问题是:

  1. joblib 是否支持此功能的任何部分?

  2. 还有其他(希望是简单的)工具吗?

python memory parallel-processing joblib

7
推荐指数
0
解决办法
575
查看次数

Julia中的并行梯度计算

不久前我被说服放弃了我舒适的matlab编程并开始在Julia编程.我一直在用神经网络工作很长时间,我认为,现在有了Julia,我可以通过并行计算梯度来更快地完成任务.

无需一次性对整个数据集计算梯度; 相反,人们可以拆分计算.例如,通过将数据集分成几部分,我们可以计算每个部分的部分梯度.然后通过将部分梯度相加来计算总梯度.

虽然原理很简单,但当我与Julia并行时,我会遇到性能下降,即一个进程比两个进程更快!我显然做错了什么......我已经咨询过论坛中提出的其他问题,但我仍然无法拼凑出答案.我认为我的问题在于有很多不必要的数据正在发生,但我无法正确修复它.

为了避免发布凌乱的神经网络代码,我发布了一个更简单的例子,它在线性回归的设置中复制了我的问题.

下面的代码块为线性回归问题创建了一些数据.代码解释了常量,但X是包含数据输入的矩阵.我们随机创建一个权重向量w ^当与乘以X创造一些目标ÿ.

######################################
## CREATE LINEAR REGRESSION PROBLEM ##
######################################

# This code implements a simple linear regression problem

MAXITER = 100   # number of iterations for simple gradient descent
N = 10000       # number of data items
D = 50          # dimension of data items
X = randn(N, D) # create random matrix of data, data items appear row-wise
Wtrue = randn(D,1) # create arbitrary weight …
Run Code Online (Sandbox Code Playgroud)

parallel-processing gradient linear-regression julia

7
推荐指数
1
解决办法
505
查看次数

块矩阵矩阵乘法的最佳块大小值

我想用以下 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

7
推荐指数
1
解决办法
3301
查看次数

如何在 Golang 中实现适当的并行性?goroutines 在 1.5+ 版本中是并行的吗?

我知道并行性和并发性之间的区别。我正在寻找如何在 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 的回答是完全有效的。出于某种原因,在 …

parallel-processing concurrency go

7
推荐指数
1
解决办法
4951
查看次数

基本的并行 python 程序在 Windows 上冻结

这是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。可能是什么问题?

python windows parallel-processing

7
推荐指数
1
解决办法
549
查看次数

多处理比 Windows 中的串行处理慢(但不是在 Linux 中)

我正在尝试并行化 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

7
推荐指数
1
解决办法
1236
查看次数

当轴 = 0 时,熊猫并行应用

我想在所有 Pandas 列上并行应用一些函数。例如,我想并行执行此操作:

def my_sum(x, a):
    return x + a


df = pd.DataFrame({'num_legs': [2, 4, 8, 0],
                   'num_wings': [2, 0, 0, 0]})
df.apply(lambda x: my_sum(x, 2), axis=0)
Run Code Online (Sandbox Code Playgroud)

我知道有一个swifter包,但它不支持axis=0应用:

NotImplementedError:Swifter 无法在大型数据集上执行 axis=0 应用。Dask 目前没有实现 axis=0 应用。更多详情请访问https://github.com/jmcarpenter2/swifter/issues/10

Dask 也不支持此功能axis=0(根据swifter 中的文档)。

我用谷歌搜索了几个来源,但找不到简单的解决方案。

不敢相信这在熊猫中如此复杂。

python parallel-processing pandas python-multiprocessing

7
推荐指数
1
解决办法
1652
查看次数

c++ 如何优雅地将 c++17 并行执行与计算整数的 for 循环一起使用?

我可以

std::vector<int> a;
a.reserve(1000);
for(int i=0; i<1000; i++)
    a.push_back(i);
std::for_each(std::execution::par_unseq, std::begin(a), std::end(a), [&](int i) {
  ... do something based on i ...
});
Run Code Online (Sandbox Code Playgroud)

但是有没有更优雅的方法来创建 for(int i=0; i<n; i++) 的并行版本,它不需要我先用升序整数填充向量?

c++ parallel-processing c++17

7
推荐指数
2
解决办法
2318
查看次数

从 8 个并行流升级的 Java 11 抛出 ClassNotFoundException

我将 Spring Boot 应用程序从 Java 8 和 tomcat 8 升级到了 java 11 和 tomcat 9。除了我在列表上使用并行流的部分外,一切似乎都运行良好。

list.addAll(items
                .parallelStream()
                .filter(item -> !SomeFilter.isOk(item.getId()))
                .map(logic::getSubItem)
                .collect(Collectors.toList()));
Run Code Online (Sandbox Code Playgroud)

前一部分代码过去在 Java 8 和 Tomcat 8 上工作得很好,但在Java 9 改变了如何使用 Fork/Join 公共池线程加载类之后,将系统类加载器作为它们的线程上下文类加载器返回。

我知道在后台并行流使用 ForkJoinPool 并且我创建了一个自定义 bean 类,但仍然没有被应用程序使用。很可能是因为它们可能是在这个 bean 之前创建的。

@Bean
public ForkJoinPool myForkJoinPool() {
    return new ForkJoinPool(threadPoolSize, makeFactory("APP"), null, false);
}

private ForkJoinPool.ForkJoinWorkerThreadFactory makeFactory(String prefix) {
    return pool -> {
        final ForkJoinWorkerThread worker = ForkJoinPool.defaultForkJoinWorkerThreadFactory.newThread(pool);
        worker.setName(prefix + worker.getPoolIndex());
        worker.setContextClassLoader(Application.class.getClassLoader());
        return worker;
    };
}
Run Code Online (Sandbox Code Playgroud)

最后,我还尝试将它包装在我的 ForkJoinPool 实例周围,但它是异步完成的,我不想这样做。我也不想使用提交和获取,因为这意味着我必须用 try/catch …

java parallel-processing spring java-stream java-11

7
推荐指数
1
解决办法
352
查看次数

优化用 TypeScript 编写的文件内容解析器类

I got a typescript module (used by a VSCode extension) which accepts a directory and parses the content contained within the files. For directories containing large number of files this parsing takes a bit of time therefore would like some advice on how to optimize it.

I don't want to copy/paste the entire class files therefore will be using a mock pseudocode containing the parts that I think are relevant.

class Parser {
    constructor(_dir: string) {
        this.dir = _dir;
    } …
Run Code Online (Sandbox Code Playgroud)

javascript parallel-processing node-worker-threads

7
推荐指数
1
解决办法
122
查看次数