标签: parallel-processing

如何并行执行多个任务?

我参加了Parallel Programming 课程,它显示了并行接口:

def parallel[A, B](taskA: => A, taskB: => B): (A, B) = {
  val ta = taskA
  val tb = task {taskB}
  (ta, tb.join())
}
Run Code Online (Sandbox Code Playgroud)

以下是错误的:

def parallel[A, B](taskA: => A, taskB: => B): (A, B) = {
  val ta = taskB
  val tb = task {taskB}.join()
  (ta, tb)
}
Run Code Online (Sandbox Code Playgroud)

更多界面见https://gist.github.com/ChenZhongPu/fe389d30626626294306264a148bd2aa

它还向我们展示了执行四个任务的正确方法:

def parallel[A, B, C, D](taskA: => A, taskB: => B, taskC: => C, taskD: => D): (A, B, C, D) = {
    val …
Run Code Online (Sandbox Code Playgroud)

parallel-processing concurrency scala

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

R 中 GLM 的并行化循环

我正在尝试编写一个并行化的 for 循环,在其中我试图以最佳方式找到最佳 GLM 以仅对具有最低 p 值的变量进行建模,以查看我是否要打网球(是/否,二进制) .

例如,我有一个包含气象数据集的表(及其数据框)。我首先通过查看这些模型中的哪一个 p 值最低来构建 GLM 模型

PlayTennis ~ Precip
PlayTennis ~ Temp, 
PlayTennis ~ Relative_Humidity
PlayTennis ~ WindSpeed)
Run Code Online (Sandbox Code Playgroud)

假设PlayTennis ~ Precip具有最低的 p 值。因此,repeat 中的下一个循环迭代是查看哪些其他变量将具有最低的 p 值。

PlayTennis ~ Precip + Temp
PlayTennis ~ Precip + Relative_Humidity 
PlayTennis ~ Precip + WindSpeed
Run Code Online (Sandbox Code Playgroud)

这将一直持续到没有更显着的变量(P 值大于 0.05)。因此我们得到了PlayTennis ~ Precip + WindSpeed(这都是假设的)的最终输出。

是否有关于如何在各种内核上并行化此代码的建议?我遇到了一个speedglm从库 speedglm调用的 glm 的新函数。这确实有所改善,但并没有太大改善。我还研究了foreach循环,但我不确定它如何与每个线程通信以了解各个运行的 p 值是更大还是更低。预先感谢您的任何帮助。

d =

Time          Precip    Temp    Relative_Humidity   WindSpeed   …   PlayTennis    
1/1/2000 0:00 …
Run Code Online (Sandbox Code Playgroud)

parallel-processing r glm parallel-foreach

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

clock_gettime() 对比 gettimeofday() 用于测量 OpenMP 执行时间

我正在编写一些 C 代码,它实现了一个三重嵌套 for 循环来计算矩阵乘法,同时使用 OpenMP 对其进行并行化。我试图准确地测量从 for 循环开始到结束所花费的时间。到目前为止,我一直在使用 gettimeofday(),但我注意到有时感觉它没有准确记录执行 for 循环所花费的时间。似乎是在说它比实际花费的时间更长。

这是原始代码:

struct timeval start end;
double elapsed;

gettimeofday(&start, NULL);
#pragma omp parallel for num_threads(threads) private(i, j, k)
for(...)
{
 ...
 for(...)
 {
  ...
  for(...)
  {
   ...
  }
 }
}

gettimeofday(&end, NULL);
elapsed = (end.tv_sec+1E-6*end.tv_usec) - (start.tv_sec+1E-6*start.tv_usec)
Run Code Online (Sandbox Code Playgroud)

这是使用clock_gettime()的相同代码:

 struct timespec start1, finish1;
 double elapsed1;

clock_gettime(CLOCK_MONOTONIC, &start1);

  #pragma omp parallel for num_threads(threads) private(i, j, k)
    for(...)
    {
     ...
     for(...)
     {
      ...
      for(...)
      {
       ...
      }
     }
    }

clock_gettime(CLOCK_MONOTONIC, &finish1);
elapsed1 …
Run Code Online (Sandbox Code Playgroud)

c parallel-processing time openmp gettimeofday

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

scala 延迟并行集合(可能吗?)

有没有办法在scala中并行运行流而不将所有对象加载到内存中?

注意:使用 par 方法,会将所有对象加载到内存中

val list = "a"::"b"::"c"::"d"::"e"::Nil //> list: List[String] = List(a, b, c, d, e)

val s = list.toStream  //> s: scala.collection.immutable.Stream[String] = Stream(a, ?)
val sq = s.par         //> sq: scala.collection.parallel.immutable.ParSeq[String] = ParVector(a, b, c, d, e)
sq.map { x => println("Map 1 "+x);x }
  .map { x => println("Map 2 "+x);x}
  .map { x => println("Map 3 "+x);x }
  .foreach { x => println("done "+x)}  
Run Code Online (Sandbox Code Playgroud)

parallel-processing scala stream

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

