我想知道两者之间的区别
1. labs
2. workers
3. cores
4. processes
Run Code Online (Sandbox Code Playgroud)
它只是语义还是它们都不同?
我正在尝试使用MPI send和recv函数发送std:vector但我没有到达哪里.我得到的错误就像
Fatal error in MPI_Recv: Invalid buffer pointer, error stack:
MPI_Recv(186): MPI_Recv(buf=(nil), count=2, MPI_INT, src=0, tag=0, MPI_COMM_WORLD, status=0x7fff9e5e0c80) failed
MPI_Recv(124): Null buffer pointer
Run Code Online (Sandbox Code Playgroud)
我尝试了多种组合
A)像用于发送数组的那些..
std::vector<uint32_t> m_image_data2; // definition of m_image_data2
m_image_data2.push_back(1);
m_image_data2.push_back(2);
m_image_data2.push_back(3);
m_image_data2.push_back(4);
m_image_data2.push_back(5);
MPI_Send( &m_image_data2[0], 2, MPI_INT, 1, 0, MPI_COMM_WORLD);
MPI_Send( &m_image_data2[2], 2, MPI_INT, 1, 0, MPI_COMM_WORLD);
MPI_Recv( &m_image_data2[0], 2, MPI_INT, 0, 0, MPI_COMM_WORLD, &status );
Run Code Online (Sandbox Code Playgroud)
B)没有[]
MPI_Send( &m_image_data2, 2, MPI_INT, 1, 0, MPI_COMM_WORLD);
MPI_Send( &m_image_data2 + 2, 2, MPI_INT, 1, 0, MPI_COMM_WORLD);
MPI_Recv( &m_image_data2, 2, …Run Code Online (Sandbox Code Playgroud) 我试图加快从h5py数据集文件中读取块(将它们加载到RAM内存中)的过程.现在我尝试通过多处理库来做到这一点.
pool = mp.Pool(NUM_PROCESSES)
gen = pool.imap(loader, indices)
Run Code Online (Sandbox Code Playgroud)
加载器功能是这样的:
def loader(indices):
with h5py.File("location", 'r') as dataset:
x = dataset["name"][indices]
Run Code Online (Sandbox Code Playgroud)
这实际上有时有效(意味着预期的加载时间除以进程数并因此并行化).但是,大部分时间它没有,加载时间只是保持与按顺序加载数据时一样高.有什么办法可以解决这个问题吗?我知道h5py支持通过mpi4py进行并行读/写,但我只想知道这对于只读也是绝对必要的.
受到statar包中的实验fuzzy_join函数的启发,我自己编写了一个函数,它结合了精确和模糊(通过字符串距离)匹配.我必须做的合并工作非常大(导致多个字符串距离矩阵,小于10亿个单元格),我的印象是函数编写效率不高(关于内存使用情况)和并行化以奇怪的方式实现(字符串距离矩阵的计算,如果存在多个模糊变量,而不是字符串距离本身的计算并行化).至于功能,想法是在可能的情况下匹配精确变量(以保持矩阵更小),然后在这个精确匹配的组内进行模糊匹配.我实际上认为这个功能是不言自明的.我在这里发布它是因为我希望得到一些反馈来改进它,因为我想我并不是唯一一个尝试在R中做类似事情的人(虽然我承认Python,SQL和类似的东西可能在这种情况下要更有效率.但是必须坚持一个人感觉最舒服的事情,并且使用相同的语言进行数据清理和准备在再现性方面是很好的) fuzzy_joinfuzzy_join
merge.fuzzy = function(a,b,.exact,.fuzzy,.weights,.method,.ncores) {
require(stringdist)
require(matrixStats)
require(parallel)
if (length(.fuzzy)!=length(.weights)) {
stop(paste0("fuzzy and weigths must have the same length"))
}
if (!any(class(a)=="data.table")) {
stop(paste0("'a' must be of class data.table"))
}
if (!any(class(b)=="data.table")) {
stop(paste0("'b' must be of class data.table"))
}
#convert everything to lower
a[,c(.fuzzy):=lapply(.SD,tolower),.SDcols=.fuzzy]
b[,c(.fuzzy):=lapply(.SD,tolower),.SDcols=.fuzzy]
a[,c(.exact):=lapply(.SD,tolower),.SDcols=.exact]
b[,c(.exact):=lapply(.SD,tolower),.SDcols=.exact]
#create ids
a[,"id.a":=as.numeric(.I),by=c(.exact,.fuzzy)]
b[,"id.b":=as.numeric(.I),by=c(.exact,.fuzzy)]
c <- unique(rbind(a[,.exact,with=FALSE],b[,.exact,with=FALSE]))
c[,"exa.id":=.GRP,by=.exact]
a <- merge(a,c,by=.exact,all=FALSE)
b <- merge(b,c,by=.exact,all=FALSE)
##############
stringdi <- function(a,b,.weights,.by,.method,.ncores) {
sdm <- list()
if (is.null(.weights)) {.weights <- …Run Code Online (Sandbox Code Playgroud) parallel-processing r fuzzy-comparison data.table stringdist
在我们的Web应用程序中,需要从数据库中的各种表中获取数据.今天,您可能会发现5个或6个数据库查询是针对单个请求串行执行的.这些查询都不依赖于来自另一个的数据,因此它们是并行执行的完美候选者.问题是众所周知的DbConcurrencyException,当针对相同的上下文执行多个查询时抛出该问题.
我们通常每个请求使用一个上下文,然后有一个存储库类,以便我们可以在各个项目中重用查询.然后,当处理控制器时,我们在请求结束时处理上下文.
下面是一个使用并行性的例子,但仍然存在问题!
var fileTask = new Repository().GetFile(id);
var filesTask = new Repository().GetAllFiles();
var productsTask = AllProducts();
var versionsTask = new Repository().GetVersions();
var termsTask = new Repository().GetTerms();
await Task.WhenAll(fileTask, filesTask, productsTask, versionsTask, termsTask);
Run Code Online (Sandbox Code Playgroud)
每个存储库都在内部创建自己的上下文,但就像现在一样,它们没有被处理掉.那是个问题.我知道我可以调用Dispose我创建的每个存储库,但这会使代码快速混乱.我可以为每个使用自己的上下文的查询创建一个包装器函数,但这感觉很麻烦,并不是解决问题的长期解决方案.
解决这个问题的最佳方法是什么?我希望客户端/消费者不必担心在并行执行多个查询的情况下处理每个存储库/上下文.
我现在唯一的想法是遵循类似于工厂模式的方法,除了我的工厂将跟踪它创建的所有对象.一旦我知道我的查询完成并且工厂可以在内部处理每个存储库/上下文,我就可以处理工厂.
我很惊讶地看到关于并行性和实体框架的这么少的讨论,所以希望来自社区的更多想法将会出现.
编辑
以下是我们的存储库的简单示例:
public class Repository : IDisposable {
public Repository() {
this.context = new Context();
this.context.Configuration.LazyLoadingEnabled = false;
}
public async Task<File> GetFile(int id) {
return await this.context.Files.FirstOrDefaultAsync(f => f.Id == id);
}
private bool disposed = false;
protected …Run Code Online (Sandbox Code Playgroud) 如果我有以下内容:
$ printf '%s\n' "${fa[@]}"
1 2 3
4 5 6
7 8 9
Run Code Online (Sandbox Code Playgroud)
其中每一行都是一个新的数组元素.我希望能够通过空格分隔符拆分元素,并将结果用作3个单独的参数并输入xargs.
例如,第一个元素是:
1 2 3
Run Code Online (Sandbox Code Playgroud)
在哪里使用我要传递的xargs 1,2并3进入一个简单的echo命令,例如:
$ echo $0
1
4
7
$ echo $1
2
5
8
$ echo $2
3
9
6
Run Code Online (Sandbox Code Playgroud)
所以我一直在尝试以下列方式:
printf '%s\n' "${fa[@]}" | cut -d' ' -f1,2,3 | xargs -d' ' -n 3 bash -c 'echo $0'
Run Code Online (Sandbox Code Playgroud)
这使:
1
2
3 4
5
6 7
8
9 10
Run Code Online (Sandbox Code Playgroud)
除了奇怪的行排序 - 尝试xargs -d' ' …
我在我的计算机上安装了Jenkins,它配置为只将主服务器作为节点(没有其他节点),执行次数为5.我创建了一个名为"myJob"的作业,我想在主服务器上运行2次同时(意思是如果我运行Builds 90和91,我不想得到"pending-Build#90已经在进行中"的消息).我还安装了Throttle Concurrent Builds插件,它允许这个作业同时运行多次.
我仍然收到"待定"消息.
谁能告诉我如何实现它?
我正在运行贝叶斯MCMC概率模型,我正在尝试并行实现它.在比较并行和串行时,我的机器性能令人困惑.我没有很多做并行处理的经验,所以有可能我做得不对.
我正在使用probit模型MCMCprobit的MCMCpack包,以及我parLapply在parallel包中使用的并行处理.
这是我的串行运行代码,结果来自system.time:
system.time(serial<-MCMCprobit(formula=econ_model,data=mydata,mcmc=10000,burnin=100))
user system elapsed
657.36 73.69 737.82
Run Code Online (Sandbox Code Playgroud)
这是我的并行运行代码:
#Setting up the functions for parLapply:
probit_modeling <- function(...) {
args <- list(...)
library(MCMCpack)
MCMCprobit(formula=args$model, data=args$data, burnin=args$burnin, mcmc=args$mcmc, thin=1)
}
probit_Parallel <- function(mc, model, data,burnin,mcmc) {
cl <- makeCluster(mc)
## To make this reproducible:
clusterSetRNGStream(cl, 123)
library(MCMCpack) # needed for c() method on master
probit.res <- do.call(c, parLapply(cl, seq_len(mc), probit_modeling, model=model, data=data,
mcmc=mcmc,burnin=burnin))
stopCluster(cl)
return(probit.res)
}
system.time(test<-probit_Parallel(model=econ_model,data=mydata,mcmc=10000,burnin=100,mc=2))
Run Code Online (Sandbox Code Playgroud)
结果来自system.time …
我有一个四核CPU.我创建了4个线程并运行了一个cpu密集型循环,它比在一个线程中以程序方式运行所需的时间长4倍.
我创建了两个要比较的项目,一个是线程,另一个没有.我将展示代码和运行时间.请注意没有线程的项目看起来很奇怪的原因是我想复制内存开销,因为我不确定它会影响运行时间.所以,这是没有线程的代码:
class TimeTest implements Runnable {
private Thread t;
private String name;
TimeTest(String name) {
this.name = name;
System.out.println("Creating class " + name);
}
public void run() {
System.out.println("Running class " + name);
int value = 100000000;
// try {
while (--value > 0) {
Math.random();
// Thread.sleep(1);
// System.out.println("Class " + name + " " + value);
}
// } catch (InterruptedException e) {
// System.out.println("Interrupted " + name);
// }
System.out.println("Class " + name + " exiting..."); …Run Code Online (Sandbox Code Playgroud) 在这里,我使用Javaparallel流来迭代List并使用每个列表元素作为输入调用REST调用.我需要将REST调用的所有结果添加到我正在使用的集合中ArrayList.下面给出的代码工作正常,只是ArrayList的非线程安全性会导致错误的结果,并且添加所需的同步会导致争用,从而破坏并行性的好处.
有人可以建议我在我的案例中使用并行流的正确方法.
public void myMethod() {
List<List<String>> partitions = getInputData();
final List<String> allResult = new ArrayList<String>();
partitions.parallelStream().forEach(serverList -> callRestAPI(serverList, allResult);
}
private void callRestAPI(List<String> serverList, List<String> allResult) {
List<String> result = //Do a REST call.
allResult.addAll(result);
}
Run Code Online (Sandbox Code Playgroud)