Neutron, notre moteur d’IA, a obtenu un score de 96,75 % sur le benchmark CyberGym de l’UC Berkeley. En savoir plus

Ingénierie

Ingénierie

ProcessPoolExecutor en Python : un pool de processus personnalisé

Chez Ostorlab, nous analysons des centaines d'applications mobiles chaque jour ; chaque scan consomme beaucoup de ressources, et nous avons dû, dès le début, optimiser le code pour la vitesse et tirer le meilleur parti des ressources cloud.

Chez Ostorlab, nous analysons des centaines d'applications mobiles chaque jour ; chaque scan consomme beaucoup de ressources, et nous avons dû, dès le début, optimiser le code pour la vitesse et tirer le meilleur parti des ressources cloud.

Sans revenir sur la concurrence et le parallélisme en Python, sujet sur lequel il existe de nombreuses excellentes ressources (ma préférée est « Fluent Python » de Luciano Ramalho), voici quelques notes rapides :

Dans CPython, le GIL empêche plusieurs threads natifs d'exécuter du bytecode Python en même temps. Le threading est donc inadapté au parallélisme limité par le CPU. Je recommande ces conférences pour en savoir plus : « Understanding the Python GIL » et « Removing Python's GIL: The Gilectomy ». Le multiprocessing est mieux adapté au parallélisme limité par le CPU, mais le lancement d'un processus entraîne un surcoût important : si l'objectif est d'exécuter plusieurs tâches à la fois de façon concurrente et parallèle, exécuter chaque tâche dans un processus distinct est très inefficace. Pour résoudre ce problème, la bibliothèque concurrent.futures a été introduite dans Python 3.2 et rétroportée vers Python 2.5. Elle permet de lancer un certain nombre de workers (par défaut, le nombre de cœurs de votre machine) auxquels vous pouvez confier des tâches à traiter. Les tâches sont envoyées via des files après avoir été sérialisées avec Pickle.

Voici un exemple tiré du dépôt GitHub public de Fluent Python :

"""Download flags of top 20 countries by population
ProcessPoolExecutor version
Sample run::
    $ python3 flags_threadpool.py    BD retrieved.    EG retrieved.    CN retrieved.    ...    PH retrieved.    
    US retrieved.    IR retrieved.    20 flags downloaded in 0.93s
"""
# BEGIN FLAGS_PROCESSPOOL

from concurrent import futures
from flags import save_flag, get_flag, show, main

MAX_WORKERS = 20

def download_one(cc):
    image = get_flag(cc)
    show(cc)
    save_flag(image, cc.lower() + '.gif')
    return cc


def download_many(cc_list):
    with futures.ProcessPoolExecutor() as executor:  # <1>        
        res = executor.map(download_one, sorted(cc_list))

    return len(list(res))


if __name__ == '__main__':
    main(download_many)
# END FLAGS_PROCESSPOOL

En l'expérimentant, nous avons rapidement rencontré de nombreuses limites dues à la sérialisation avec Pickle ; les deux principaux problèmes étaient les suivants :

Certaines tâches étaient décorées par des décorateurs paramétrés, ce qui provoquait une erreur semblable à celle-ci : Error pickling <function>. Les arguments transmis étaient parfois très volumineux et provoquaient une erreur de recursion limit. Pour le second problème, nous avons dû trouver un moyen de passer les arguments aux tâches sans passer par la file ni par le processus de sérialisation.

Avant d'utiliser ProcessPoolExecutor, nous utilisions de simples Process qui recevaient les arguments via le paramètre args :

p = multiprocessing.Process(
        target=_process_worker,
        args=(self._call_queue,
              self._result_queue,
              self._initial_args,
              self._initial_kwargs))

Passer les arguments de cette façon ne posait aucun problème de taille et ne nécessite évidemment aucune sérialisation. Python permet de passer des files, des pipes et même de la mémoire partagée avec les types multiprocessing.Value et multiprocessing.Array (voir multiprocessing).

La particularité de notre code était que les arguments transmis étaient toujours les mêmes et de taille importante ; nous avons donc décidé de modifier ProcessPoolExecutor pour qu'il reçoive ces arguments à l'initialisation et les transmette ensuite aux processus de travail qu'il lance.

Les tâches seront exécutées avec ces arguments comme paramètres initiaux. Voici les principales modifications, avec l'ajout des attributs initial_args et initial_kwargs :

class ProcessPoolExecutor(_base.Executor):
    def __init__(self, max_workers=None, *args, **kwargs):
        """Initializes a new ProcessPoolExecutor instance.
        Args:            max_workers: The maximum number of processes that can be used to               
         execute the given calls. If None or not given then as many                
         worker processes will be created as the machine has processors.        
        """        
        _check_system_limits()

        if max_workers is None:
            self._max_workers = multiprocessing.cpu_count()
        else:
            self._max_workers = max_workers

        # Make the call queue slightly larger than the number of processes to        
        # prevent the worker processes from idling. But don't make it too big        
        # because futures in the call queue cannot be cancelled.        
        self._call_queue = multiprocessing.Queue(self._max_workers + EXTRA_QUEUED_CALLS)
        self._result_queue = multiprocessing.Queue()
        self._work_ids = queue.Queue()
        self._queue_management_thread = None        
        self._processes = set()

        # Shutdown is a two-step process.        self._shutdown_thread = False        self._shutdown_lock = threading.Lock()
        self._queue_count = 0  
        self._pending_work_items = {}
        self._initial_args = args
        self._initial_kwargs = kwargs

Transmettre les attributs aux processus lancés :

def _adjust_process_count(self):
    for _ in range(len(self._processes), self._max_workers):
        p = multiprocessing.Process(
                target=_process_worker,
                args=(self._call_queue,
                      self._result_queue,
                      self._initial_args,
                      self._initial_kwargs))
        p.start()
        self._processes.add(p)

Utiliser les paramètres pour exécuter chaque tâche :

def _process_worker(call_queue, result_queue, args, kwargs):
    """Evaluates calls from call_queue and places the results in result_queue.
    This worker is run in a separate process.
    Args:        call_queue: A multiprocessing.Queue of _CallItems that will be read and            
    evaluated by the worker.        
    result_queue: A multiprocessing.Queue of _ResultItems that will written            
    to by the worker.        
    shutdown: A multiprocessing.Event that will be set as a signal to the            
    worker that it should exit when call_queue is empty.    
    """    
    while True:
        call_item = call_queue.get(block=True)
        if call_item is None:
            # Wake up queue management thread
            result_queue.put(None)
            return        
        try:
            args += call_item.args
            kwargs.update(**call_item.kwargs)
            r = call_item.fn(*args, **kwargs)
        except BaseException:
            e = sys.exc_info()[1]
            result_queue.put(_ResultItem(call_item.work_id, exception=e))
        else:
            result_queue.put(_ResultItem(call_item.work_id, result=r))

Le gain a été encore plus important, car nous avons supprimé tout le traitement superflu lié à la sérialisation du même objet.

Le même code est passé de plus de 20 minutes d'exécution à moins d'une minute grâce au multiprocessing, à un surcoût de lancement des processus désormais fixe et à l'absence de toute sérialisation superflue.

Pour aller plus loin, nous prévoyons d'expérimenter avec Go et les Unikernels afin de voir si nous pouvons gagner en performances.