草庐IT

python - ProcessPoolExecutor 中的 ThreadPoolExecutor

coder 2023-08-16 原文

我是新来的 the futures module并且有一项可以从并行化中受益的任务;但我似乎无法确切地弄清楚如何为线程设置函数和为进程设置函数。我很感激任何人都可以在这个问题上提供的帮助。

我正在运行 particle swarm optimization (PSO) .在不深入了解 PSO 本身的情况下,以下是我的代码的基本布局:

有一个Particle类,带有 getFitness(self)方法(计算一些指标并将其存储在 self.fitness 中)。一个 PSO 模拟有多个粒子实例(很容易超过 10 个;对于某些模拟,100 秒甚至 1000 秒)。
每隔一段时间,我就必须计算粒子的适应度。目前,我在 for 循环中执行此操作:

for p in listOfParticles:
  p.getFitness(args)

但是,我注意到每个粒子的适应度可以相互独立计算。这使得这种适应度计算成为并行化的主要候选者。确实,我可以做map(lambda p: p.getFitness(args), listOfParticles) .

现在,我可以使用 futures.ProcessPoolExecutor 轻松做到这一点。 :
with futures.ProcessPoolExecutor() as e:
  e.map(lambda p: p.getFitness(args), listOfParticles)

由于调用 p.getFitness 的副作用存储在每个粒子本身,我不必担心从 futures.ProcessPoolExecutor() 得到返回.

到现在为止还挺好。但现在我注意到 ProcessPoolExecutor创建新进程,这意味着它复制内存,这很慢。我希望能够共享内存 - 所以我应该使用线程。这很好,直到我意识到在每个进程中运行具有多个线程的多个进程可能会更快,因为多个线程仍然只在我可爱的 ​​8 核机器的一个处理器上运行。

这是我遇到麻烦的地方:
根据我看到的例子,ThreadPoolExecutorlist 上运行. ProcessPoolExecutor也是如此.所以我不能在 ProcessPoolExecutor 中做任何迭代农场到 ThreadPoolExecutor因为那时 ThreadPoolExecutor将得到一个单一的对象来处理(见我的尝试,贴在下面)。
另一方面,我不能切片 listOfParticles我自己,因为我想要ThreadPoolExecutor发挥自己的魔力来弄清楚需要多少线程。

所以,大问题(终于) :
我应该如何构建我的代码,以便我可以使用进程和线程有效地并行化以下内容:
for p in listOfParticles:
  p.getFitness()

这是我一直在尝试的,但我不敢尝试运行它,因为我知道它行不通:
>>> def threadize(func, L, mw):
...     with futures.ThreadpoolExecutor(max_workers=mw) as executor:
...             for i in L:
...                     executor.submit(func, i)
... 

>>> def processize(func, L, mw):
...     with futures.ProcessPoolExecutor() as executor:
...             executor.map(lambda i: threadize(func, i, mw), L)
...

我很感激关于如何解决这个问题,甚至如何改进我的方法的任何想法

万一重要,我在python3.3.2上

最佳答案

我将为您提供将进程与线程混合以解决问题的工作代码,但这不是您所期望的 ;-) 第一件事是制作一个不会危及您的真实数据的模拟程序。尝试一些无害的东西。所以这是开始:

class Particle:
    def __init__(self, i):
        self.i = i
        self.fitness = None
    def getfitness(self):
        self.fitness = 2 * self.i

现在我们有东西可以玩了。接下来是一些常量:
MAX_PROCESSES = 3
MAX_THREADS = 2 # per process
CHUNKSIZE = 100

摆弄那些来品尝。 CHUNKSIZE稍后会解释。

第一个让你惊讶的是我的最低级工作函数做了什么。那是因为你在这里过于乐观了:

Since the side-effects of calling p.getFitness are stored in each particle itself, I don't have to worry about getting a return from futures.ProcessPoolExecutor().



唉,在工作进程中所做的任何事情都不会对 Particle 产生任何影响。主程序中的实例。工作进程处理 Particle 的副本实例,是否通过 fork() 的写时复制实现或者因为它正在处理一个通过解酸 Particle 制作的副本泡菜通过进程。

所以如果你想让你的主程序看到健身结果,你需要安排将信息发送回主程序。因为我对你的实际程序了解不够,这里我假设 Particle().i是一个唯一的整数,并且主程序可以轻松地将整数映射回 Particle实例。考虑到这一点,这里的最低级工作函数需要返回一对:唯一整数和适应度结果:
def thread_worker(p):
    p.getfitness()
    return (p.i, p.fitness)

鉴于此,很容易传播 Particle 的列表。 s 跨线程,并返回 (particle_id, fitness) 的列表结果:
def proc_worker(ps):
    import concurrent.futures as cf
    with cf.ThreadPoolExecutor(max_workers=MAX_THREADS) as e:
        result = list(e.map(thread_worker, ps))
    return result

