Python 3 asyncio avec aioboto3 semble séquentiel
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
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