标签: parallel-processing

Java Parallel Stream 与 ExecutorService 的性能

假设我们有一个列表并且想要选择满足某个属性的所有元素(比如一些函数 f)。有 3 种方法可以并行执行此过程。

一 :

listA.parallelStream.filter(element -> f(element)).collect(Collectors.toList());
Run Code Online (Sandbox Code Playgroud)

二:

listA.parallelStream.collect(Collectors.partitioningBy(element -> f(element))).get(true);
Run Code Online (Sandbox Code Playgroud)

三:

ExecutorService executorService = Executors.newFixedThreadPool(nThreads);
//separate the listA into several batches
for each batch {
     Future<List<T>> result = executorService.submit(() -> {
          // test the elements in this batch and return the validate element list
     });
}
//merge the results from different threads.
Run Code Online (Sandbox Code Playgroud)

假设测试功能是 CPU 密集型任务。我想知道哪种方法更有效。非常感谢。

java parallel-processing java-stream java-threads

5
推荐指数
1
解决办法
6877
查看次数

使用比核心更多的工作进程

这个来自 PYMOTW 的例子给出了一个例子,multiprocessing.Pool()其中processes传递的参数(工作进程数)是机器内核数的两倍。

pool_size = multiprocessing.cpu_count() * 2
Run Code Online (Sandbox Code Playgroud)

(否则该类将默认为 just cpu_count()。)

这有什么道理吗?创建比核心数更多的工人有什么影响?是否有理由这样做,或者它可能会在错误的方向上施加额外的开销?我很好奇为什么它会一直包含在我认为是信誉良好的网站的示例中。

在最初的测试中,它实际上似乎会减慢速度:

$ python -m timeit -n 25 -r 3 'import double_cpus; double_cpus.main()'
25 loops, best of 3: 266 msec per loop
$ python -m timeit -n 25 -r 3 'import default_cpus; default_cpus.main()'
25 loops, best of 3: 226 msec per loop
Run Code Online (Sandbox Code Playgroud)

double_cpus.py

import multiprocessing

def do_calculation(n):
    for i in range(n):
        i ** 2