笔记:
  • 这是每个工作进程将运行的功能。
  • 我使用的是 Python 3,所以使用 list()给力e.map()将所有结果具体化到一个列表中。
  • 正如评论中提到的,在 CPython 下,跨线程传播 CPU 绑定(bind)任务比在单个线程中完成所有任务要慢。

  • 剩下的就是写代码传播Particle的列表了s 跨进程,并检索结果。用 multiprocessing 很容易做到这一点,这就是我要使用的。我不知道是否concurrent.futures可以做到(考虑到我们也在线程中混合),但不在乎。但是因为我给你工作代码,你可以玩它并报告回来;-)
    if __name__ == "__main__":
        import multiprocessing
    
        particles = [Particle(i) for i in range(100000)]
        # Note the code below relies on that particles[i].i == i
        assert all(particles[i].i == i for i in range(len(particles)))
    
        pool = multiprocessing.Pool(MAX_PROCESSES)
        for result_list in pool.imap_unordered(proc_worker,
                          (particles[i: i+CHUNKSIZE]
                           for i in range(0, len(particles), CHUNKSIZE))):
            for i, fitness in result_list:
                particles[i].fitness = fitness
    
        pool.close()
        pool.join()
    
        assert all(p.fitness == 2*p.i for p in particles)
    

    笔记:
  • 我正在打破 Particle 的名单s 成块“手工”。就是这样CHUNKSIZE是为了。那是因为一个工作进程需要一个 Particle 的列表。 s 继续工作,反过来,这是因为这就是 futures map()功能要。无论如何,将工作分块是一个好主意,这样您就可以真正物有所值,以换取每次调用的进程间开销。
  • imap_unordered()不保证返回结果的顺序。这使实现可以更自由地尽可能高效地安排工作。我们不关心这里的顺序,所以没关系。
  • 请注意,循环检索 (particle_id, fitness)结果,并修改 Particle相应的实例。也许你真正的.getfitnessParticle 进行其他突变实例 - 无法猜测。无论如何,主程序永远不会看到任何“魔法”在 worker 身上发生的变化——你必须明确地安排。在限制中,您可以返回 (particle_id, particle_instance)成对,并替换 Particle主程序中的实例。然后它们会反射(reflect)在工作进程中进行的所有更改。

  • 玩得开心 :-)

    future 一路下跌

    结果证明它很容易更换 multiprocessing .以下是变化。这也(如前所述)取代了原来的 Particle实例,以便捕获所有突变。不过,这里有一个权衡:酸洗一个实例比酸洗单个“适应度”结果需要“多得多”的字节。更多的网络流量。选择你的毒药;-)

    返回变异的实例只需要替换 thread_worker() 的最后一行,像这样:
    return (p.i, p)
    

    然后用这个替换所有的“ main ”块:
    def update_fitness():
        import concurrent.futures as cf
        with cf.ProcessPoolExecutor(max_workers=MAX_PROCESSES) as e:
            for result_list in e.map(proc_worker,
                          (particles[i: i+CHUNKSIZE]
                           for i in range(0, len(particles), CHUNKSIZE))):
                for i, p in result_list:
                    particles[i] = p
    
    if __name__ == "__main__":
        particles = [Particle(i) for i in range(500000)]
        assert all(particles[i].i == i for i in range(len(particles)))
    
        update_fitness()
    
        assert all(particles[i].i == i for i in range(len(particles)))
        assert all(p.fitness == 2*p.i for p in particles)
    

    该代码与 multiprocessor 非常相似舞蹈。就个人而言,我会使用 multiprocessing版本,因为 imap_unordered是有值(value)的。这是简化界面的一个问题:他们通常以隐藏有用的可能性为代价来购买简单性。

    关于python - ProcessPoolExecutor 中的 ThreadPoolExecutor,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/19994478/

    有关python - ProcessPoolExecutor 中的 ThreadPoolExecutor的更多相关文章

    1. ruby - 如何从 ruby​​ 中的字符串运行任意对象方法? - 2

      总的来说,我对ruby​​还比较陌生,我正在为我正在创建的对象编写一些rspec测试用例。许多测试用例都非常基础,我只是想确保正确填充和返回值。我想知道是否有办法使用循环结构来执行此操作。不必为我要测试的每个方法都设置一个assertEquals。例如:describeitem,"TestingtheItem"doit"willhaveanullvaluetostart"doitem=Item.new#HereIcoulddotheitem.name.shouldbe_nil#thenIcoulddoitem.category.shouldbe_nilendend但我想要一些方法来使用

    2. ruby - 其他文件中的 Rake 任务 - 2

      我试图在一个项目中使用rake,如果我把所有东西都放到Rakefile中,它会很大并且很难读取/找到东西,所以我试着将每个命名空间放在lib/rake中它自己的文件中,我添加了这个到我的rake文件的顶部:Dir['#{File.dirname(__FILE__)}/lib/rake/*.rake'].map{|f|requiref}它加载文件没问题,但没有任务。我现在只有一个.rake文件作为测试,名为“servers.rake”,它看起来像这样:namespace:serverdotask:testdoputs"test"endend所以当我运行rakeserver:testid时

    3. ruby-on-rails - Ruby net/ldap 模块中的内存泄漏 - 2

      作为我的Rails应用程序的一部分,我编写了一个小导入程序,它从我们的LDAP系统中吸取数据并将其塞入一个用户表中。不幸的是,与LDAP相关的代码在遍历我们的32K用户时泄漏了大量内存,我一直无法弄清楚如何解决这个问题。这个问题似乎在某种程度上与LDAP库有关,因为当我删除对LDAP内容的调用时,内存使用情况会很好地稳定下来。此外,不断增加的对象是Net::BER::BerIdentifiedString和Net::BER::BerIdentifiedArray,它们都是LDAP库的一部分。当我运行导入时,内存使用量最终达到超过1GB的峰值。如果问题存在,我需要找到一些方法来更正我的代

    4. python - 如何使用 Ruby 或 Python 创建一系列高音调和低音调的蜂鸣声? - 2

      关闭。这个问题是opinion-based.它目前不接受答案。想要改进这个问题?更新问题,以便editingthispost可以用事实和引用来回答它.关闭4年前。Improvethisquestion我想在固定时间创建一系列低音和高音调的哔哔声。例如:在150毫秒时发出高音调的蜂鸣声在151毫秒时发出低音调的蜂鸣声200毫秒时发出低音调的蜂鸣声250毫秒的高音调蜂鸣声有没有办法在Ruby或Python中做到这一点?我真的不在乎输出编码是什么(.wav、.mp3、.ogg等等),但我确实想创建一个输出文件。

    5. ruby-on-rails - Rails 3 中的多个路由文件 - 2

      Rails2.3可以选择随时使用RouteSet#add_configuration_file添加更多路由。是否可以在Rails3项目中做同样的事情? 最佳答案 在config/application.rb中:config.paths.config.routes在Rails3.2(也可能是Rails3.1)中,使用:config.paths["config/routes"] 关于ruby-on-rails-Rails3中的多个路由文件,我们在StackOverflow上找到一个类似的问题

    6. ruby-on-rails - Rails - 一个 View 中的多个模型 - 2

      我需要从一个View访问多个模型。以前,我的links_controller仅用于提供以不同方式排序的链接资源。现在我想包括一个部分(我假设)显示按分数排序的顶级用户(@users=User.all.sort_by(&:score))我知道我可以将此代码插入每个链接操作并从View访问它,但这似乎不是“ruby方式”,我将需要在不久的将来访问更多模型。这可能会变得很脏,是否有针对这种情况的任何技术?注意事项:我认为我的应用程序正朝着单一格式和动态页面内容的方向发展,本质上是一个典型的网络应用程序。我知道before_filter但考虑到我希望应用程序进入的方向,这似乎很麻烦。最终从任何

    7. ruby-on-rails - Rails 3.2.1 中 ActionMailer 中的未定义方法 'default_content_type=' - 2

      我在我的项目中添加了一个系统来重置用户密码并通过电子邮件将密码发送给他,以防他忘记密码。昨天它运行良好(当我实现它时)。当我今天尝试启动服务器时,出现以下错误。=>BootingWEBrick=>Rails3.2.1applicationstartingindevelopmentonhttp://0.0.0.0:3000=>Callwith-dtodetach=>Ctrl-CtoshutdownserverExiting/Users/vinayshenoy/.rvm/gems/ruby-1.9.3-p0/gems/actionmailer-3.2.1/lib/action_mailer

    8. ruby-on-rails - Rails 应用程序中的 Rails : How are you using application_controller. rb 是新手吗? - 2

      刚入门rails,开始慢慢理解。有人可以解释或给我一些关于在application_controller中编码的好处或时间和原因的想法吗?有哪些用例。您如何为Rails应用程序使用应用程序Controller?我不想在那里放太多代码,因为据我了解,每个请求都会调用此Controller。这是真的? 最佳答案 ApplicationController实际上是您应用程序中的每个其他Controller都将从中继承的类(尽管这不是强制性的)。我同意不要用太多代码弄乱它并保持干净整洁的态度,尽管在某些情况下ApplicationContr

    9. ruby-on-rails - form_for 中不在模型中的自定义字段 - 2

      我想向我的Controller传递一个参数,它是一个简单的复选框,但我不知道如何在模型的form_for中引入它,这是我的观点:{:id=>'go_finance'}do|f|%>Transferirde:para:Entrada:"input",:placeholder=>"Quantofoiganho?"%>Saída:"output",:placeholder=>"Quantofoigasto?"%>Nota:我想做一个额外的复选框,但我该怎么做,模型中没有一个对象,而是一个要检查的对象,以便在Controller中创建一个ifelse,如果没有检查,请帮助我,非常感谢,谢谢

    10. ruby - rspec 需要 .rspec 文件中的 spec_helper - 2

      我注意到像bundler这样的项目在每个specfile中执行requirespec_helper我还注意到rspec使用选项--require,它允许您在引导rspec时要求一个文件。您还可以将其添加到.rspec文件中,因此只要您运行不带参数的rspec就会添加它。使用上述方法有什么缺点可以解释为什么像bundler这样的项目选择在每个规范文件中都需要spec_helper吗? 最佳答案 我不在Bundler上工作,所以我不能直接谈论他们的做法。并非所有项目都checkin.rspec文件。原因是这个文件,通常按照当前的惯例,只

    随机推荐