我在7500多个对象上运行一个Parallel.For循环.在for循环中,我正在为每个对象做很多事情,特别是调用两个Web服务和两个内部方法.Web服务只是检查对象,处理并返回一个字符串,然后我将其设置为对象上的属性.两种内部方法也是如此.
我没有写任何东西到磁盘或从磁盘读取.
我还在带有标签和进度条的winforms应用程序中更新UI,以便让用户知道它在哪里.这是代码:
var task = Task.Factory.StartNew(() =>
{
Parallel.For(0, upperLimit, (i, loopState) =>
{
if (cancellationToken.IsCancellationRequested)
loopState.Stop();
lblProgressBar.Invoke(
(Action)
(() => lblProgressBar.Text = string.Format("Processing record {0} of {1}.", (progressCounter++), upperLimit)));
progByStep.Invoke(
(Action)
(() => progByStep.Value = (progressCounter - 1)));
CallSvc1(entity[i]);
Conversion1(entity[i]);
CallSvc2(entity[i]);
Conversion2(entity[i]);
});
}, cancellationToken);
Run Code Online (Sandbox Code Playgroud)
这是在Win7 32位机器上进行的.
关于为什么当增量器大约在1370左右时突然冻结的任何想法(这是1361,1365和1371)?
关于如何调试这个并看看有什么锁定的任何想法?
编辑:
以下评论的一些答案:
@BrokenGlass - 不,没有互操作.我将尝试x86编译并让你知道.
@chibacity - 因为它是在后台任务上,所以它不会冻结UI.直到它冻结的时间,进度条和标签每秒大约2点.当它冻结时,它就会停止移动.我可以验证它停止的号码是否已被处理,但不再处理.双核2.2GHz的CPU使用率在运行期间最低,每次3-4%,冻结后1-2%.
@Henk Holterman - 到达1360需要大约10-12分钟,是的,我可以验证所有这些记录是否已经处理但不是剩余的记录.
@CodeInChaos - 谢谢,我会试试!如果我拿出并行代码,代码确实有用,它只需要一天又一天.我没有尝试过限制线程数,但是会.
编辑2:
关于Web服务发生了什么的一些细节
基本上,Web服务正在发生的是它们传递一些数据并接收数据(XmlNode).然后在Conversion1进程中使用该节点,该进程又在实体上设置另一个属性,该属性被发送到CallSvc2方法,依此类推.它看起来像这样:
private void CallSvc1(Entity entity)
{
var svc = new MyWebService();
var …Run Code Online (Sandbox Code Playgroud) 我在Makefile中有这个:
run:
for x in *.bin ; do ./$$x ; done
Run Code Online (Sandbox Code Playgroud)
这样它就可以逐个启动所有可执行文件.我想做这个:
run:
for x in *.bin ; do ./$$x &; done
Run Code Online (Sandbox Code Playgroud)
这样它就会启动每个可执行文件并将其放在后台.当我放入&符号时,上面的语句出现语法错误.
我不想调用make,make &因为这将在后台运行进程但仍然是一个接一个,而我希望单个可执行文件在后台运行,这样在任何时刻我都有多个可执行文件在运行.
先感谢您.
我正在研究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 == …
我正在处理一个需要并行计算以获得比经典“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
我正在处理相当大的 Pandas DataFrame - 我的数据集类似于以下df设置:
import pandas as pd
import numpy as np
#--------------------------------------------- SIZING PARAMETERS :
R1 = 20 # .repeat( repeats = R1 )
R2 = 10 # .repeat( repeats = R2 )
R3 = 541680 # .repeat( repeats = [ R3, R4 ] )
R4 = 576720 # .repeat( repeats = [ R3, R4 ] )
T = 55920 # .tile( , T)
A1 = np.arange( 0, 2708400, 100 ) # ~ 20x re-used
A2 …Run Code Online (Sandbox Code Playgroud) 我知道关于 Julia 中多线程性能的问题已经被问过(例如这里),但它们涉及相当复杂的代码,其中可能有很多东西在起作用。
在这里,我使用 Julia v1.5.3 在多个线程上运行一个非常简单的循环,与使用例如 Chapel 运行相同的循环相比,加速似乎并没有很好地扩展。
我想知道我做错了什么,以及如何更有效地在 Julia 中运行多线程。
using BenchmarkTools
function slow(n::Int, digits::String)
total = 0.0
for i in 1:n
if !occursin(digits, string(i))
total += 1.0 / i
end
end
println("total = ", total)
end
@btime slow(Int64(1e8), "9")
Run Code Online (Sandbox Code Playgroud)
时间:8.034s
Threads.@threads4 个线程上的共享内存并行性using BenchmarkTools
using Base.Threads
function slow(n::Int, digits::String)
total = Atomic{Float64}(0)
@threads for i in 1:n
if !occursin(digits, string(i))
atomic_add!(total, 1.0 / i)
end
end
println("total = ", total)
end
@btime slow(Int64(1e8), …Run Code Online (Sandbox Code Playgroud) 正如我在 中看到的,一次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) c# ×2
.net ×1
.net-4.0 ×1
apache-spark ×1
c++ ×1
c++20 ×1
chapel ×1
concurrency ×1
dask ×1
foreach ×1
freeze ×1
iterator ×1
java ×1
java-stream ×1
julia ×1
makefile ×1
mpi ×1
pandas ×1
performance ×1
postgresql ×1
python ×1
r ×1
std-ranges ×1
timescaledb ×1