在从EPFL并行编程过程中,4个抽象数据并行提到:Iterator,Builder,Combiner,和Splitter.
我很熟悉Iterator,但从未使用过其他三个.我所看到的其他特征Builder,Combiner和Splitter下scala.collection包.不过,我知道如何在现实世界的发展使用它们,特别是如何在合作与其他收藏品一样使用它们List,Array,ParArray等任何人都可以请给我一些指导和实例?
谢谢!
我正在使用两个data.frames列表,目前运行类似于此的东西(我正在做的简化版本):
df1 <- data.frame("a","a1","L","R","b","c",1,2,3,4)
df2 <- data.frame("a","a1","L","R","b","c",4,4,4,4,4,44)
df3 <- data.frame(7,7,7,7)
df4 <- data.frame(5,5,5,5,9,9)
L1 <- list(df1,df2)
L2 <- list(df3,df4)
myfun <- function(x,y) {
difa = rowSums(abs(x[c(T,F)] - x[c(F,T)]))
difb=sum(abs(as.numeric(y[-c(1:6)])[c(T,F)] - as.numeric(y[-c(1:6)])[c(F,T)]))
diff <- difa + difb
return(diff)
}
output1 <- mapply(myfun, x = L2, y = L1)
Run Code Online (Sandbox Code Playgroud)
每个列表中的数据帧数相同,一个列表中的每个数据帧对应另一个列表中的数据帧.一个列表中的数据帧包含单个行,而第二个列表中的其他数据帧包含动态行数; 因此使用sum和rowSums.数字列的数量也是动态的,但在相应的数据帧之间始终相同.
我希望在处理每个列表1-10万个数据帧时使用并行处理来加速计算.我尝试了以下方法:
library(parallel)
if(detectCores() > 1) {no_cores <- detectCores() - 1}
if(.Platform$OS.type == "unix") {ptype <- "FORK"}
cl <- makeCluster(no_cores, type = ptype)
clusterMap(cl, myfun, x = L2, y = L1)
stopCluster(cl) …Run Code Online (Sandbox Code Playgroud) 在LINQ查询中,我可以正确地(如:编译器不会抱怨)调用.AsParallel(),如下所示:
(from l in list.AsParallel() where <some_clause> select l).ToList();
Run Code Online (Sandbox Code Playgroud)
或者像这样:
(from l in list where <some_clause> select l).AsParallel().ToList();
Run Code Online (Sandbox Code Playgroud)
究竟有什么区别?
从官方文档来看,我几乎总是看到第一种方法,所以我认为这是要走的路.
今天,我试图自己运行一些基准测试,结果令人惊讶.这是我运行的代码:
var list = new List<int>();
var rand = new Random();
for (int i = 0; i < 100000; i++)
list.Add(rand.Next());
var treshold= 1497234;
var sw = new Stopwatch();
sw.Restart();
var result = (from l in list.AsParallel() where l > treshold select l).ToList();
sw.Stop();
Console.WriteLine($"call .AsParallel() before: {sw.ElapsedMilliseconds}");
sw.Restart();
result = (from l in list where …Run Code Online (Sandbox Code Playgroud) 我有一个AWS实例.我想运行一堆任务,一些内存和CPU密集.理想情况下,我想计算每项任务的时间信息.如果我连续运行它们,它会计算准确的计时信息,但速度很慢.如果我并行运行它们,整个事情就会更快,但是单个任务的速度会更慢,正如壁时间和线程CPU时间所报告的那样.
随着线程数量增加到CPU数量,这种减速会增加
粗略检查ghc-events-analyze并+RTS -s暗示减速的来源(不出所料)GC暂停.使用RTS选项显示+RTS -qg -qb -qa -A256m(禁用并行GC,禁用负载平衡GC,禁用线程迁移以及增加GC分配区域)可以改善这一点,但并不能完全消除它.
我正在使用线程运行forkIO,但除了打印进度信息之外,线程是独立且纯粹的.我正在使用parallel-io来管理正在运行的线程的数量,但是当我简单地尝试一种更常规的方法来获得一个固定的线程池和一个任务队列时,我仍然遇到了这个问题.
有关如何调试的任何建议?
编辑:
@jberryman问了一个例子.每个任务看起来像下面的代码
computation params = do
!x <- force params
print $ "Starting computation on " ++ show params
t1 <- getCPUTime
!y <- fmap force $ do $
...some work with x ...
t2 <- getCPUTime
print $ "Finished computation on " ++ show params
return (t2 - t1, y)
Run Code Online (Sandbox Code Playgroud) parallel-processing multithreading garbage-collection haskell
据我所知,C++ 17将带有Parallelism.但是,我无法理解的是它是一种特定的硬件并行性(默认为CPU)?或者它可以扩展到具有多个计算单元的任何硬件?
换句话说,我们会看到类似于"nVidia C++标准编译器"的东西,它将编译要在GPU上执行的并行部分吗?
例如,它是OpenCL的一些标准替代品吗?
注意:当然,我不是在问"nVidia会这么做吗?".我在问C++ 17标准是否允许,以及理论上是否可行.
我现在正在处理大型数据集,某些功能可能需要数小时才能处理.我想知道如何通过进度条或数字(1,2,3,...,100)显示代码的进度.我想将结果存储为具有两列的数据框.这是一个例子.谢谢.
require(foreach)
require(doParallel)
require(Kendall)
cores=detectCores()
cl <- makeCluster(cores-1)
registerDoParallel(cl)
mydata=matrix(rnorm(8000*500),ncol = 500)
result=as.data.frame(matrix(nrow = 8000,ncol = 2))
pb <- txtProgressBar(min = 1, max = 8000, style = 3)
foreach(i=1:8000,.packages = "Kendall",.combine = rbind) %dopar%
{
abc=MannKendall(mydata[i,])
result[i,1]=abc$tau
result[i,2]=abc$sl
setTxtProgressBar(pb, i)
}
close(pb)
stopCluster(cl)
Run Code Online (Sandbox Code Playgroud)
但是,当我运行代码时,我没有看到任何进度条显示,结果不正确.有什么建议吗?谢谢.
我有一个我想要并行执行的进程,但是由于一些奇怪的错误我失败了.现在我正在考虑组合,并计算主CPU上的失败任务.但是我不知道如何为.combine编写这样的函数.
怎么写?
我知道如何编写它们,例如这个答案提供了一个例子,但它没有提供如何处理失败的任务,也没有重复在主服务器上重复任务.
我会做的事情如下:
foreach(i=1:100, .combine = function(x, y){tryCatch(?)} %dopar% {
long_process_which_fails_randomly(i)
}
Run Code Online (Sandbox Code Playgroud)
但是,如何在.combine函数中使用该任务的输入(如果可以的话)?或者我应该在内部提供%dopar%返回标志或列表来计算它?
我可以改变我的循环
for (int i = 0; i < something; i++)
Run Code Online (Sandbox Code Playgroud)
至:
Parallel.For(0, something, i =>
Run Code Online (Sandbox Code Playgroud)
但是如何用这个循环做到这一点?:
for (i = 3; i <= something / 2; i = i + 2)
Run Code Online (Sandbox Code Playgroud)
谢谢你的回答.
我们使用更重量级的控制台应用程序,使用HTTP触发器和消费服务计划测试了Azure功能的横向扩展功能.所以我们期望通过扩展来实现并行执行.我们在新的AppDomain中执行控制台应用程序,因为funcion实例在同一进程中运行.在控制台应用程序中,我们在内存数据库中执行sqlite数据库操作.
首先我们只执行一次该功能,并测量执行时间.让它成为x :)我们不断开始增加并行线程的数量.我们经历过在这些情况下1个函数app实例的执行时间是x*num_of_threads.好像函数实例已被序列化并且不是并行执行的.
谢谢你的帮助.
编辑:我的应用程序的基本源代码:
using System.Net;
using System;
public static HttpResponseMessage Run(HttpRequestMessage req, TraceWriter log, ExecutionContext context)
{
string testThreadId = req.GetQueryNameValuePairs()
.FirstOrDefault(q => string.Compare(q.Key, "id", true) == 0)
.Value;
var funcId = context.InvocationId.ToString();
var homePath = Environment.GetEnvironmentVariable("HOME");
var folderName = Path.Combine(homePath,@"site\wwwroot\JanoRunTime2");
var fileName = Path.Combine(folderName,"AzureFunctionTest.exe");
var configFile = Path.Combine(folderName,"AzureFunctionTest.exe.config");
var setup = new AppDomainSetup();
setup.ApplicationBase = folderName;
setup.ConfigurationFile = configFile;
var newDomain = AppDomain.CreateDomain("JanoTestExecutorDomain_" + funcId, null, setup );
try{
newDomain.ExecuteAssembly(fileName, new []{testThreadId, funcId});
return req.CreateResponse(HttpStatusCode.OK …Run Code Online (Sandbox Code Playgroud) 我们在整合Spark-Kafka流时遇到了性能问题.
项目设置:我们使用带有3个分区的Kafka主题,并在每个分区中生成3000条消息,并在Spark直接流式处理中进行处理.
我们面临的问题:在处理结束时,我们采用Spark直接流方法来处理相同的问题.根据以下文档.Spark应该创建与主题中的分区数量一样多的并行直接流(在本例中为3).但是在阅读时我们可以看到来自分区1的所有消息首先被处理,然后是第二个然后是第三个.任何帮助为什么它不处理并行?根据我的理解,如果它同时从所有分区并行读取,那么消息输出应该是随机的.
r ×3
c# ×2
.net ×1
azure ×1
c++ ×1
c++17 ×1
collections ×1
for-loop ×1
foreach ×1
haskell ×1
linq ×1
parallel.for ×1
plinq ×1
progress-bar ×1
scala ×1