标签: parallel-processing

使用fork/join可以跨线程边界安全地移植非线程安全值吗?

我有一些不是线程安全的类:

class ThreadUnsafeClass {
  long i;

  long incrementAndGet() { return ++i; }
}
Run Code Online (Sandbox Code Playgroud)

(我在long这里使用了a 作为字段,但我们应该将其字段视为一些线程不安全的类型).

我现在有一个看起来像这样的课程

class Foo {
  final ThreadUnsafeClass c;

  Foo(ThreadUnsafeClass c) {
    this.c = c;
  }
}
Run Code Online (Sandbox Code Playgroud)

也就是说,线程不安全类是它的最后一个字段.现在我要这样做:

public class JavaMM {
  public static void main(String[] args) {
    final ForkJoinTask<ThreadUnsafeClass> work = ForkJoinTask.adapt(() -> {
      ThreadUnsafeClass t = new ThreadUnsafeClass();
      t.incrementAndGet();
      return new FC(t);
    });

    assert (work.fork().join().c.i == 1); 
  }
}
Run Code Online (Sandbox Code Playgroud)

也就是说,从thread T(main),我调用了一些工作T'(fork-join-pool),它创建并改变了我的不安全类的实例,然后返回包含在a中的结果Foo.请注意,我的线程不安全类的所有变异都发生在一个线程上T'.

问题1:我是否保证thread-unsafe-class实例的结束状态安全地移植到?的 …

java parallel-processing fork-join java-memory-model java-stream

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

R在外环中嵌套foreach%dopar%,在内环中嵌套%do%

我在R中运行以下脚本.如果我使用%do%而不是%dopar%,则脚本运行正常.但是,如果在外部循环中我使用%dopar%,则循环将永远运行而不会抛出任何错误(内存使用量会不断增加,直到内存不足为止).我正在使用16个核心.

library(parallel)
library(foreach)
library(doSNOW)
library(dplyr)


NumberOfCluster <- 16 
cl <- makeCluster(NumberOfCluster) 
registerDoSNOW(cl) 