如何从源代码为任何应用程序创建数据流图 (DFG/SDFG)

我进行了大量研究,以弄清楚如何从应用程序的源代码为应用程序创建 DFG。对于某些应用程序,例如 MP3 解码器、JPEG 压缩和 H.263 解码器,可以在线使用 DFG。

我无法弄清楚如何从源代码中为 HEVC 等应用程序创建 DFG?是否有任何工具可以为此类复杂的应用程序立即生成数据流图,还是必须手动完成?

请就此事给我建议。

编辑:我将 Doxygen 用于 HEVC,我可以看到不同的功能如何相互交互。然而,每个函数都有许多入口和出口点,一段时间后 Doxygen 的输出变得太混乱而无法理解。

我还看了 StreamIt:http ://camlunity.ru/swap/Library/Conflux/Stream%20Programming/streamit-cc_stream_graph_programming_language.pdf

它看起来很方便,但它为更简单的应用程序(如 MP3 解码器)生成的图表太复杂了。为了生成连贯的 DFG,我是否必须重新编写整个源代码?

algorithm parallel-processing open-source dataflow data-structures

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

使用 rtweet get_timeline() 避免速率限制

有没有办法阻止我的循环被速率限制中断?如果可能的话,我希望我的代码等到时间限制过后才执行。

一个附带问题:我考虑过并行化 for 循环。你认为这会是个好主意吗?我不确定是否有机会将数据写入错误的文件。

library(rtweet)
create_token(app="Arconic Influential Followers",consumer_key,consumer_secret) 

flw <- get_followers("arconic")
fds <- get_friends("arconic")
usrs <- lookup_users(c(flw$user_id, fds$user_id))

for(i in 1:length(usrs$user_id)){

    a<-tryCatch({get_timeline(usrs$user_id[i])},
                error=function(e){message(e)}
       )
    tryCatch({save_as_csv(a,usrs$user_id[i])},
                error=function(e){message(e)}
       )

}
Run Code Online (Sandbox Code Playgroud)

twitter parallel-processing r

5
推荐指数
2
解决办法
1840
查看次数

并行处理的最佳内核数是多少?

假设我有一个 8 核 CPU。doParallel在 R 中使用,当我注册时makeCluster(x),理想的核心数是x多少?

是尽可能多的核心吗?或者使用 7 核会比使用 6 核慢吗?有什么规则可以解决这个问题吗?

parallel-processing r doparallel

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

在 Spark mapParitions 中使用 Java 8 parallelStream

我试图了解 Spark 并行性中 Java 8 并行流的行为。当我运行下面的代码时,我期望的输出大小listOfThings与输入大小相同。但事实并非如此,我有时会在输出中丢失项目。这种行为是不一致的。如果我只是遍历迭代器而不是使用parallelStream,那么一切都很好。每次都计算匹配。

// listRDD.count = 10
JavaRDD test = listRDD.mapPartitions(iterator -> {
    List listOfThings = IteratorUtils.toList(iterator);
    return listOfThings.parallelStream.map(
        //some stuff here
    ).collect(Collectors.toList());
});
// test.count = 9
// test.count = 10
// test.count = 8
// test.count = 7
Run Code Online (Sandbox Code Playgroud)

parallel-processing java-8 apache-spark spark-streaming

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

parLapply 多个参数 R

我正在尝试通过hargreaves方法计算蒸发量package SPEI。这涉及使用最低温度 ( TMIN) 和最高温度 ( TMAX)。鉴于此Tmin,并行计算是我最好的选择,并且Tmax rasterstacks拥有500,000 cells and 100 layers each. Hargreaves functionTminTmaxlatitudeeach grid作为输入。以下是我的第一个猜测如何解决这个问题:

library(SPEI)
# go parallel 
library(parallel)
clust <- makeCluster(detectCores())

#har <- hargreaves(TMIN,TMAX,lat=37.6475) # get evaporation for a station. 
Run Code Online (Sandbox Code Playgroud)

但是,我的数据是网格化的。

Tmin并且Tmax是名单中,每个数据帧Tmin,并Tmax有一个$latitude连接到它。在 中petk$d是 Tmin,k$d是 Tmax(也许我应该在petegfunction(k,y)而不是仅仅提供两个参数k?) …

parallel-processing r function raster

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

Java Stream 有状态行为示例

包的摘要java.util.stream规定如下:

有状态 lambda 的一个例子是map()in的参数:

Set<Integer> seen = Collections.synchronizedSet(new HashSet<>());
stream.parallel().map(e -> { if (seen.add(e)) return 0; else return e; })...
Run Code Online (Sandbox Code Playgroud)

在这里,如果映射操作是并行执行的,由于线程调度差异,相同输入的结果可能会因运行而异,而对于无状态 lambda 表达式,结果将始终相同。

我不明白为什么这不会产生一致的结果,因为该集合是同步的并且一次只能处理一个元素。您能否以一种演示结果如何因并行化而变化的方式完成上述示例?

java parallel-processing synchronization java-stream

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