标签: parallel-processing

如何使用 Pandas 并行读取 .xls?

我想使用 Pandas 并行读取一个大的 .xls 文件。目前我正在使用这个:

LARGE_FILE = "LARGEFILE.xlsx"
CHUNKSIZE = 100000 # processing 100,000 rows at a time

def process_frame(df):
      # process data frame
      return len(df)

if __name__ == '__main__':
      reader = pd.read_excel(LARGE_FILE, chunksize=CHUNKSIZE)
      pool = mp.Pool(4) # use 4 processes

      funclist = []
      for df in reader:
              # process each data frame
              f = pool.apply_async(process_frame,[df])
              funclist.append(f)

      result = 0
      for f in funclist:
              result += f.get(timeout=10) # timeout in 10 seconds
Run Code Online (Sandbox Code Playgroud)

虽然这会运行,但我认为它实际上并没有加快读取文件的过程。有没有更有效的方法来实现这一目标?

python parallel-processing pandas

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

并行 foreach 与并行适用于 R 吗?

在哪些情况下哪个更有效?在某些情况下,哪一个根本无法工作?

我试图使一些通用代码更有效,并且很好奇哪个更好,因为据我所知,它们不能结合使用。

以供参考:

library(doParallel)
library(foreach)

foreach (i = list) %dopar% {
...
}
Run Code Online (Sandbox Code Playgroud)

对比

library(parallel)

parLapply(cl, X = list, fun = function)
Run Code Online (Sandbox Code Playgroud)

parallel-processing performance foreach r vectorization

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

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
查看次数

并行 STL 是否处理插入迭代器,例如 std::back_insert_iterator?

并行 STL 算法是否符合std::back_insert_iterator??

我可能误解了std::par和之间的区别std::par_vec,是否std::par_vec意味着需要预先分配输出范围?

代码示例:

auto numbers = {1,2,3,4,5,6};
auto squared = std::vector<int>{};
std::transform(
  **std::par/std::par_vec,**
  numbers.begin(),
  numbers.end(),
  std::back_inserter(squared),
  [](auto val) { 
    return val*val; 
  }
);
Run Code Online (Sandbox Code Playgroud)

更新

简化问题作为我的第一个问题是误读文章的结果。

c++ parallel-processing stl stl-algorithm c++17

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