【问题标题】:Launch MPI worker processes in a Django app在 Django 应用程序中启动 MPI 工作进程
【发布时间】:2014-12-11 04:17:07
【问题描述】:

我们想在我们的 Django 应用程序中运行一些后台进程。看起来 Celery 是最常见的解决方案,但我们团队对 MPI 更熟悉,所以我正在尝试它。我想创建一个 Django 管理命令来启动 MPI 工作人员池,所以我阅读了 Django admin commands 和 MPI4py 的 dynamic process management。

我编写了一个管理命令来运行车队管理器和一个管理命令来运行一个工人。车队经理成功使用MPI.COMM_SELF.Spawn() 启动工作人员,但他们无法相互通信。经理和第一个工人的等级都是 0,所以看起来他们正在使用不同的通信器。

如何让经理和员工使用同一个沟通器?

【问题讨论】:

    标签: python django mpi mpi4py


    【解决方案1】:

    诀窍是合并两个通信器,如this answer to a C question 中所述。在MPI4py documentation 的帮助下,我将其转换为 Python:

    # myproject/myapp/management/commands/runfleet.py
    from mpi4py import MPI
    from optparse import make_option
    
    from django.core.management.base import BaseCommand
    
    import sys
    
    class Command(BaseCommand):
        help = 'Launches the example manager and workers.'
    
        option_list = BaseCommand.option_list + (
            make_option('--workers', '-w', type='int', default=1), )
    
        def handle(self, *args, **options):
            self.stdout.write("Manager launched.")
    
            worker_count = options['workers']
            manage_script = sys.argv[0]
            comm = MPI.COMM_SELF.Spawn(sys.executable,
                                       args=[manage_script, 'fleetworker'],
                                       maxprocs=worker_count).Merge()
            self.stdout.write('Manager rank {}.'.format(comm.Get_rank()))
    
            start_data = [None] # First item is sent to manager and ignored
            for i in range(worker_count):
                start_data.append("Item {}".format(i))
            comm.scatter(start_data)
    
            end_data = comm.gather()
            self.stdout.write('Manager received data {!r}.'.format(end_data))
    
            comm.Disconnect()
    

    worker 命令如下所示:

    # myproject/myapp/management/commands/fleetworker.py
    from mpi4py import MPI
    
    from django.core.management.base import BaseCommand
    
    class Command(BaseCommand):
        help = 'Example worker process.'
    
        def handle(self, *args, **options):
            self.stdout.write("Worker launched.")
    
            comm = MPI.Comm.Get_parent().Merge()
            rank = comm.Get_rank()
            self.stdout.write('Worker rank {}.'.format(rank))
    
            data = comm.scatter()
            result = "{!r}, {!r}".format(rank, data)
    
            comm.gather(result)
            self.stdout.write("Finished worker.")
    
            comm.Disconnect()
    

    【讨论】:

      猜你喜欢
      • 2022-11-11
      • 2011-02-18
      • 2014-10-12
      • 2012-12-22
      • 2013-06-22
      • 1970-01-01
      • 1970-01-01
      • 2015-02-17
      • 1970-01-01
      相关资源
      最近更新 更多