标签: parallel-processing

txtProgressBar用于并行引导程序无法正常显示

下面是我的问题的MWE:我已经使用引导程序(通过引导程序包中的引导功能)为某些功能编写了进度条.

只要我不使用并行处理(res_1core下面),这样就可以正常工作.如果我想通过设置parallel = "multicore"和使用并行处理ncpus = 2,则进度条显示不正确(res_2core如下).

library(boot)

rsq <- function(formula, data, R, parallel = c("no", "multicore", "snow"), ncpus = 1) {
  env <- environment()
  counter <- 0
  progbar <- txtProgressBar(min = 0, max = R, style = 3)
  bootfun <- function(formula, data, indices) {
    d <- data[indices,]
    fit <- lm(formula, data = d)
    curVal <- get("counter", envir = env)
    assign("counter", curVal + 1, envir = env)
    setTxtProgressBar(get("progbar", envir = env), curVal + …
Run Code Online (Sandbox Code Playgroud)

parallel-processing r progress-bar statistics-bootstrap

6
推荐指数
1
解决办法
567
查看次数

为什么在4核超线程CPU上使用8个线程比4个线程更快?

我有一个四核i7 920 CPU.它是超线程的,因此计算机认为它有8个核心.

从我在interweb上看到的,在执行并行任务时,我应该使用物理内核的数量,而不是超线程内核的数量.

所以我做了一些时间,并且惊讶地发现在并行循环中使用8个线程比使用4个线程更快.

为什么是这样?我的示例代码太长了,无法在此处发布,但可以通过运行以下示例找到:https://github.com/jsphon/MTVectorizer

性能图表在这里:

在此输入图像描述

python parallel-processing numpy numba

6
推荐指数
1
解决办法
1183
查看次数

在Swift中使用Grand Central Dispatch来并行化并加速"for"循环?

我试图围绕如何使用GCD来并行化和加速蒙特卡罗模拟.大多数/所有简单示例都是针对Objective C提供的,我真的需要一个Swift的简单示例,因为Swift是我的第一个"真正的"编程语言.

Swift中蒙特卡罗模拟的最小工作版本将是这样的:

import Foundation

import Cocoa
var winner = 0
var j = 0
var i = 0
var chance = 0
var points = 0
for j=1;j<1000001;++j{
    var ability = 500

    var player1points = 0

    for i=1;i<1000;++i{
        chance = Int(arc4random_uniform(1001))
        if chance<(ability-points) {++points}
        else{points = points - 1}
    }
    if points > 0{++winner}
}
    println(winner)
Run Code Online (Sandbox Code Playgroud)

代码可以直接粘贴到xcode 6.1中的命令行程序项目中

最内层的循环不能并行化,因为变量"points"的新值在下一个循环中使用.但最外面的只是运行最里面的模拟1000000次并计算结果,应该是并行化的理想候选者.

所以我的问题是如何使用GCD并行化最外层的for循环?

parallel-processing macos grand-central-dispatch swift

6
推荐指数
1
解决办法
2609
查看次数

通过分离#omp parallel和#omp for来减少OpenMP fork/join开销

我正在阅读Peter S. Pacheco 对并行编程的介绍.在5.6.2节中,它提供了一个关于减少fork/join开销的有趣讨论.考虑奇偶换位排序算法:

for(phase=0; phase < n; phase++){
    if(phase is even){
#       pragma omp parallel for default(none) shared(n) private(i)
        for(i=1; i<n; i+=2){//meat}
    }
    else{
#       pragma omp parallel for default(none) shared(n) private(i)
        for(i=1; i<n-1; i+=2){//meat}
    }
}
Run Code Online (Sandbox Code Playgroud)

作者认为上面的代码有一些高的fork/join开销.因为线程在外循环的每次迭代中分叉并连接.因此,他提出以下版本:

# pragma omp parallel default(none) shared(n) private(i, phase)
for(phase=0; phase < n; phase++){
    if(phase is even){
#       pragma omp for
        for(i=1; i<n; i+=2){//meat}
    }
    else{
#       pragma omp for
        for(i=1; i<n-1; i+=2){//meat}
    }
}
Run Code Online (Sandbox Code Playgroud)

根据作者的说法,第二个版本在外部循环开始之前分叉线程,并为每次迭代重用线程,从而产生更好的性能.

但是,我怀疑第二个版本的正确性.在我的理解中,#pragma omp parallel指令启动一组线程并让线程并行执行以下结构化块.在这种情况下,结构化块应该是整个外部for循环 …

parallel-processing multithreading openmp

6
推荐指数
1
解决办法
3976
查看次数

Java 8流和parallelStream

假设我们有Collection这样的:

Set<Set<Integer>> set = Collections.newSetFromMap(new ConcurrentHashMap<>());
for (int i = 0; i < 10; i++) {
    Set<Integer> subSet = Collections.newSetFromMap(new ConcurrentHashMap<>());
    subSet.add(1 + (i * 5));
    subSet.add(2 + (i * 5));
    subSet.add(3 + (i * 5));
    subSet.add(4 + (i * 5));
    subSet.add(5 + (i * 5));
    set.add(subSet);
}
Run Code Online (Sandbox Code Playgroud)

并处理它:

set.stream().forEach(subSet -> subSet.stream().forEach(System.out::println));
Run Code Online (Sandbox Code Playgroud)

要么

set.parallelStream().forEach(subSet -> subSet.stream().forEach(System.out::println));
Run Code Online (Sandbox Code Playgroud)

要么

set.stream().forEach(subSet -> subSet.parallelStream().forEach(System.out::println));
Run Code Online (Sandbox Code Playgroud)

要么

set.parallelStream().forEach(subSet -> subSet.parallelStream().forEach(System.out::println));
Run Code Online (Sandbox Code Playgroud)

所以,有人可以解释我:

  • 他们之间有什么区别?
  • 哪一个更好?快点?更安全?
  • 哪一个适合大量收藏?
  • 当我们想对每个项目应用繁重的流程时,哪一个是好的?

java collections parallel-processing java-8 java-stream

6
推荐指数
1
解决办法
822
查看次数

Go(lang)的地址空间是多少?

我尝试了解Go中并发编程的基础知识.几乎所有文章都使用术语"地址空间",例如:"所有goroutines共享相同的地址空间".这是什么意思?

我试图从wiki中理解以下主题,但它没有成功:

但是目前我很难理解,因为我在内存管理和并发编程等方面的知识非常差.有许多未知的单词,如段,页面,相对/绝对地址,VAS等.

有人可以向我解释问题的基础吗?可能有一些有用的文章,我找不到.

memory parallel-processing concurrency memory-management go

6
推荐指数
1
解决办法
1203
查看次数

如何保持线程执行,直到异步线程返回回调

我有如下图所示的场景

在此输入图像描述

这里主线程是my java application.it打开一个WM线程执行.WM处理任务执行.他需要调用执行任务的数量.假设它包含任务T1,T2,T3

T3取决于T2,T2取决于T1.WM首先调用RM来执行T1的任务执行.T1可以在寻呼时给出响应,也可以在T1完成后给出响应.

问题是我如何等待T1完成然后开始T2的执行.当T1部分完成时,如何在分页中发送数据时通知WM.

这是一个简单的场景,但在T1,T2,T3,T4的情况下.T3取决于T1和T2.

代码:

public class TestAsync implements TaskCallBack {
    public static ExecutorService exService = Executors.newFixedThreadPool(5);
    public static void main(String args[]) throws InterruptedException, ExecutionException{
        Task t1 = new Task();
        t1.doTask(new TestAsync());

    }

    public static ExecutorService getPool(){
        return exService;
    }

    @Override
    public void taskCompleted(String obj) {
        System.out.println(obj);
    }
}

class Task {
 public void doTask(TaskCallBack tcb) throws InterruptedException, ExecutionException{
     FutureTask<String> ft = new FutureTask<>(new Task1());
     TestAsync.getPool().execute(ft);
     tcb.taskCompleted(ft.get());
 }

}

class Task1 implements Callable<String>{

    @Override
    public String call() …
Run Code Online (Sandbox Code Playgroud)

java parallel-processing concurrency multithreading executorservice

6
推荐指数
1
解决办法
2734
查看次数

并行预测

我试图predict()在我的Windows机器上并行运行.这适用于较小的数据集,但不能很好地扩展,因为每个进程都会创建新的数据框副本.有没有办法如何并行运行而不制作临时副本?

我的代码(这个原始代码只有少量修改):

library(foreach)
library(doSNOW)

fit <- lm(Employed ~ ., data = longley)
scale <- 100
longley2 <- (longley[rep(seq(nrow(longley)), scale), ])

num_splits <-4
cl <- makeCluster(num_splits)
registerDoSNOW(cl)  

split_testing<-sort(rank(1:nrow(longley))%%num_splits)

predictions<-foreach(i= unique(split_testing),
                     .combine = c, .packages=c("stats")) %dopar% {
                       predict(fit, newdata=longley2[split_testing == i, ])
                     }
stopCluster(cl)
Run Code Online (Sandbox Code Playgroud)

我正在使用简单的数据复制来测试它.有scale10或1000它正在工作,但我想让它运行scale <- 1000000- 具有16M行的数据帧(1.86GB数据帧,如object_size()from所示pryr.注意,必要时我也可以使用Linux机器,如果这是唯一的选择.

parallel-processing r predict

6
推荐指数
1
解决办法
1245
查看次数

多处理python没有并行运行

我一直在尝试使用python中的多处理模块来实现计算成本高昂的任务的并行性.

我能够执行我的代码,但它并不是并行运行的.我一直在阅读多处理的手册页和foruns,以找出它为什么不工作,我还没有想出来.

我认为这个问题可能与执行我创建和导入的其他模块的某种锁有关.

这是我的代码:

main.py:

##import my modules
import prepare_data
import filter_part
import wrapper_part
import utils
from myClasses import ML_set
from myClasses import data_instance

n_proc = 5

def main():
    if __name__ == '__main__':
        ##only main process should run this
        data = prepare_data.import_data() ##read data from file  
        data = prepare_data.remove_and_correct_outliers(data)
        data = prepare_data.normalize_data_range(data)
        features = filter_part.filter_features(data)

        start_t = time.time()
        ##parallelism will be used on this part
        best_subset = wrapper_part.wrapper(n_proc, data, features)

        print time.time() - start_t


main()
Run Code Online (Sandbox Code Playgroud)

wrapper_part.py:

##my modules
from myClasses …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing multiprocessing python-multiprocessing

6
推荐指数
1
解决办法
4936
查看次数

无法使用MPI_Send和MPI_Recv发送std :: vector

我正在尝试使用MPI send和recv函数发送std:vector但我没有到达哪里.我得到的错误就像

Fatal error in MPI_Recv: Invalid buffer pointer, error stack:
MPI_Recv(186): MPI_Recv(buf=(nil), count=2, MPI_INT, src=0, tag=0, MPI_COMM_WORLD, status=0x7fff9e5e0c80) failed
MPI_Recv(124): Null buffer pointer
Run Code Online (Sandbox Code Playgroud)

我尝试了多种组合

A)像用于发送数组的那些..

 std::vector<uint32_t> m_image_data2; // definition of  m_image_data2
     m_image_data2.push_back(1);
     m_image_data2.push_back(2);
     m_image_data2.push_back(3);
     m_image_data2.push_back(4);
     m_image_data2.push_back(5);

MPI_Send( &m_image_data2[0], 2, MPI_INT, 1, 0, MPI_COMM_WORLD);
MPI_Send( &m_image_data2[2], 2, MPI_INT, 1, 0, MPI_COMM_WORLD);

MPI_Recv( &m_image_data2[0], 2, MPI_INT, 0, 0, MPI_COMM_WORLD, &status );
Run Code Online (Sandbox Code Playgroud)

B)没有[]

MPI_Send( &m_image_data2, 2, MPI_INT, 1, 0, MPI_COMM_WORLD);
MPI_Send( &m_image_data2 + 2, 2, MPI_INT, 1, 0, MPI_COMM_WORLD);

MPI_Recv( &m_image_data2, 2, …
Run Code Online (Sandbox Code Playgroud)

c++ parallel-processing hpc vector mpi

6
推荐指数
1
解决办法
1万
查看次数