Python 3 asyncio с aioboto3 кажется последовательным
Я портирую простой скрипт Python 3 на AWS Lambda. Скрипт прост: он собирает информацию из десятка объектов S3 и возвращает результаты.
Скрипт, используемый multiprocessing.Poolдля параллельного сбора всех файлов. Хотя multiprocessingне может использоваться в среде AWS Lambda, поскольку /dev/shmотсутствует. Поэтому я подумал, что вместо того, чтобы писать грязную multiprocessing.Process/ multiprocessing.Queueзамену, я бы попробовал asyncio.
Я использую последнюю версию aioboto3(8.0.5) на Python 3.8.
Моя проблема в том, что я не могу добиться каких-либо улучшений между наивной последовательной загрузкой файлов и циклом событий asyncio, мультиплексирующим загрузки.
Вот две версии моего кода.
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())
Время для обоих run_sequential()и run_concurrent()довольно похоже (~ 3 секунды для дюжины файлов размером 10 МБ). Я убежден, что параллельной версии нет по нескольким причинам:
- Я попытался переключиться на
Process/ThreadPoolExecutor, и у меня процессы / потоки порождались на время выполнения функции, хотя они ничего не делают - Время между последовательным и параллельным очень близко к одному и тому же, хотя мой сетевой интерфейс определенно не насыщен, и ЦП тоже не привязан.
- Время, затрачиваемое параллельной версией, линейно увеличивается с количеством файлов.
Я уверен, что чего-то не хватает, но я просто не могу понять, что именно.
Есть идеи?
Ответы
Потратив несколько часов на то, чтобы понять, как aioboto3правильно пользоваться, я решил просто переключиться на свое решение для резервного копирования. Я закончил тем, что скатал свою наивную версию multiprocessing.Poolдля использования в лямбда-среде AWS.
Если кто-то наткнется на эту ветку в будущем, вот она. Он далек от совершенства, но в multiprocessing.Poolмоих простых случаях его достаточно легко заменить как есть.
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