Python 3 assíncio com aioboto3 parece sequencial

Aug 28 2020

Estou portando um script python 3 simples para o AWS Lambda. O script é simples: ele reúne informações de uma dúzia de objetos S3 e retorna os resultados.

O script usado multiprocessing.Poolpara reunir todos os arquivos em paralelo. Embora multiprocessingnão possa ser usado em um ambiente AWS Lambda, pois /dev/shmestá ausente. Então, pensei, em vez de escrever um sujo multiprocessing.Process/ multiprocessing.Queuesubstituto, tentaria asyncio.

Estou usando a versão mais recente aioboto3(8.0.5) no Python 3.8.

Meu problema é que não consigo obter nenhuma melhora entre um download sequencial ingênuo dos arquivos e um loop de evento de asyncio multiplexando os downloads.

Aqui estão as duas versões do meu 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())

O tempo para ambos run_sequential()e run_concurrent()é bastante semelhante (~ 3 segundos para uma dúzia de arquivos de 10 MB). Estou convencido de que a versão simultânea não é, por vários motivos:

  • Tentei alternar para Process/ThreadPoolExecutore os processos / threads gerados durante a função, embora eles não estejam fazendo nada
  • O tempo entre sequencial e simultâneo é muito próximo do mesmo, embora minha interface de rede definitivamente não esteja saturada e a CPU também não esteja vinculada
  • O tempo gasto pela versão simultânea aumenta linearmente com o número de arquivos.

Tenho certeza de que algo está faltando, mas simplesmente não consigo entender o quê.

Alguma ideia?

Respostas

NewbiZ Aug 28 2020 at 11:50

Depois de perder algumas horas tentando entender como usar aioboto3corretamente, decidi apenas mudar para minha solução de backup. Acabei lançando minha própria versão ingênua do multiprocessing.Poolpara uso em um ambiente lambda da AWS.

Se alguém topar com este tópico no futuro, aqui está. Está longe de ser perfeito, mas é fácil de substituir multiprocessing.Poolcomo está em meus 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