Python 3 asyncio avec aioboto3 semble séquentiel

Aug 28 2020

Je porte un simple script python 3 sur AWS Lambda. Le script est simple: il rassemble les informations d'une douzaine d'objets S3 et renvoie les résultats.

Le script utilisé multiprocessing.Poolpour rassembler tous les fichiers en parallèle. Bien multiprocessingqu'il ne puisse pas être utilisé dans un environnement AWS Lambda car il /dev/shmest manquant. Alors j'ai pensé qu'au lieu d'écrire un sale multiprocessing.Process/ multiprocessing.Queueremplacement, j'essaierais à la asyncioplace.

J'utilise la dernière version de aioboto3(8.0.5) sur Python 3.8.

Mon problème est que je ne parviens pas à gagner d'amélioration entre un téléchargement séquentiel naïf des fichiers et une boucle d'événement asyncio multiplexant les téléchargements.

Voici les deux versions de mon code.

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

Le timing pour les deux run_sequential()et run_concurrent()est assez similaire (~ 3 secondes pour une douzaine de fichiers de 10 Mo). Je suis convaincu que la version simultanée ne l'est pas, pour plusieurs raisons:

  • J'ai essayé de passer à Process/ThreadPoolExecutor, et les processus / threads sont générés pendant la durée de la fonction, bien qu'ils ne fassent rien
  • Le timing entre séquentiel et concurrent est très proche du même, bien que mon interface réseau ne soit certainement pas saturée et que le processeur ne soit pas lié non plus
  • Le temps nécessaire à la version simultanée augmente linéairement avec le nombre de fichiers.

Je suis sûr qu'il manque quelque chose, mais je ne peux tout simplement pas comprendre quoi.

Des idées?

Réponses

NewbiZ Aug 28 2020 at 11:50

Après avoir perdu quelques heures à essayer de comprendre comment l'utiliser aioboto3correctement, j'ai décidé de simplement passer à ma solution de sauvegarde. J'ai fini par lancer ma propre version naïve de multiprocessing.Poolpour une utilisation dans un environnement AWS lambda.

Si quelqu'un tombe sur ce fil à l'avenir, le voici. Il est loin d'être parfait, mais assez facile à remplacer multiprocessing.Pooltel quel pour mes cas 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