Python 3 asyncio con aioboto3 parece secuencial

Aug 28 2020

Estoy portando un script de Python 3 simple a AWS Lambda. El script es simple: recopila información de una docena de objetos S3 y devuelve los resultados.

El script utilizado multiprocessing.Poolpara recopilar todos los archivos en paralelo. Aunque multiprocessingno se puede utilizar en un entorno AWS Lambda porque /dev/shmfalta. Así que pensé que en lugar de escribir un sucio multiprocessing.Process/ multiprocessing.Queuereemplazo, lo intentaría asyncio.

Estoy usando la última versión de aioboto3(8.0.5) en Python 3.8.

Mi problema es que parece que no puedo obtener ninguna mejora entre una descarga secuencial ingenua de los archivos y un bucle de eventos asincio que multiplexa las descargas.

Aquí están las dos versiones de mi código.

import sys
import asyncio
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

import boto3
import aioboto3

BUCKET = 'some-bucket'
KEYS = [
    'some/key/1',
    [...]
    'some/key/10',
]

async def download_aio():
    """Concurrent download of all objects from S3"""
    async with aioboto3.client('s3') as s3:
        objects = [s3.get_object(Bucket=BUCKET, Key=k) for k in KEYS]
        objects = await asyncio.gather(*objects)
        buffers = await asyncio.gather(*[o['Body'].read() for o in objects])

def download():
    """Sequentially download all objects from S3"""
    s3 = boto3.client('s3')
    for key in KEYS:
        object = s3.get_object(Bucket=BUCKET, Key=key)
        object['Body'].read()

def run_sequential():
    download()

def run_concurrent():
    loop = asyncio.get_event_loop()
    #loop.set_default_executor(ProcessPoolExecutor(10))
    #loop.set_default_executor(ThreadPoolExecutor(10))
    loop.run_until_complete(download_aio())

El tiempo para ambos run_sequential()y run_concurrent()es bastante similar (~ 3 segundos para una docena de archivos de 10 MB). Estoy convencido de que la versión concurrente no lo es, por múltiples razones:

  • Intenté cambiar a Process/ThreadPoolExecutor, y los procesos / subprocesos se generaron durante la duración de la función, aunque no están haciendo nada
  • El tiempo entre secuencial y concurrente es muy similar, aunque mi interfaz de red definitivamente no está saturada y la CPU tampoco está vinculada
  • El tiempo que tarda la versión simultánea aumenta linealmente con el número de archivos.

Estoy seguro de que falta algo, pero no puedo entender qué.

¿Algunas ideas?

Respuestas

NewbiZ Aug 28 2020 at 11:50

Después de perder algunas horas tratando de entender cómo usarlo aioboto3correctamente, decidí simplemente cambiar a mi solución de respaldo. Terminé lanzando mi propia versión ingenua de multiprocessing.Poolpara usarla dentro de un entorno AWS lambda.

Si alguien se encuentra con este hilo en el futuro, aquí está. Está lejos de ser perfecto, pero es lo suficientemente fácil de reemplazar multiprocessing.Poolcomo está para mis casos simples.

from multiprocessing import Process, Pipe
from multiprocessing.connection import wait


class Pool:
    """Naive implementation of a process pool with mp.Pool API.

    This is useful since multiprocessing.Pool uses a Queue in /dev/shm, which
    is not mounted in an AWS Lambda environment.
    """

    def __init__(self, process_count=1):
        assert process_count >= 1
        self.process_count = process_count

    @staticmethod
    def wrap_pipe(pipe, index, func):
        def wrapper(args):
            try:
                result = func(args)
            except Exception as exc:  # pylint: disable=broad-except
                result = exc
            pipe.send((index, result))
        return wrapper

    def __enter__(self):
        return self

    def __exit__(self, exc_type, exc_value, exc_traceback):
        pass

    def map(self, function, arguments):
        pending = list(enumerate(arguments))
        running = []
        finished = [None] * len(pending)
        while pending or running:
            # Fill the running queue with new jobs
            while len(running) < self.process_count:
                if not pending:
                    break
                index, args = pending.pop(0)
                pipe_parent, pipe_child = Pipe(False)
                process = Process(
                    target=Pool.wrap_pipe(pipe_child, index, function),
                    args=(args, ))
                process.start()
                running.append((index, process, pipe_parent))
            # Wait for jobs to finish
            for pipe in wait(list(map(lambda t: t[2], running))):
                index, result = pipe.recv()
                # Remove the finished job from the running list
                running = list(filter(lambda x: x[0] != index, running))
                # Add the result to the finished list
                finished[index] = result

        return finished