使用流计算笛卡尔积时,我可以并行生成它们,并按顺序使用它们,以下代码演示:
int min = 0;
int max = 9;
Supplier<IntStream> supplier = () -> IntStream.rangeClosed(min, max).parallel();
supplier.get()
.flatMap(a -> supplier.get().map(b -> a * b))
.forEachOrdered(System.out::println);
Run Code Online (Sandbox Code Playgroud)
这将按顺序完美打印所有内容,现在考虑以下代码,我想将其添加到列表中,同时保留顺序.
int min = 0;
int max = 9;
Supplier<IntStream> supplier = () -> IntStream.rangeClosed(min, max).parallel();
List<Integer> list = supplier.get()
.flatMap(a -> supplier.get().map(b -> a * b))
.boxed()
.collect(Collectors.toList());
list.forEach(System.out::println);
Run Code Online (Sandbox Code Playgroud)
现在它不按顺序打印!
鉴于我没有要求保留订单,这是可以理解的.
现在的问题是:有没有办法collect()或有没有Collector保留秩序?
我有以下代码:
FTP ... do |ftp|
files.each do |file|
...
ftp.put(file)
sleep 1
end
end
Run Code Online (Sandbox Code Playgroud)
我想以单独的线程或某种并行的方式运行每个文件.这样做的正确方法是什么?这是对的吗?
这是我对并行宝石的尝试
FTP ... do |ftp|
Parallel.map(files) do |file|
...
ftp.put(file)
sleep 1
end
end
Run Code Online (Sandbox Code Playgroud)
并行的问题是put/outputs可以同时发生,如下所示:
as = [1,2,3,4,5,6,7,8]
results = Parallel.map(as) do |a|
puts a
end
Run Code Online (Sandbox Code Playgroud)
我怎样才能强制看跌,就像他们通常会分开一样.
我打算将一些计算卸载到Xeon Phi,但是想先测试不同的API和不同的并行编程.
是否有适用于Xeon Phi(Windows或Linux)的模拟器/模拟器?
假设我有一个处理100万句话的任务.
对于每个句子,我需要对它做一些事情,无论处理它们的具体顺序如何.
在我的Java程序中,我有一组从我的主要工作块中划分出来的一组未来,它用一个可调用来定义要在一大块句子上完成的工作单元,我正在寻找一种优化线程数量的方法分配工作通过大块的句子,然后重新组合每个线程的所有结果.
在我看到收益递减之前,我可以使用的最大线程数是多少?
另外,是什么原因导致逻辑分配的线程越多,即一次完成的线程越多,就越不正确?
我一直在使用R中的库'doParallel'来提高一组函数的速度.但是,我遇到了一个我无法解决的错误.我相信以下代码隔离了问题的精髓:
library(Matrix)
library(doParallel)
test_mat = Matrix(c(0,1,2,NA,0,0,2,NA,1,NA,1,2,2,NA,0,1,0,2,2,2,0,0,NA,NA,1,2,1,1,2,1,rep(NA,5)), ncol=7, byrow=TRUE, sparse=TRUE)
par_func <- function(mat, ncores)
{
cl <- makePSOCKcluster(ncores)
clusterSetRNGStream(cl)
registerDoParallel(cl, cores = ncores)
df = data.frame(1:7, NA)
temp_vec = foreach(i=iter(df, by='row'), .combine=rbind) %dopar%
{
i[,2] <- sum(mat[,i[,1]] == 1, na.rm = TRUE) + 1
}
stopCluster(cl)
return(temp_vec)
}
par_func(mat=test_mat, ncores=5)
Run Code Online (Sandbox Code Playgroud)
这会产生以下错误消息:
Error in { : task 1 failed - "object of type 'S4' is not subsettable"
Run Code Online (Sandbox Code Playgroud)
如果'mat'是'matrix'类而不是'dgCMatrix',则此函数有效,因此问题似乎是由于稀疏矩阵的子集化.我有什么选择可以解决这个问题吗?矩阵"mat"可以非常大并且可以包含许多零,因此我想继续使用稀疏矩阵.
我有一个包含+100,000个文件的输入文件夹.
我想对它们进行批量操作,即以某种方式重命名所有这些操作,或者根据每个文件名称中的信息将它们移动到新路径.
我想使用Spark来做到这一点,但不幸的是,当我尝试下面这段代码时:
final org.apache.hadoop.fs.FileSystem ghfs = org.apache.hadoop.fs.FileSystem.get(new java.net.URI(args[0]), new org.apache.hadoop.conf.Configuration());
org.apache.hadoop.fs.FileStatus[] paths = ghfs.listStatus(new org.apache.hadoop.fs.Path(args[0]));
List<String> pathsList = new ArrayList<>();
for (FileStatus path : paths) {
pathsList.add(path.getPath().toString());
}
JavaRDD<String> rddPaths = sc.parallelize(pathsList);
rddPaths.foreach(new VoidFunction<String>() {
@Override
public void call(String path) throws Exception {
Path origPath = new Path(path);
Path newPath = new Path(path.replace("taboola","customer"));
ghfs.rename(origPath,newPath);
}
});
Run Code Online (Sandbox Code Playgroud)
我得到一个错误,hadoop.fs.FileSystem不是Serializable(因此可能不能用于并行操作)
知道如何解决它或以其他方式完成它吗?
来自插入符R包的parRF不适合我使用多个核心,这是非常具有讽刺意味的,因为parRF中的par表示并行.我在Windows机器上,如果这是一个相关的信息.我检查过我正在使用最新的关于插入符号和doParallel的最新内容.
我做了一个最小的例子并给出了下面的结果.有任何想法吗?
源代码
library(caret)
library(doParallel)
trCtrl <- trainControl(
method = "repeatedcv"
, number = 2
, repeats = 5
, allowParallel = TRUE
)
# WORKS
registerDoParallel(1)
train(form = Species~., data=iris, trControl = trCtrl, method="parRF")
closeAllConnections()
# FAILS
registerDoParallel(2)
train(form = Species~., data=iris, trControl = trCtrl, method="parRF")
closeAllConnections()
Run Code Online (Sandbox Code Playgroud)
产量
> library(caret)
> library(doParallel)
>
> trCtrl <- trainControl(
+ method = "repeatedcv"
+ , number = 2
+ , repeats = 5
+ , allowParallel = TRUE
+ …Run Code Online (Sandbox Code Playgroud) 当我跑来make -j3并行构建时,我明白了
warning: -jN forced in submake: disabling jobserver mode.
Run Code Online (Sandbox Code Playgroud)
在文档中我发现了警告
如果make检测到子系统可以通信的系统上与并行处理相关的错误条件.
这些错误情况是什么?我该怎么做才能治愈它们或抑制错误信息?
makefile是从CMake生成的,所以我不能(=我不想)编辑makefile.
我目前正在制定一个开放式的提议,为我正在开发的项目带来并行功能,但我遇到了一个障碍find_end.
现在find_end可以描述为:
一种算法,用于搜索[first,last]范围内元素[s_first,s_last]的最后一个子序列.第一个版本使用operator ==来比较元素,第二个版本使用给定的二元谓词p.
它的要求由cppreference列出.现在我没有问题并行find/ findif/ findifnot等等.这些可以很容易地分成异步执行的单独分区,我没有遇到任何麻烦.问题find_end是将算法拆分成块不是解决方案,因为如果我们说一个向量:
1 2 3 4 5 1 2 3 8
我们想要搜索1 2.
好的,首先我将矢量异步分隔成块,然后只搜索每个块中的范围吧?看起来很容易,但是如果由于某种原因只有3个可用内核会发生什么,所以向量分为3个块:
1 2 3| 4 5 1|2 3 8
现在我遇到了问题,第二个1 2范围被分成不同的分区.这将导致许多无效结果,因为有些x核心最终会将搜索结果拆分为y不同的分区.我想我会search chunks -> merge y chunks into y/2 chunks -> search ->在递归样式搜索中做某种事情,但这看起来效率很低,这个算法的重点是提高效率.我也许会过度思考这种折磨
tl; dr,有没有办法以find_end我不想的方式并行化?
我试图使用R中的并行包向四个不同的处理器发送四个不同的函数调用,但我真的迷失了如何分配不同的内核来做不同的工作.我已经阅读了R中并行包,doParallel,Rmpi和foreach的文档.我看过很多帖子使用mclapply来调用具有相同参数的不同函数.我想用不同的参数调用相同的函数.
这是我想要完成的伪代码:
BEGIN parallel (core)
if(core == 1)
foo(5, 4, 1/2, 3, "a")
if(core == 2)
foo(5, 3, 1/3, 1, "b")
if(core == 3)
foo(5, 4, 1/4, 1, "c")
if(core == 4)
foo(5, 2, 1/5, 0, "d")
END parallel
Run Code Online (Sandbox Code Playgroud)
这似乎是并行计算的完美应用,因为这四个独立的函数调用可以独立地解决我正在处理的问题.我不知道如何在R中这样做.