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.