标签: parallel-processing

并行进程覆盖进度条 (tqdm)

我正在用 Python 3.7 编写一个脚本,该脚本使用multiprocessing.Process(每个内核一个任务)启动多个并行任务。为了跟踪每个进程的进度,我使用了tqdm实现进度条的库。我的代码如下所示:

with tqdm(total=iterator_size) as progress_bar:
     for row in tqdm(batch):
         process_batch(batch)
         progress_bar.update(1)
Run Code Online (Sandbox Code Playgroud)

进度条确实会相应地更新,但由于多个进程运行上面的代码,每个进程都会覆盖控制台上的进度条,如下面的屏幕截图所示。

在此处输入图片说明

完成后,控制台正确显示完成的进度条:

在此处输入图片说明

我的目标是让进度条更新而不会相互覆盖。有没有办法实现这一目标?

一个可能的解决方案是只在需要最长的进程上显示进度条(我事先知道哪个是),但最好的情况是根据第二张图片为每个进程更新一个。

所有解决方案在线地址multiprocess.Pool,但我不打算改变我的架构,因为我可以充分利用multiprocess.Process.

parallel-processing concurrency python-3.x tqdm

5
推荐指数
1
解决办法
2319
查看次数

如何同时加入一个 multiprocessing.Process() 列表?

给定一个list()正在运行multiprocessing.ProcessProcess.join实例,我如何加入所有实例并在没有超时和循环的情况下在一个退出时立即返回?

例子

from multiprocessing import Process
from random import randint
from time import sleep
def run():
    sleep(randint(0,5))
running = [ Process(target=run) for i in range(10) ]

for p in running:
    p.start()
Run Code Online (Sandbox Code Playgroud)

我怎样才能阻止,直到至少一个Processp退出?

我不想做的是:

exit = False
while not exit:
    for p in running:
        p.join(0)
        if p.exitcode is not None:
            exit = True
            break
Run Code Online (Sandbox Code Playgroud)

python parallel-processing multiprocessing python-3.x python-multiprocessing

5
推荐指数
1
解决办法
1820
查看次数

如何配置批处理脚本以使用 future.batchtools (SLURM) 并行化 R 脚本

我试图使用 future.batchtools 包在 SLURM HPC 上并行化 R 文件。虽然脚本在多个节点上执行,但它只使用 1 个 CPU 而不是可用的 12 个。

到目前为止,我尝试了不同的配置(参见附上的代码),但都没有达到预期的结果。我的 bash 文件配置如下:

#!/bin/bash
#SBATCH --nodes=2
#SBATCH --cpus-per-task=12

R CMD BATCH test.R output
Run Code Online (Sandbox Code Playgroud)

在 R 中,我使用 foreach 循环:

# First level = cluster
# Second level = multiprocess 
# https://cran.r-project.org/web/packages/future.batchtools/vignettes/future.batchtools.html
plan(list(batchtools_slurm, multiprocess))

# Parallel for loop
result <- foreach(i in 100) %dopar% {
       Sys.sleep(100)
return(i) 
}
Run Code Online (Sandbox Code Playgroud)

如果有人可以指导我如何为多个节点和多个核心配置代码,我将不胜感激。

parallel-processing hpc future r slurm

5
推荐指数
1
解决办法
737
查看次数

为什么 scikit-learn 切换到 SequentialBackend?

我尝试在具有 16 个可用 CPU 的机器上运行以下代码:

def tokenizer(text):
    return text.split()

param_grid = [{'vect__stop_words': [None, stop],
               'vect__binary': [True, False]}]

bow = CountVectorizer(ngram_range=(1,1), tokenizer=tokenizer)  
multinb_bow = Pipeline([('vect', bow), ('clf', MultinomialNB())])

gs_multinb_bow = GridSearchCV(multinb_bow, param_grid, scoring='f1_macro', 
                              cv=3, verbose=1, n_jobs=-1)

gs_multinb_bow.fit(X_train, y_train)
Run Code Online (Sandbox Code Playgroud)

