假设我们有一个列表并且想要选择满足某个属性的所有元素(比如一些函数 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 密集型任务。我想知道哪种方法更有效。非常感谢。
这个来自 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
我想使用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
我有以下这些工作:
我想要
我试过 async await 但它使它们同时完成(如预期的那样)。我认为这个例子可能是学习并行编程概念的好方法,比如信号量互斥锁或自旋锁。但这太复杂了,我无法理解。
我应该如何使用 Kotlin 协程来实现这些?
我正在尝试实现一个数组列表的 lambda foreach 并行流,以提高现有应用程序的性能。
到目前为止,没有并行 Stream的foreach 迭代创建了写入数据库的预期数据量。
但是当我切换到 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) 我想使用 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) 我目前正在开发一个包,假设它被称为 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) 我正在研究并行编程并在排序算法上对其进行测试。我发现最简单的方法是使用 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) 我正在寻找一种在 PHP 上执行多线程的方法,并遇到了 pthreads PHP API,我认为这很容易实现(但是我必须找出如何安装支持 ZTS 的 PHP 版本对 Debian) .
问题是,当我查看 pthreads php.net 文档时,我发现了这个:
提示 考虑改用并行。
我不知道的。
我的目标是获取项目列表,并为每个项目打开一个 websocket,它永远监听某些更新。所以线程应该永远存在,如果被杀死,或者停止,或者这样,它应该重新启动(但是我认为我可以在外部处理这个)。
我不确定哪一种最适合这种情况。有什么推荐吗?
我想在 Python 3.6 中导入线程包。但是出现这个错误:
import thread
ModuleNotFoundError: No module named 'thread'
Run Code Online (Sandbox Code Playgroud)