foreach(i = UNSPSC_list, .packages = c('data.table', 'dplyr'), .verbose = TRUE) %dopar% 
    { 
      terms <- as.data.table(unique(gsub(" ", "", unlist(terms_list_by_UNSPSC$Terms[which(substr(terms_list_by_UNSPSC$UNSPSC,1,6) == i)])))) 
      temp <- inner_join(N_of_UNSPSCs_by_Term, terms, on = 'V1') 
      temp$V2 <- 1/as.numeric(temp$V2)
      temp <- temp[order(temp$V2, decreasing = TRUE),]
      names(temp) <- c('Term','Imp')
      ABNs <- unique(UNSPSCs_per_ABN[which(substr(UNSPSCs_per_ABN$UNSPSC,1,4) == substr(i,1,4)), 1])

      predictions <- as.numeric(vector()) 
      predictions <- foreach (j = seq(1 : nrow(train)), .combine = 'c', .packages = 'dplyr')  %do% 
      { 
        descr <- names(which(!is.na(train[j,]) == …
Run Code Online (Sandbox Code Playgroud)

parallel-processing foreach r

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

在Akka Dispatcher上启动后,为什么Futures中的Futures按顺序运行

当我们尝试从参与者的接收方法中启动多个期货时,我们观察到一种奇怪的行为。如果我们将配置的调度程序用作ExecutionContext,则期货将在同一线程上按顺序运行。如果我们使用ExecutionContext.Implicits.global,则期货将按预期并行运行。

我们将代码简化为以下示例(下面是一个更完整的示例):

implicit val ec = context.getDispatcher

Future{ doWork() } // <-- all running parallel
Future{ doWork() }
Future{ doWork() }
Future{ doWork() }

Future {
   Future{ doWork() } 
   Future{ doWork() } // <-- NOT RUNNING PARALLEL!!! WHY!!!
   Future{ doWork() }
   Future{ doWork() }
}
Run Code Online (Sandbox Code Playgroud)

一个可编译的示例如下:

import akka.actor.ActorSystem
import scala.concurrent.{ExecutionContext, Future}

object WhyNotParallelExperiment extends App {

  val actorSystem = ActorSystem(s"Experimental")   

  // Futures not started in future: running in parallel
  startFutures(runInFuture = false)(actorSystem.dispatcher)
  Thread.sleep(5000)

  // Futures started in future: running in …
Run Code Online (Sandbox Code Playgroud)

parallel-processing scala future akka akka-dispatcher

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

oracle中的文件ORA_DUMMY_FILE.f是什么?

oracle版本:12.2.0.1

如您所知,这些是oracle中并行服务器的unix进程:

ora_p000_ora12c
ora_p001_ora12c
....
ora_p???_ora12c
Run Code Online (Sandbox Code Playgroud)

它们也可以在视图中看到:gv $ px_process.可以从那里获得每个并行服务器的spid.

然后我在这里寻找与te并行服务器相关的打开文件:

ls -l /proc/<spid>/fd
Run Code Online (Sandbox Code Playgroud)

我正在为几个与此相同的并行服务器获取大约500-10000个文件描述符:

991 -> /u01/app/oracle/admin/ora12c/dpdump/676185682F2D4EA0E0530100007FFF5E/ORA_DUMMY_FILE.f (deleted)
Run Code Online (Sandbox Code Playgroud)

我已经删除了它们:(实际上我已经创建了一个小脚本,因为它有数千个)

gdb -p <spid>
    gdb> p close(<fd_id>)
Run Code Online (Sandbox Code Playgroud)

但几个小时后,文件描述符又开始被创建(每天数百个)

如果它们没有被删除,那么最终达到linux限制并且任何并行查询都会抛出这样的错误:

ORA-12801: error signaled in parallel query server P001
ORA-01116: error in opening database file 132
ORA-01110: data file 132: '/u02/oradata/ora12c/pdbname/tablespacenaname_ts_1.dbf'
ORA-27077: too many files open
Run Code Online (Sandbox Code Playgroud)

有没有人知道如何以及为什么要创建这个文件描述符,以及如何避免它?

编辑:添加了一些可能有用的信息.我已经测试过,当创建一个新的PDB时,会在其中创建一个目录DATA_PUMP_DIR(select*from all_directories),该目录指向:

/u01/app/oracle/admin/ora12c/dpdump/<xxxxxxxxxxxxx>
Run Code Online (Sandbox Code Playgroud)

linux目录也是创建的.还会创建一个文件描述符,指向新dpdump子目录中的ORA_DUMMY_FILE.f,就像最初描述的那样

lsof | grep "ORA_DUMMY_FILE.f (deleted)"

/u01/app/oracle/admin/ora12c/dpdump/<xxxxxxxxxxxxx>/ORA_DUMMY_FILE.f (deleted)
Run Code Online (Sandbox Code Playgroud)

这可能没问题,我面临的问题是指向ORA_DUMMY_FILE的文件描述符的持续增长达到linux限制.

oracle parallel-processing file-descriptor oracle12c

6
推荐指数
0
解决办法
413
查看次数

运行时的不同执行策略

在C++ 17中,algorithm标题中的许多函数现在可以采用执行策略.我可以举例来定义和调用这样的函数:

template <class ExecutionPolicy>
void f1(const std::vector<std::string>& vec, const std::string& elem, ExecutionPolicy&& policy) {
    const auto it = std::find(
        std::forward<ExecutionPolicy>(policy),
        vec.cbegin(), vec.cend(), elem
    );
}

std::vector<std::string> vec;
f1(vec, "test", std::execution::seq);
Run Code Online (Sandbox Code Playgroud)

但是我没有找到在运行时使用不同策略的好方法.例如,当我想根据某些输入文件使用不同的策略时.

我玩弄变种,但最后问题总是不同的类型std::execution::seq,std::execution::parstd::execution::par_unseq.

一个工作但繁琐的解决方案看起来像这样:

void f2(const std::vector<std::string>& vec, const std::string& elem, const int policy) {
    const auto it = [&]() {
        if (policy == 0) {
            return std::find(
                std::execution::seq,
                vec.cbegin(), vec.cend(), elem
            );
        }
        else if (policy == 1) {
            return std::find( …
Run Code Online (Sandbox Code Playgroud)

c++ parallel-processing

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

生成排列可以并行完成吗?

我想知道我是否可以加快排列的产生.具体来说,我正在使用[az]中的8个,我想使用[a-zA-Z]中的8个和[a-zA-Z0-9]中的8个.我所知道的将很快占用大量的时间和空间.

即使仅使用小写ASCII字符的长度为8的排列也需要一段时间并生成千兆字节.我的问题是我不理解底层算法,所以我无法弄清楚我是否可以将问题分解成比我以后可以加入的更小的任务.

我用来生成排列列表的python脚本:

import string
import itertools
from itertools import permutations

comb = itertools.permutations(string.ascii_lowercase, 8)

f = open('8letters.txt', 'w')
for x in comb:
        y = ''.join(x)
        f.write(y + '\n')

f.close()
Run Code Online (Sandbox Code Playgroud)

有谁知道如何将其划分为子任务并将它们放在一起?有可能吗?

我可能只是尝试(可能)更快的方式,但我遇到了C++及其std :: next_permutation()的问题,所以我无法验证它是否可以加速甚至一点点.

如果我可以将它分成16个任务,并在16个Xeon CPU上运行,那么加入结果,这将是很棒的.

python parallel-processing permutation combinatorics

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

shell脚本运行多个文件

我想在for循环中使用shell脚本,并行运行100个文件.

目前,我有一个以下格式的shell脚本:

#!/bin/bash
NUM=10
python a1.py $((NUM + 0)) &
python a2.py $((NUM + 2)) &
python a3.py $((NUM + 4)) &
python a4.py $((NUM + 6)) &
python a5.py $((NUM + 8)) &
Run Code Online (Sandbox Code Playgroud)

现在,如果我有a1.py,a2.py,a3.py... ... a100.py,我想并行运行它们,我怎么做,在for循环?

parallel-processing bash shell

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

Create_Matrix'RTextTools'包的并行计算

我正在创建一个DocumentTermMatrix使用create_matrix()RTextTools创建containermodel基于它.它适用于极大的数据集.

我为每个类别(因子级别)执行此操作.因此,对于每个类别,它必须运行矩阵,容器和模型.当我运行下面的代码(例如16核/ 64 GB)时 - 它只在一个核心中运行,并且使用的内存小于10%.

有没有办法加快这个过程?也许用doparallel&foreach?任何信息肯定会有所帮助.

#import the required libraries
library("RTextTools")
library("hash")
library(tm)

for ( n in 1:length(folderaddress)){
    #Initialize the variables
    traindata = list()
    matrix = list()
    container = list()
    models = list()
    trainingdata = list()
    results = list()
    classifiermodeldiv = 0.80`

    #Create the directory to place the models and the output files
    pradd = paste(combinedmodelsaveaddress[n],"SelftestClassifierModels",sep="")
    if (!file.exists(pradd)){
        dir.create(file.path(pradd))
    }  
    Data$CATEGORY <- as.factor(Data$CATEGORY)

    #Read the …
Run Code Online (Sandbox Code Playgroud)

parallel-processing foreach text-processing r doparallel

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

如何确定numba的prange实际上是否正常工作?

在另一个Q + A(我可以在pandas中执行动态cumsum?)我对使用prange这个代码的正确性做了评论(这个答案):

from numba import njit, prange

@njit
def dynamic_cumsum(seq, index, max_value):
    cumsum = []
    running = 0
    for i in prange(len(seq)):
        if running > max_value:
            cumsum.append([index[i], running])
            running = 0
        running += seq[i] 
    cumsum.append([index[-1], running])

    return cumsum
Run Code Online (Sandbox Code Playgroud)

评论是:

我不建议并行化一个不纯的循环.在这种情况下,running变量使其不纯.有4种可能的结果:(1)numba决定它不能并行处理它只是处理循环cumsum而不是prange(2)它可以将变量提升到循环之外并在余数上使用并行化(3)numba错误地插入并行执行和结果之间的同步可能是虚假的(4)numba在运行时插入必要的同步,这可能会比通过并行化首先获得更多的开销

而后来的补充:

当然,runningcumsum变量都使循环"不纯",而不仅仅是前面评论中所述的运行变量

然后我被问到:

这可能听起来像一个愚蠢的问题,但我怎么能弄清楚它做了哪4件事并改进了呢?我真的想用numba变得更好!

鉴于它可能对未来的读者有用,我决定在这里创建一个自我回答的Q + A. 掠夺者:我无法真正回答4个结果中的哪一个产生的问题(或者如果numba产生完全不同的结果),所以我非常鼓励其他答案.

python parallel-processing numba

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

如何在python中加快嵌套交叉验证?

从我发现的内容来看,还有一个其他问题(加速嵌套交叉验证),但是尝试在此站点和Microsoft上提出了一些修复建议后,安装MPI对我也不起作用,所以我希望有另一个软件包或回答这个问题。

我正在寻找比较多种算法和gridsearch各种参数(也许参数太多?)的方法,除了mpi4py之外还有什么方法可以加快我的代码的运行速度?据我了解,我不能使用n_jobs = -1,因为那是不嵌套的?

还要注意,我无法在下面尝试查看的许多参数上运行它(运行时间超过了我的时间)。如果我给每个模型仅两个参数进行比较,则只有2小时后才会有结果。另外,我在252行和25个特征列以及4个类别变量的数据集上运行此代码,以预测(“确定”,“可能”,“可能”或“未知”)某个基因(具有252个基因)是否影响疾病。使用SMOTE将样本大小增加到420,然后将其投入使用。

dataset= pd.read_csv('data.csv')
data = dataset.drop(["gene"],1)
df = data.iloc[:,0:24]
df = df.fillna(0)
X = MinMaxScaler().fit_transform(df)

le = preprocessing.LabelEncoder()
encoded_value = le.fit_transform(["certain", "likely", "possible", "unlikely"])
Y = le.fit_transform(data["category"])

sm = SMOTE(random_state=100)
X_res, y_res = sm.fit_resample(X, Y)

seed = 7
logreg = LogisticRegression(penalty='l1', solver='liblinear',multi_class='auto')
LR_par= {'penalty':['l1'], 'C': [0.5, 1, 5, 10], 'max_iter':[500, 1000, 5000]}

rfc =RandomForestClassifier()
param_grid = {'bootstrap': [True, False],
              'max_depth': [10, 20, 30, 40, 50, 60, 70, 80, 90, 100, None],
              'max_features': ['auto', 'sqrt'], …
Run Code Online (Sandbox Code Playgroud)

python parallel-processing scikit-learn cross-validation dask

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