我设置n_jobs-1,但scikit-learn切换到SequentialBackend,即使我添加了上下文管理器with parallel_backend('loky'):并且脚本仍然仅使用 1 个并发工作器运行。

Fitting 3 folds for each of 4 candidates, totalling 12 fits
[Parallel(n_jobs=-1)]: Using backend SequentialBackend with 1 concurrent workers.
Run Code Online (Sandbox Code Playgroud)

如果我为 指定不同的值,同样的结果仍然存在n_jobs

为什么会这样?我最近在类似的任务上运行了似乎是相同的代码,并且网格搜索在多个 CPU 上并行工作,如n_jobs使用LokyBackend …

python parallel-processing machine-learning scikit-learn grid-search

5
推荐指数
0
解决办法
685
查看次数

关闭ssh连接后在后台通过java执行bash脚本

我正在使用 java 在远程 linux 机器上执行一个简单的 bash 脚本。

名为“shortoracle.bash”的 bash 脚本有这个脚本:

#!/bin/sh
runsql() {
   i="$1"
   end=$((SECONDS+360))
   SECONDS=0
   while (( SECONDS < end )); do
   echo "INSERT into table_$i (col1) values (CURRENT_TIMESTAMP);" | sqlplus username/password
   sleep 1
   done
}


for i in $(seq 1 10); do
 echo "DROP TABLE table_$i;" | sqlplus username/password
 echo "CREATE TABLE table_$i (col1 TIMESTAMP WITH TIME ZONE);" | sqlplus username/password
 runsql $i &
done
wait
Run Code Online (Sandbox Code Playgroud)

简单来说:创建 10 个并行连接,执行 360 秒的查询。

从我的 Java 程序中,我执行以下命令:

sshconnection.execute("nohup su - oracle …
Run Code Online (Sandbox Code Playgroud)

java linux parallel-processing ssh bash

5
推荐指数
1
解决办法
182
查看次数

如何使用 Neo4j Cypher APOC 并行处理

我们有很多 (n1:EffortUser)-[r1:EFFORT]->(n2:EffortObject) 需要按天和周计算,即有多少 EffortObject:Email 做了 EffortUser SENT。如果您有很多用户和电子邮件可能需要相当长的时间,那么我们希望并行化此查询。

现在我们正在使用:

match(n1:EffortUser)-[r1:EFFORT]-(n2:EffortObject) 
where r1.Effort = 'yes' and r1.TimeEvent>='2017-01-01' and r1.TimeEvent<='2017-12-31'
return distinct n1.Name as User, date(datetime(r1.TimeEvent)) as date, count(distinct r1.IdUnique) as count
order by user, date
Run Code Online (Sandbox Code Playgroud)

似乎有几个选项可以并行化/优化它,但所有的文档都相当糟糕。

我做了一些研究,发现了以下 APOC 函数,但我可能无法让它们工作(而且我在 Stackoverflow 上也找不到太多东西)。以下哪个选项是最好的,包括 使用上述示例代码的示例?这让我发疯。我们有 4 个内核和 32 GB 的内存,所以这应该运行得非常快,但我就是无法让它工作。

https://neo4j.com/docs/labs/apoc/current/cypher-execution/

CALL apoc.cypher.runMany('cypher;\nstatements;',{params},{config})
Run Code Online (Sandbox Code Playgroud)

运行每个分号分隔的语句并返回摘要 - 当前没有模式操作

CALL apoc.cypher.mapParallel(fragment, params, list-to-parallelize) yield value
Run Code Online (Sandbox Code Playgroud)

并行批量执行片段,并将列表段分配给 _

https://neo4j.com/docs/labs/apoc/current/cypher-execution/running-cypher/

apoc.cypher.mapParallel(fragment :: STRING?, params :: MAP?, list :: LIST? OF ANY?) :: (value :: MAP?)
Run Code Online (Sandbox Code Playgroud)

apoc.cypher.mapParallel(fragment, params, list-to-parallelize) yield value …

parallel-processing neo4j cypher neo4j-apoc

5
推荐指数
0
解决办法
1047
查看次数

达斯克VS急流。急流提供哪些 dask 没有?

我想了解 dask 和 Rapids 之间的区别是什么,rapids 提供哪些 dask 没有的好处。

Rapids 内部是否使用 dask 代码?如果是这样,那么为什么我们有 dask,因为即使 dask 也可以与 GPU 交互。

parallel-processing gpu machine-learning dask rapids

5
推荐指数
2
解决办法
1097
查看次数

为什么parallelStream 使用的是ForkJoinPool,而不是普通的线程池?

参考Java 的 Fork/Join vs ExecutorService - 何时使用哪个?,一个传统的线程池通常用于处理很多独立的请求;和 aForkJoinPool用于处理连贯/递归任务,其中一个任务可能会产生另一个子任务并稍后加入。

那么,为什么默认parallelStream使用Java-8ForkJoinPool而不是传统的执行器呢?

在很多情况下,我们forEach()stream()orparallelStream()之后使用,然后提交一个功能接口作为参数。在我看来,这些任务是独立的,不是吗?

parallel-processing concurrency threadpool forkjoinpool java-stream

5
推荐指数
1
解决办法
520
查看次数

SSIS 包全表加载缓慢

我们有一个 SSIS 包,显然被开发团队称为“慢”。由于他们没有使用 SSIS ETL 的人,因此作为 DBA,我尝试深入研究。以下是我找到的信息:SQL Server 已从 2014 版本升级到 2017,因此它具有两个版本的 SSIS。

  1. 他们将大小为 200 GB 的 SQL Server 表加载到 SSIS 中,然后使用命令行压缩功能将数据压缩到平面文件中。
  2. 数据流任务很简单select * from view——视图只是包含没有其他花哨连接的表。
  3. 在进行故障排除时,我发现在 SQL Server 上几乎没有任何负载,可能是因为 select 命令在单线程中运行而不使用 SQL Server 内核。
  4. 当我运行相同的 select * 命令(仅 5 秒,因为它是 200 GB 表)时,即使我的命令也是单线程的。
  5. 该包具有 SQL 作业显示的配置文件(这是包的运行方式)以及一些连接设置。
  6. 在 BIDS 中打开包显示 defaultBufferMaxRows 仅为 10000(可能是默认值)(因为配置文件或任何变量没有客户值,我猜这也是包正在使用的)。

SQL 和 SSIS 都在同一台服务器上。SQL 已分配最大内存,为 SSIS 和 OS 留下大约 100 GB。

请分享有关如何强制 SQL Server 使用多个线程运行此 select 命令以便整个表更快地进入 SSIS 缓冲池的任何想法。

编辑:我知道bcp可以比任何进程更快地读取数据并将其保存到平面文件,但此时对 SSIS 包的更改必须保持最少,并探索可以合并到 …

sql-server parallel-processing performance ssis etl

5
推荐指数
1
解决办法
1607
查看次数

在python中为多个参数并行运行单个函数的最快方法

假设我只有一个函数processing。我想为多个参数并行运行相同的函数多次,而不是一个接一个地依次运行。

def processing(image_location):
    
    image = rasterio.open(image_location)
    ...
    ...
    return(result)

#calling function serially one after the other with different parameters and saving the results to a variable.
results1 = processing(r'/home/test/image_1.tif')
results2 = processing(r'/home/test/image_2.tif')
results3 = processing(r'/home/test/image_3.tif')

Run Code Online (Sandbox Code Playgroud)

例如,如果我运行delineation(r'/home/test/image_1.tif')然后delineation(r'/home/test/image_2.tif'),然后delineation(r'/home/test/image_3.tif'),如图上面的代码,这将顺序运行一前一后,并且如果需要5分钟一个函数来运行然后运行这三个将采取5X3 = 15分钟。因此,我想知道我是否可以并行/尴尬地并行运行这三个,以便对所有三个不同参数执行该函数只需要 5 分钟。

帮助我以最快的方式完成这项工作。该脚本应该能够利用默认情况下可用的所有资源/CPU/ram 来完成此任务。

python parallel-processing function multiprocessing embarrassingly-parallel

5
推荐指数
1
解决办法
694
查看次数