def main():
    with multiprocessing.Pool(
        processes=multiprocessing.cpu_count() * 2,
        maxtasksperchild=2, …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing optimization multiprocessing python-multiprocessing

5
推荐指数
1
解决办法
2656
查看次数

与 multiprocessing.Pool 共享一个计数器

我想使用multiprocessing.Value+multiprocessing.Lock在不同的进程之间共享一个计数器。例如:

import itertools as it
import multiprocessing

def func(x, val, lock):
    for i in range(x):
        i ** 2
    with lock:
        val.value += 1
        print('counter incremented to:', val.value)

if __name__ == '__main__':
    v = multiprocessing.Value('i', 0)
    lock = multiprocessing.Lock()

    with multiprocessing.Pool() as pool:
        pool.starmap(func, ((i, v, lock) for i in range(25)))
    print(counter.value())
Run Code Online (Sandbox Code Playgroud)

这将引发以下异常:

RuntimeError:同步对象只能通过继承在进程之间共享

我最困惑的是,一个相关的(虽然不是完全类似的)模式适用于multiprocessing.Process()

if __name__ == '__main__':
    v = multiprocessing.Value('i', 0)
    lock = multiprocessing.Lock()

    procs = [multiprocessing.Process(target=func, args=(i, v, lock))
             for i in …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing multiprocessing python-3.x python-multiprocessing

5
推荐指数
1
解决办法
1553
查看次数

如何使用 kotlin 协程并行运行两个作业但等待另一个作业完成

我有以下这些工作:

  1. 异步膨胀复杂的日历视图
  2. 下载事件 A
  3. 下载事件 B
  4. 将日历视图添加到它的父级
  5. 将 A 事件添加到日历
  6. 将 B 事件添加到日历

我想要

  • 1, 2, 3 必须启动异步,
  • 4 必须等待 1
  • 5 必须等待 2 和 4
  • 6 必须等待 3 和 4
  • 6 和 5 不应该相互依赖,可以在不同的时间运行。
  • 4只依赖于1,所以它可以在2或3完成之前运行。

我试过 async await 但它使它们同时完成(如预期的那样)。我认为这个例子可能是学习并行编程概念的好方法,比如信号量互斥锁或自旋锁。但这太复杂了,我无法理解。

我应该如何使用 Kotlin 协程来实现这些?

parallel-processing coroutine kotlin

5
推荐指数
1
解决办法
3445
查看次数

lambda foreach parallelStream 创建的数据少于预期

我正在尝试实现一个数组列表的 lambda foreach 并行流,以提高现有应用程序的性能。

到目前为止,没有并行 Streamforeach 迭代创建了写入数据库的预期数据量。

但是当我切换到 parallelStream 时,它总是向数据库中写入更少的行。假设从预期的 10.000 行开始,将近 7000 行,但结果在这里有所不同。

知道我在这里缺少什么,数据竞争条件,还是必须使用锁和同步?

代码基本上是这样的:

// Create Persons from an arraylist of data

arrayList.parallelStream()
          .filter(d -> d.personShouldBeCreated())
          .forEach(d -> {

   // Create a Person
   // Fill it's properties
   // Update object, what writes it into a DB

  }
);
Run Code Online (Sandbox Code Playgroud)

到目前为止我尝试过的事情

将结果收集到一个新列表中...

collect(Collectors.toList())
Run Code Online (Sandbox Code Playgroud)

...然后迭代新列表并执行第一个代码片段中描述的逻辑。新“收集”的ArrayList的大小与预期结果匹配,但最后在数据库中创建数据仍然较少

更新/解决方案:

根据我在该代码中标记的关于非线程安全部分的答案(以及评论中的提示) ,我将其实现如下,最终给了我预期的数据量。性能有所提升,现在只需要执行之前的 1/3。

StringBuffer sb …
Run Code Online (Sandbox Code Playgroud)

java parallel-processing lambda java-8

5
推荐指数
1
解决办法
1261
查看次数

使用 gradle 并行运行 JUnit4 测试

我想使用 Gradle 并行运行测试。我的实验表明 Gradle 在类级别上并行运行测试。我需要在方法级别执行它们。

给定两个这样的测试类:

public class OneTest {
    @Test
    public void should_sleep_4_seconds() throws Exception {
        System.out.println("Start 4 seconds");
        Thread.sleep(4000);
        System.out.println("Done 4 seconds");
    }
}
Run Code Online (Sandbox Code Playgroud)
public class ManyTest {

    @Test
    public void should_sleep_1_seconds() throws Exception {
        System.out.println("Start 1 seconds");
        Thread.sleep(1000);
        System.out.println("Done 1 seconds");
    }

    @Test
    public void should_sleep_2_seconds() throws Exception {
        System.out.println("Start 2 seconds");
        Thread.sleep(2000);
        System.out.println("Done 2 seconds");
    }
}
Run Code Online (Sandbox Code Playgroud)

和这样的构建脚本:

plugins {
    id 'java'
}

sourceCompatibility = '1.8'
targetCompatibility = '1.8'

dependencies {
    testCompile 'junit:junit:4.12'
}

test { …
Run Code Online (Sandbox Code Playgroud)

java parallel-processing junit4 gradle

5
推荐指数
1
解决办法
4061
查看次数

furrr 没有找到自己的包

我目前正在开发一个包,假设它被称为 myPack。我有一个名为 myFunc1 的函数和另一个名为 myFunc2 的函数,它看起来像这样:

myFunc2 <- function(x, parallel = FALSE) { 
if(parallel) future::plan(future::multiprocess)
values <- furrr::future_map(x, myFunc1)
values
}
Run Code Online (Sandbox Code Playgroud)

现在,如果我在不并行的情况下调用 myFunc2,它会起作用。但是,如果我使用 parallel = TRUE 调用它,则会出现以下错误:

Error: Unexpected result (of class ‘snow-try-error’ != ‘FutureResult’)
retrieved for MultisessionFuture future (label = ‘<none>’, expression = 
‘{; do.call(function(...) {; ...future.f.env <- environment(...future.f); 
if (!is.null(...future.f.env$`~`)) {; if 
(is_bad_rlang_tilde(...future.f.env$`~`)) {; ...future.f.env$`~` <- 
base::`~`; ...; }); }, args = future.call.arguments); }’): there is no 
package called 'myPack'. This suggests that the communication with 
MultisessionFuture worker …
Run Code Online (Sandbox Code Playgroud)

parallel-processing r package r-future furrr

5
推荐指数
1
解决办法
513
查看次数

我的归并排序算法使用 OpenMP 时速度较慢,如何使其比序列化形式更快?

我正在研究并行编程并在排序算法上对其进行测试。我发现最简单的方法是使用 OpenMP,因为它提供了一种实现线程的简单方法。我做了研究,发现其他人已经这样做了,然后我尝试了一些代码。但是,当我perf stat -r 10 -d在 Linux 上测试它时,我得到的时间比序列化代码更糟糕(在某些情况下是两倍)。我尝试在数组上使用不同数量的元素,我使用的最大数量是 1.000.000 个数字,就好像我使用更多我收到错误一样。


void merge(int aux[], int left, int middle, int right){
    int temp[middle-left+1], temp2[right-middle];
    for(int i=0; i<(middle-left+1); i++){
        temp[i]=aux[left+i];
    }
    for(int i=0; i<(right-middle); i++){
        temp2[i]=aux[middle+1+i];
    }
    int i=0, j=0, k=left;
    while(i<(middle-left+1) && j<(right-middle))
    {
        if(temp[i]<temp2[j]){
            aux[k++]=temp[i++];
        }
        else{
            aux[k++]=temp2[j++];
        }
    }
    while(i<(middle-left+1)){
        aux[k++]=temp[i++];
    }
    while(j<(right-middle)){
        aux[k++]=temp2[j++];
    }
}

void mergeSort (int aux[], int left, int right){
    if (left < right){
        int middle = (left + right)/2;
        omp_set_num_threads(2);
        #pragma omp parallel …
Run Code Online (Sandbox Code Playgroud)

c++ parallel-processing mergesort openmp

5
推荐指数
1
解决办法
478
查看次数

Pthreads 与 Parallel 在 PHP 上的无限循环

我正在寻找一种在 PHP 上执行多线程的方法,并遇到了 pthreads PHP API,我认为这很容易实现(但是我必须找出如何安装支持 ZTS 的 PHP 版本对 Debian) .

问题是,当我查看 pthreads php.net 文档时,我发现了这个:

提示 考虑改用并行。

我不知道的。

我的目标是获取项目列表,并为每个项目打开一个 websocket,它永远监听某些更新。所以线程应该永远存在,如果被杀死,或者停止,或者这样,它应该重新启动(但是我认为我可以在外部处理这个)。

我不确定哪一种最适合这种情况。有什么推荐吗?

php parallel-processing pthreads

5
推荐指数
1
解决办法
1818
查看次数

如何在 Python 3 中导入线程包?

我想在 Python 3.6 中导入线程包。但是出现这个错误:

import thread
ModuleNotFoundError: No module named 'thread'
Run Code Online (Sandbox Code Playgroud)

parallel-processing multithreading python-3.x

5
推荐指数
1
解决办法
5554
查看次数