我对MapReduce非常陌生,我完成了一个Hadoop字数统计示例.
在该示例中,它生成单词计数的未排序文件(具有键值对).那么是否可以通过将另一个MapReduce任务与前一个任务相结合来按字出现次数对其进行排序?
我想在Haskell中编写一个尽可能高效的并行映射函数.我最初的尝试,似乎是目前最好的,只是写,
pmap :: (a -> b) -> [a] -> [b]
pmap f = runEval . parList rseq . map f
Run Code Online (Sandbox Code Playgroud)
但是,我没有看到完美的CPU划分.如果这可能与火花的数量有关,我可以编写一个将列表划分为#ppus段的pmap ,因此创建了最少的火花吗?我试过以下,但是性能(和火花的数量)要差得多,
pmap :: (a -> b) -> [a] -> [b]
pmap f xs = concat $ runEval $ parList rseq $ map (map f) (chunk xs) where
-- the (len / 4) argument represents the size of the sublists
chunk xs = chunk' ((length xs) `div` 4) xs
chunk' n xs | length xs <= n = …Run Code Online (Sandbox Code Playgroud) 我刚刚在CUDA开始了一个小项目.
我需要知道以下内容:是否可以在不使用/购买Microsoft Visual Studio的情况下编译CUDA代码?使用Nvcc.exe我收到错误" 无法在路径中找到编译器cl.exe ".
我曾尝试为NetBeans 安装CUDA 插件,但它不起作用.(使用当前版本的NetBeans)
平台:Windows 7
提前致谢.
我正在研究For循环中的Parallelism Break.
我希望这段代码:
Parallel.For(0, 10, (i,state) =>
{
Console.WriteLine(i); if (i == 5) state.Break();
}
Run Code Online (Sandbox Code Playgroud)
在得到最 6号(0..6).他不仅没有这样做,而且结果长度不同:
02351486
013542
0135642
Run Code Online (Sandbox Code Playgroud)
很烦人.(这里的地狱是Break(){5之后}这里??)
所以我看了msdn
Break可以用于与循环通信,在当前迭代之后不需要运行其他迭代.如果从for循环的第100次迭代调用Break从0到1000并行迭代,则仍应运行小于100的所有迭代,但不需要从101到1000的迭代.
Quesion #1 :
哪个迭代?整个迭代计数器?还是每个帖子?我很确定这是每个帖子.请批准.
Question #2 :
让我们假设我们使用并行+范围分区(由于元素之间没有cpu成本变化),因此它在线程之间划分数据.因此,如果我们有4个核心(并且它们之间有完美的划分):
core #1 got 0..250
core #2 got 251..500
core #3 got 501..750
core #4 got 751..1000
Run Code Online (Sandbox Code Playgroud)
所以线程core #1会在value=100某个时候遇到并且会中断.这将是他的迭代号 100.但是线程core #4得到了更多的量子,他900现在正在进行中.他超越了他的100'th迭代.他没有指数少于100被停止!! - 所以他会向他们展示所有.
我对吗 ?这就是我在我的例子中获得超过5个元素的原因吗?
Question #3 :
我真的打破了什么时候(i == …
有人可以提供一些关于如何并行化PyMC MCMC代码的一般性说明.我试图LASSO按照这里给出的例子运行回归.我在某处读到默认情况下并行采样,但是我是否还需要使用类似的功能Parallel Python来使其工作?
这是一些我希望能够在我的机器上并行化的参考代码.
x1 = norm.rvs(0, 1, size=n)
x2 = -x1 + norm.rvs(0, 10**-3, size=n)
x3 = norm.rvs(0, 1, size=n)
X = np.column_stack([x1, x2, x3])
y = 10 * x1 + 10 * x2 + 0.1 * x3
beta1_lasso = pymc.Laplace('beta1', mu=0, tau=1.0 / b)
beta2_lasso = pymc.Laplace('beta2', mu=0, tau=1.0 / b)
beta3_lasso = pymc.Laplace('beta3', mu=0, tau=1.0 / b)
@pymc.deterministic
def y_hat_lasso(beta1=beta1_lasso, beta2=beta2_lasso, beta3=beta3_lasso, x1=x1, x2=x2, x3=x3):
return beta1 * x1 …Run Code Online (Sandbox Code Playgroud) 我有一个.sln包含大量项目(~50)的大型Visual Studio 2010解决方案文件().每个项目都包含粗略20 .cpp和.h文件.在快速的Intel i7计算机上构建整个项目大约需要2个小时,所有依赖项都在快速SSD上本地缓存.
我试图复制我在这个主题上发现的现有实验,需要澄清如何:
MSBuild.exe.首先,/m和/maxcpucount选项是一样的吗?我已经阅读了一篇关于该/m选项的文章,该文章似乎并行构建了更多项目.这似乎与我通过GUI中的以下操作设置的选项相同:
Maximum number of parallel project builds还有另一种选择,/MP我可以通过以下方式访问GUI:
Multi-processor Compilation如果我错了,请纠正我,但这似乎表明该项目将.cpp并行构建多个文件,但没有选项指定多少文件,除非我必须在文本框中手动设置它(即:/MP 4或者/MP4.
似乎为了/MP工作,我需要禁用最小重建(即:)Enabled Minimal Rebuild: No (/GM-),它也不适用于预编译的头文件.我已经解决了预编译头文件问题,方法是将预编译头文件作为一个专用项目,该解决方案首先构建在所有其他项目之前.
问题:我如何在解决方案中为并行项目构建实现这些选项,并在项目中构建并行文件(.cpp),通过命令行使用MSBuild,并确认它们按预期工作?我的上述假设是否也正确?我的最终目标是测试构建项目的最多使用的核心总数,或者如果我的构建过程是I/O绑定而不是CPU绑定,甚至可能重载它.我可以在Linux中通过以下方式轻松完成此操作:
NUMCPUS=`grep …Run Code Online (Sandbox Code Playgroud) msbuild parallel-processing visual-studio-2010 multiprocessing visual-studio
我有一个多处理工作,我正在排队只读numpy数组,作为生产者消费者管道的一部分.
目前他们正在被腌制,因为这是默认行为multiprocessing.Queue会降低性能.
是否有任何pythonic方法将引用传递给共享内存而不是pickle数组?
不幸的是,在消费者启动之后会生成数组,并且没有简单的方法.(所以全局变量方法会很难看......).
[注意,在下面的代码中,我们不期望并行计算h(x0)和h(x1).相反,我们看到并行计算的h(x0)和g(h(x1))(就像CPU中的流水线一样).
from multiprocessing import Process, Queue
import numpy as np
class __EndToken(object):
pass
def parrallel_pipeline(buffer_size=50):
def parrallel_pipeline_with_args(f):
def consumer(xs, q):
for x in xs:
q.put(x)
q.put(__EndToken())
def parallel_generator(f_xs):
q = Queue(buffer_size)
consumer_process = Process(target=consumer,args=(f_xs,q,))
consumer_process.start()
while True:
x = q.get()
if isinstance(x, __EndToken):
break
yield x
def f_wrapper(xs):
return parallel_generator(f(xs))
return f_wrapper
return parrallel_pipeline_with_args
@parrallel_pipeline(3)
def f(xs):
for x in xs:
yield x + 1.0
@parrallel_pipeline(3)
def g(xs):
for x in xs:
yield x …Run Code Online (Sandbox Code Playgroud) 我正在处理一个需要并行计算以获得比经典“for 循环”更快的结果的问题。
问题是这样的:
我需要为列表对象内的数据帧中包含的 198135 个结果变量生成线性模型。我必须将模型中每个预测变量的所有 beta 和 p 值以及它们的拟合优度度量存储在数据框中。
我编写了一个功能性“for 循环”,可以正确完成该任务,但完成它需要超过 35 个小时。我知道 R 使用了我的 8 核 CPU 的不到 20%,但我想全部使用。问题是我不知道如何将 for 循环转换为 foreach 循环以利用并行计算。
这是我的问题的一些较小规模的示例代码:
library(tidyverse)
library(broom)
## Example data
outcome_list <- list(as.data.frame(cbind(rnorm(32), dataframe_id = c(1))),
as.data.frame(cbind(rnorm(32), dataframe_id = c(2))),
as.data.frame(cbind(rnorm(32), dataframe_id = c(3)))) ## This represents my list of 198135 dataframes
mtcars <- mtcars #I will use the explanatory variables from here
## Below this line is my current solution with a for loop that works fine
x …Run Code Online (Sandbox Code Playgroud) 我有这个:
Stream<CompletableFuture<List<Item>>>
Run Code Online (Sandbox Code Playgroud)
我怎样才能将它转换为
Stream<CompletableFuture<Item>>
Run Code Online (Sandbox Code Playgroud)
其中:第二个流由第一个流中每个列表内的每个项目组成。
我研究了一下thenCompose,但这解决了一个完全不同的问题,也称为“扁平化”。
如何以流方式高效地完成此操作,而不阻塞或过早消耗不必要的流项目?
这是迄今为止我最好的尝试:
ExecutorService pool = Executors.newFixedThreadPool(PARALLELISM);
Stream<CompletableFuture<List<IncomingItem>>> reload = ... ;
@SuppressWarnings("unchecked")
CompletableFuture<List<IncomingItem>> allFutures[] = reload.toArray(CompletableFuture[]::new);
CompletionService<List<IncomingItem>> queue = new ExecutorCompletionService<>(pool);
for(CompletableFuture<List<IncomingItem>> item: allFutures) {
queue.submit(item::get);
}
List<IncomingItem> THE_END = new ArrayList<IncomingItem>();
CompletableFuture<List<IncomingItem>> ender = CompletableFuture.allOf(allFutures).thenApply(whatever -> {
queue.submit(() -> THE_END);
return THE_END;
});
queue.submit(() -> ender.get());
Iterable<List<IncomingItem>> iter = () -> new Iterator<List<IncomingItem>>() {
boolean checkNext = true;
List<IncomingItem> next = null;
@Override
public boolean hasNext() {
if(checkNext) {
try { …Run Code Online (Sandbox Code Playgroud) java parallel-processing concurrency java-stream completable-future
正如我在 中看到的,一次pg_stat_activiry只有一个命令执行。COPY正如我在专栏中看到的那样,其他查询处于锁定状态wait_event_type。
如何COPY mytable FROM STDIN在不锁定表的情况下并行运行多个?
附:mytable是TimescaleDB 2.5.0的超表。
UPD
CREATE TABLE "public"."mytable" (
"q_time" timestamp,
"symbol_id" int,
"o" decimal(24,12),
"c" decimal(24,12),
"h" decimal(24,12),
"l" decimal(24,12),
"v" bigint,
CONSTRAINT mytable_ts_pkey PRIMARY KEY (symbol_id, "q_time")
);
SELECT create_hypertable('mytable', 'q_time', 'symbol_id', 1,
create_default_indexes => false,
chunk_time_interval => '7 days'::interval);
Run Code Online (Sandbox Code Playgroud)
UPD2
我并行运行下一个命令:
out, err := exec.Command("bash", "-c", "cat file01.gz | gunzip | psql -d db -U user -c "\copy mytable from stdin HEADER DELIMITER …Run Code Online (Sandbox Code Playgroud) python ×2
.net ×1
.net-4.0 ×1
c# ×1
concurrency ×1
cuda ×1
foreach ×1
hadoop ×1
haskell ×1
java ×1
java-stream ×1
mapreduce ×1
msbuild ×1
multicore ×1
numpy ×1
performance ×1
postgresql ×1
pymc ×1
pymc3 ×1
r ×1
timescaledb ×1
windows ×1
word-count ×1