Python 3 asyncio con aioboto3 parece secuencial
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
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