Python 3 assíncio com aioboto3 parece sequencial
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
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