parallel.futures:如何在2个单独的线程中提交任务和处理Future?

如何解决parallel.futures:如何在2个单独的线程中提交任务和处理Future?

我正在尝试改善小型CLI程序的交互式输出,使其在目录中处理文件,并使用Rich进度条来显示任务的进度。

此刻,我分两个步骤进行操作:

  • pool.submit()所有任务
  • for future in as_completed(xxxx)等待下一个未来。

问题在于,第一步(pool.submit)可能要花一些时间(因为我正在走目录),并且UI也没有更新,即使将来已经可用。

因此,我试图提出一个将在我的池中提交的线程,而主线程将等待下一个Future并更新UI:

"""
Usage: walker.py [options] <file/directory>...

Options:
    -r --recursive                  Walk directories recursively
    -w WORKERS --workers=WORKERS    Specify the number of process pool workers [default: 4]
    -d --debug                      Enable debug output
    -h --help                       Display this message
"""
import os
import threading
import time
from concurrent.futures._base import as_completed
from concurrent.futures.process import ProcessPoolExecutor
from pathlib import Path
from random import randint
from typing import List

from docopt import docopt
from rich.console import Console
from rich.progress import BarColumn,Progress,TextColumn


def walk_filepath_list(filepath_list: List[Path],recursive: bool = False):
    for path in filepath_list:
        if path.is_dir() and not path.is_symlink():
            if recursive:
                for f in os.scandir(path):
                    yield from walk_filepath_list([Path(f)],recursive)
            else:
                yield from (Path(f) for f in os.scandir(path))
        elif path.is_file():
            yield path


def process_task(filepath):
    rand = randint(0,1)
    time.sleep(rand)


def thread_submit(pool,filepath_list,recursive,future_to_filepath):
    for filepath in walk_filepath_list(filepath_list,recursive):
        future = pool.submit(process_task,filepath)
        # update shared dict
        future_to_filepath[future] = filepath


def main(args):
    filepath_list = [Path(entry) for entry in args["<file/directory>"]]
    debug = args["--debug"]
    workers = int(args["--workers"])
    recursive = args["--recursive"]

    console = Console()

    process_bar = Progress(
        TextColumn("[bold blue]Processing...",justify="left"),BarColumn(bar_width=None),"{task.completed}/{task.total}","•","[progress.percentage]{task.percentage:>3.1f}%",console=console,)
    process_bar.start()

    # we need to consume the iterator once to get the total
    # for the progress bar
    count = sum(1 for i in walk_filepath_list(filepath_list,recursive))
    task_process_bar = process_bar.add_task("Main task",total=count)
    with ProcessPoolExecutor(max_workers=workers) as pool:
        # shared dict between threads
        # [Future] => [filepath]
        future_to_filepath = {}
        submit_thread = threading.Thread(
            target=thread_submit,args=(pool,future_to_filepath)
        )
        submit_thread.start()
        while len(future_to_filepath.keys()) != count:
            for future in as_completed(future_to_filepath):
                filepath = future_to_filepath[future]
                # print(f"processing future: {filepath}")
                try:
                    data = future.result()
                finally:
                    # update progress bar
                    process_bar.update(task_process_bar,advance=1)
    process_bar.stop()


def entrypoint():
    args = docopt(__doc__)
    main(args)


if __name__ == "__main__":
    entrypoint()

但是,进度栏未按预期更新。 更糟糕的是,在某些情况下处理似乎没有结束。

  • 我更新字典future_to_filepath时是比赛条件吗?
  • 您将如何使用concurrent.futures创建一个 submit 线程和一个 process_results 线程?

非常感谢!

解决方法

查看我对您问题的评论,然后:

更改:

submit_thread = threading.Thread(
    target=thread_submit,args=(pool,filepath_list,recursive,future_to_filepath)
)
submit_thread.start()

收件人:

thread_submit(pool,future_to_filepath)

(对此函数名称进行更改,因为它不再作为单独的线程运行,这是件好事-create_futures呢?)

并删除外部循环:

while len(future_to_filepath.keys()) != count:

最后,不清楚您的实际process_task将如何处理该文件,但是肯定有可能受I / O约束。在这种情况下,您可能会受益于使用ThreadPoolExecutor类,而该类可以很容易地替换为ProcessPoolExecutor类,在这种情况下,您应该考虑指定更多的工人,可能等于{{1} }。由于您当前的count除了睡觉外没有其他事情,因此可能会通过与更多的工人一起穿线而受益。

更新

