我是 R 中并行计算的新手,希望使用并行包来加速我的计算(这比下面的示例更复杂)。然而,与通常的 lapply 函数相比,使用 mclapply 函数时的计算时间要长得多。
\n\n我在笔记本电脑上安装了全新的 Ubuntu 18.04.2 LTS,它具有 7.7 GB 内存和 Intel\xc2\xae Core\xe2\x84\xa2 i7-4500U CPU @ 1.80GHz \xc3\x97 4 处理器。我在 R studio 上运行 R。
\n\nrequire(parallel)\n\na <- seq(0, 1, length.out = 110) #data\nb <- seq(0, 1, length.out = 110)\nc <- replicate(1000, sample(1:100,size=10), simplify=FALSE)\n\nfunction_A <- function(i, j, k) { # some random function to examplify the problem\n i+ j * pmax(i-k,0) \n}\n\n#running it with mclapply \nptm_mc <- proc.time() \noutput <- mclapply(1:NROW(c), function(o){ \n mclapply(1:NROW(a),function(p) function_A(a[p], b, c[[o]]))})\ntime_mclapply …Run Code Online (Sandbox Code Playgroud) 我一直在Parallel.ForEach对项目集合进行一些耗时的处理。该处理实际上是由外部命令行工具处理的,我无法更改它。然而,似乎Parallel.ForEach会“卡在”集合中长期运行的项目上。我已经将问题提炼出来,并且可以表明Parallel.ForEach,事实上,等待这个漫长的过程完成并且不允许任何其他人通过。我编写了一个控制台应用程序来演示该问题:
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
namespace testParallel
{
class Program
{
static int inloop = 0;
static int completed = 0;
static void Main(string[] args)
{
// initialize an array integers to hold the wait duration (in milliseconds)
var items = Enumerable.Repeat(10, 1000).ToArray();
// set one of the items to 10 seconds
items[50] = 10000;
// Initialize our line for reporting status
Console.Write(0.ToString("000") + " Threads, " + …Run Code Online (Sandbox Code Playgroud) c# parallel-processing blocking task-parallel-library parallel.foreach
模块joblib提供了一个简单的帮助程序类来使用多处理编写并行 for 循环。
此代码使用列表理解来完成这项工作:
import time
from math import sqrt
from joblib import Parallel, delayed
start_t = time.time()
list_comprehension = [sqrt(i ** 2) for i in range(1000000)]
print('list comprehension: {}s'.format(time.time() - start_t))
Run Code Online (Sandbox Code Playgroud)
大约需要0.51s
list comprehension: 0.5140271186828613s
Run Code Online (Sandbox Code Playgroud)
此代码使用joblib.Parallel()构造函数:
start_t = time.time()
list_from_parallel = Parallel(n_jobs=2)(delayed(sqrt)(i ** 2) for i in range(1000000))
print('Parallel: {}s'.format(time.time() - start_t))
Run Code Online (Sandbox Code Playgroud)
大约需要31秒
Parallel: 31.3990638256073s
Run Code Online (Sandbox Code Playgroud)
这是为什么?不应该Parallel()比非并行计算更快吗?
这是其中的一部分cpuinfo:
processor : 0
vendor_id : GenuineIntel
cpu family : 6
model : 79
model …Run Code Online (Sandbox Code Playgroud) 尝试找出使用python multiprocessing运行的正确并行进程数。
下面的脚本在 8 核、32 GB (Ubuntu 18.04) 计算机上运行。(在测试以下内容时,仅运行系统进程和基本用户进程。)
经过测试multiprocessing.Pool并apply_async具有以下内容:
from multiprocessing import current_process, Pool, cpu_count
from datetime import datetime
import time
num_processes = 1 # vary this
print(f"Starting at {datetime.now()}")
start = time.perf_counter()
print(f"# CPUs = {cpu_count()}") # 8
num_procs = 5 * cpu_count() # 40
def cpu_heavy_fn():
s = time.perf_counter()
print(f"{datetime.now()}: {current_process().name}")
x = 1
for i in range(1, int(1e7)):
x = x * i
x = x / i
t_taken = …Run Code Online (Sandbox Code Playgroud) python parallel-processing multiprocessing python-3.x python-3.6
到目前为止,我发现lapply在 R 中使用并行的最简单方法是通过以下示例代码:
library(parallel)
library(pbapply)
cl <- makeCluster(10)
clusterExport(cl = cl, {...})
clusterEvalQ(cl = cl, {...})
results <- pblapply(1:100, FUN = function(x){rnorm(x)}, cl = cl)
Run Code Online (Sandbox Code Playgroud)
它有一个非常有用的功能,即为结果提供进度条,并且当不需要并行计算时,通过设置 可以很容易地重用相同的代码cl = NULL。
然而,我注意到的一个问题是它pblapply正在批量循环遍历列表。例如,如果一个工人在某项任务上停留了很长时间,那么剩下的工人将等待该任务完成,然后再开始一批新的工作。对于某些任务,这会为工作流程增加大量不必要的时间。
我的问题:
是否有任何类似的并行框架允许工作人员独立运行?进度条和重用代码的能力cl=NULL将是一个很大的优势。
也许可以修改现有代码pbapply来添加此选项/功能?
1.我有一个函数var。我想知道通过利用系统拥有的所有处理器、内核、线程和 RAM 内存进行多处理/并行处理来快速运行此函数中的循环的最佳方法。
import numpy
from pysheds.grid import Grid
xs = 82.1206, 72.4542, 65.0431, 83.8056, 35.6744
ys = 25.2111, 17.9458, 13.8844, 10.0833, 24.8306
a = r'/home/test/image1.tif'
b = r'/home/test/image2.tif'
def var(interest):
variable_avg = []
for (x,y) in zip(xs,ys):
grid = Grid.from_raster(interest, data_name='map')
grid.catchment(data='map', x=x, y=y, out_name='catch')
variable = grid.view('catch', nodata=np.nan)
variable = numpy.array(variable)
variablemean = (variable).mean()
variable_avg.append(variablemean)
return(variable_avg)
Run Code Online (Sandbox Code Playgroud)
2.var如果我可以针对给定的函数多个参数并行运行函数和循环,那就太好了。var(a)例如:var(b)同时。因为它比单独并行化循环消耗的时间要少得多。
如果没有意义,请忽略 2。
python parallel-processing multithreading multiprocessing python-asyncio
我已经用简单的查询测试了并行性,但我不明白结果。
我检查了以下参数pg_settings:
max_parallel_workers = 8
max_parallel_workers_per_gather = 2
Run Code Online (Sandbox Code Playgroud)
我运行以下查询(该表包含约 16M 行):
explain analyze
select *
from tbl
where value<>-1
Run Code Online (Sandbox Code Playgroud)
结果:
Gather (cost=1000.00 .. 1136714.86 rows=580941 width=78 actual time=0.495..3057.813 rows = 587886 loops=1)
workers planned: 2
workers launched: 2
-> parallel seq scan on tbl (cost=0.00..10776.76 rows=242059 width=718) (actual time=0.095..2968.77 rows=195962 loops=3)
filter: (value<>-1::integer)
rows removed by filter: 5389091
plain time: 0.175ms
exection time: 3086.243ms
Run Code Online (Sandbox Code Playgroud)
max_parallel_workers和 和有什么不一样max_parallel_workers_per_gather?何时使用每个值?我需要如何运行脚本的安排是首先使用该函数并行运行 4 个 R 脚本rstudioapi::jobRunScript()。并行运行的每个脚本不会从任何环境导入任何内容,而是将创建的数据帧导出到全局环境。我的第 5 个 R 脚本基于并行运行的 4 个 R 脚本创建的数据帧,并且第 5 个脚本也在控制台中运行。如果有一种方法可以在前 4 个 R 脚本并行运行完成后在后台而不是在控制台中运行第 5 个脚本,那就会好很多。我还试图减少整个过程的总运行时间。
尽管我能够弄清楚如何并行运行前 4 个 R 脚本,但我的任务尚未完全完成,因为我找不到如何触发运行第 5 个 R 脚本的方法。希望大家能在这里帮助我
我想在 Apache Flink 中实现以下场景:
给定一个具有 4 个分区的 Kafka 主题,我想根据事件的类型使用不同的逻辑在 Flink 中独立处理分区内数据。
特别是,假设输入 Kafka 主题包含前面图像中描述的事件。每个事件都有不同的结构:分区 1 有字段“ a ”作为键,分区 2 有字段“ b ”作为键,等等。在 Flink 中,我想根据事件应用不同的业务逻辑,所以我想我应该以某种方式分割流。为了实现图中所描述的效果,我想只使用一个消费者来做类似的事情(我不明白为什么我应该使用更多):
FlinkKafkaConsumer<..> consumer = ...
DataStream<..> stream = flinkEnv.addSource(consumer);
stream.keyBy("a").map(new AEventMapper()).addSink(...);
stream.keyBy("b").map(new BEventMapper()).addSink(...);
stream.keyBy("c").map(new CEventMapper()).addSink(...);
stream.keyBy("d").map(new DEventMapper()).addSink(...);
Run Code Online (Sandbox Code Playgroud)
(一)正确吗?另外,如果我想并行处理每个 Flink 分区,因为我只想按顺序处理按同一 Kafka 分区排序的事件,而不是全局考虑它们,(b) 我该怎么办?我知道该方法的存在setParallelism(),但我不知道在这种情况下将其应用到哪里。
我正在寻找有关标记(a)和(b)的问题的答案。先感谢您。
parallel-processing partitioning apache-kafka apache-flink kafka-topic
python ×3
r ×3
c# ×2
apache-flink ×1
apache-kafka ×1
blocking ×1
furrr ×1
kafka-topic ×1
lapply ×1
nunit ×1
partitioning ×1
postgresql ×1
python-3.6 ×1
python-3.x ×1
rscript ×1
rstudioapi ×1
windows ×1