我想使用 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)
虽然这会运行,但我认为它实际上并没有加快读取文件的过程。有没有更有效的方法来实现这一目标?
在哪些情况下哪个更有效?在某些情况下,哪一个根本无法工作?
我试图使一些通用代码更有效,并且很好奇哪个更好,因为据我所知,它们不能结合使用。
以供参考:
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) 有没有办法在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) 我进行了大量研究,以弄清楚如何从应用程序的源代码为应用程序创建 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
有没有办法阻止我的循环被速率限制中断?如果可能的话,我希望我的代码等到时间限制过后才执行。
一个附带问题:我考虑过并行化 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) 假设我有一个 8 核 CPU。doParallel在 R 中使用,当我注册时makeCluster(x),理想的核心数是x多少?
是尽可能多的核心吗?或者使用 7 核会比使用 6 核慢吗?有什么规则可以解决这个问题吗?
我试图了解 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) 我正在尝试通过hargreaves方法计算蒸发量package SPEI。这涉及使用最低温度 ( TMIN) 和最高温度 ( TMAX)。鉴于此Tmin,并行计算是我最好的选择,并且Tmax rasterstacks拥有500,000 cells and 100 layers each.
Hargreaves function取Tmin,Tmax并latitude在each 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连接到它。在 中pet,k$d是 Tmin,k$d是 Tmax(也许我应该在petegfunction(k,y)而不是仅仅提供两个参数k?) …
该包的摘要的java.util.stream规定如下:
有状态 lambda 的一个例子是
map()in的参数:Run Code Online (Sandbox Code Playgroud)Set<Integer> seen = Collections.synchronizedSet(new HashSet<>()); stream.parallel().map(e -> { if (seen.add(e)) return 0; else return e; })...在这里,如果映射操作是并行执行的,由于线程调度差异,相同输入的结果可能会因运行而异,而对于无状态 lambda 表达式,结果将始终相同。
我不明白为什么这不会产生一致的结果,因为该集合是同步的并且一次只能处理一个元素。您能否以一种演示结果如何因并行化而变化的方式完成上述示例?
并行 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)
更新
简化问题作为我的第一个问题是误读文章的结果。
r ×4
algorithm ×1
apache-spark ×1
c++ ×1
c++17 ×1
dataflow ×1
doparallel ×1
foreach ×1
function ×1
java ×1
java-8 ×1
java-stream ×1
open-source ×1
pandas ×1
performance ×1
python ×1
raster ×1
scala ×1
stl ×1
stream ×1
twitter ×1