在 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() 从一些示例作业中计算出内存使用情况,自己估计分布,并使用固定上限。如果内存使用情况变化很大,那么这要么是危险的,要么是低效的。
如果我要并行化一个特定的函数,我只需对其进行分析即可完成。但是,对于我的应用程序,我不会提前知道该功能。
我不想重新发明这个轮子。所以,我的问题是:
joblib 是否支持此功能的任何部分?
还有其他(希望是简单的)工具吗?
不久前我被说服放弃了我舒适的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) 我想用以下 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
我知道并行性和并发性之间的区别。我正在寻找如何在 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 的回答是完全有效的。出于某种原因,在 …
这是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。可能是什么问题?
我正在尝试并行化 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
我想在所有 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 中的文档)。
我用谷歌搜索了几个来源,但找不到简单的解决方案。
不敢相信这在熊猫中如此复杂。
我可以
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++) 的并行版本,它不需要我先用升序整数填充向量?
我将 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 …
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)