您可以减少运行process_task 的时间的一件事是将函数修改为传递单个路径而不是列表,并处理每个路径在原始列表中同时存在于单独的线程中。在下面的代码中,为了方便起见,我使用walk_filepath_list ThreadPoolExecutor函数,这实际上要求对(新重命名的)map函数的参数进行反转,以便可以使用{{1} }以“编码”所有调用的第一个参数walk_filepath

functools.partial

更新2

上述代码的基准测试表明,它不能节省时间(也许目录是否位于不同的物理驱动器上?)。

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


依赖报错 idea导入项目后依赖报错,解决方案:https://blog.csdn.net/weixin_42420249/article/details/81191861 依赖版本报错:更换其他版本 无法下载依赖可参考:https://blog.csdn.net/weixin_42628809/a
错误1:代码生成器依赖和mybatis依赖冲突 启动项目时报错如下 2021-12-03 13:33:33.927 ERROR 7228 [ main] o.s.b.d.LoggingFailureAnalysisReporter : *************************** APPL
错误1:gradle项目控制台输出为乱码 # 解决方案:https://blog.csdn.net/weixin_43501566/article/details/112482302 # 在gradle-wrapper.properties 添加以下内容 org.gradle.jvmargs=-Df
错误还原:在查询的过程中,传入的workType为0时,该条件不起作用 &lt;select id=&quot;xxx&quot;&gt; SELECT di.id, di.name, di.work_type, di.updated... &lt;where&gt; &lt;if test=&qu
报错如下,gcc版本太低 ^ server.c:5346:31: 错误:‘struct redisServer’没有名为‘server_cpulist’的成员 redisSetCpuAffinity(server.server_cpulist); ^ server.c: 在函数‘hasActiveC
解决方案1 1、改项目中.idea/workspace.xml配置文件,增加dynamic.classpath参数 2、搜索PropertiesComponent,添加如下 &lt;property name=&quot;dynamic.classpath&quot; value=&quot;tru
删除根组件app.vue中的默认代码后报错:Module Error (from ./node_modules/eslint-loader/index.js): 解决方案:关闭ESlint代码检测,在项目根目录创建vue.config.js,在文件中添加 module.exports = { lin
查看spark默认的python版本 [root@master day27]# pyspark /home/software/spark-2.3.4-bin-hadoop2.7/conf/spark-env.sh: line 2: /usr/local/hadoop/bin/hadoop: No s
使用本地python环境可以成功执行 import pandas as pd import matplotlib.pyplot as plt # 设置字体 plt.rcParams[&#39;font.sans-serif&#39;] = [&#39;SimHei&#39;] # 能正确显示负号 p
错误1:Request method ‘DELETE‘ not supported 错误还原:controller层有一个接口,访问该接口时报错:Request method ‘DELETE‘ not supported 错误原因:没有接收到前端传入的参数,修改为如下 参考 错误2:cannot r
错误1:启动docker镜像时报错:Error response from daemon: driver failed programming external connectivity on endpoint quirky_allen 解决方法:重启docker -&gt; systemctl r
错误1:private field ‘xxx‘ is never assigned 按Altʾnter快捷键,选择第2项 参考:https://blog.csdn.net/shi_hong_fei_hei/article/details/88814070 错误2:启动时报错,不能找到主启动类 #
报错如下,通过源不能下载,最后警告pip需升级版本 Requirement already satisfied: pip in c:\users\ychen\appdata\local\programs\python\python310\lib\site-packages (22.0.4) Coll
错误1:maven打包报错 错误还原:使用maven打包项目时报错如下 [ERROR] Failed to execute goal org.apache.maven.plugins:maven-resources-plugin:3.2.0:resources (default-resources)
错误1:服务调用时报错 服务消费者模块assess通过openFeign调用服务提供者模块hires 如下为服务提供者模块hires的控制层接口 @RestController @RequestMapping(&quot;/hires&quot;) public class FeignControl
错误1:运行项目后报如下错误 解决方案 报错2:Failed to execute goal org.apache.maven.plugins:maven-compiler-plugin:3.8.1:compile (default-compile) on project sb 解决方案:在pom.
参考 错误原因 过滤器或拦截器在生效时,redisTemplate还没有注入 解决方案:在注入容器时就生效 @Component //项目运行时就注入Spring容器 public class RedisBean { @Resource private RedisTemplate&lt;String
使用vite构建项目报错 C:\Users\ychen\work&gt;npm init @vitejs/app @vitejs/create-app is deprecated, use npm init vite instead C:\Users\ychen\AppData\Local\npm-