Neutron, nuestro motor de IA, obtuvo un 96.75% en el benchmark CyberGym de UC Berkeley. Más información

Ingeniería

Ingeniería

ProcessPoolExecutor de Python: un pool de procesos personalizado

En Ostorlab escaneamos cientos de aplicaciones móviles cada día; cada escaneo consume muchos recursos, pero desde el principio tuvimos que optimizar el código para la velocidad y aprovechar al máximo los recursos de la nube.

En Ostorlab escaneamos cientos de aplicaciones móviles cada día; cada escaneo consume muchos recursos, pero al mismo tiempo, desde el principio, hemos tenido que optimizar el código para la velocidad y aprovechar al máximo los recursos de la nube.

Sin volver a discutir la concurrencia y el paralelismo en Python, ya que existen muchos recursos excelentes sobre el tema, mi favorito es «Fluent Python» de Luciano Ramalho.

Pero aquí van algunas notas rápidas:

En CPython, el GIL impide que varios hilos nativos ejecuten bytecode de Python al mismo tiempo. Esto significa que el uso de hilos no es adecuado para el paralelismo limitado por CPU. Recomiendo estas charlas para obtener más información sobre el tema: «Understanding the Python GIL» y «Removing Python's GIL: The Gilectomy». El procesamiento mediante procesos se adapta mejor al paralelismo limitado por CPU; sin embargo, crear un proceso supone una sobrecarga considerable, y si el objetivo es ejecutar varias tareas tanto de forma concurrente como en paralelo, ejecutar cada tarea en un proceso independiente es muy ineficiente. Para resolver ese problema, la biblioteca concurrent.futures se introdujo en Python 3.2 y se portó a Python 2.5. La biblioteca permite crear un número de workers (por defecto, el número de núcleos de su equipo) a los que se pueden pasar tareas para su procesamiento. Las tareas se envían a través de colas después de serializarse con Pickle.

Este es un ejemplo del GitHub público 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

Al experimentar con esto, nos topamos rápidamente con muchas limitaciones debidas a las restricciones de la serialización con Pickle; los dos problemas principales fueron:

Algunas tareas estaban decoradas con decoradores parametrizados, lo que provocaba un error similar a este: Error pickling <function>. Los argumentos que se pasaban eran en ocasiones bastante grandes y provocaban un error de recursion limit. Para el segundo problema, tuvimos que pensar en una forma de pasar argumentos a las tareas sin pasar por la cola y por el proceso de serialización.

Antes de usar ProcessPoolExecutor, utilizábamos Process simples que recibían los argumentos en el parámetro args:

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

Pasar los argumentos de esta forma no planteaba ningún problema de tamaño y, por supuesto, no requiere serialización. Python permite pasar colas, pipes e incluso memoria compartida mediante los tipos multiprocessing.Value y multiprocessing.Array (véase multiprocessing).

La particularidad de nuestro código era que los argumentos pasados siempre eran los mismos y de gran tamaño, así que decidimos modificar ProcessPoolExecutor para pasar estos argumentos en la inicialización, y que este los transmitiera a los procesos de trabajo creados.

Las tareas se ejecutarán con estos argumentos como parámetros iniciales; estos son los cambios principales, que añaden los atributos initial_args y 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

Pase los atributos a los procesos creados:

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)

Use los parámetros para ejecutar cada tarea:

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))

El beneficio fue aún mayor, ya que eliminamos todo el procesamiento adicional debido a la serialización del mismo objeto.

El mismo código pasó de ejecutarse durante más de 20 minutos a menos de un minuto gracias al uso de multiprocessing, a una sobrecarga fija de creación de procesos y a evitar toda la serialización adicional.

Como siguiente paso, planeamos experimentar con Go y Unikernels para ver si podemos obtener más rendimiento.

Etiquetas:

Python, Performance