Python 3 asyncio mit aioboto3 scheint sequentiell zu sein
Ich portiere ein einfaches Python 3-Skript nach AWS Lambda. Das Skript ist einfach: Es sammelt Informationen aus einem Dutzend S3-Objekten und gibt die Ergebnisse zurück.
Das Skript, multiprocessing.Poolmit dem alle Dateien parallel erfasst wurden. Obwohl multiprocessingnicht in einer AWS Lambda - Umgebung verwendet werden , da /dev/shmfehlt. Also dachte ich, anstatt einen Dirty multiprocessing.Process/ multiprocessing.QueueErsatz zu schreiben , würde ich es asynciostattdessen versuchen .
Ich verwende die neueste Version von aioboto3(8.0.5) unter Python 3.8.
Mein Problem ist, dass ich zwischen einem naiven sequentiellen Download der Dateien und einer asynchronen Ereignisschleife, die die Downloads multiplext, keine Verbesserung zu erzielen scheint.
Hier sind die beiden Versionen meines Codes.
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())
Das Timing für beide run_sequential()und run_concurrent()ist ziemlich ähnlich (~ 3 Sekunden für ein Dutzend 10-MB-Dateien). Ich bin überzeugt, dass die gleichzeitige Version aus mehreren Gründen nicht funktioniert:
- Ich habe versucht, zu zu wechseln
Process/ThreadPoolExecutor, und ich habe die Prozesse / Threads für die Dauer der Funktion erzeugt, obwohl sie nichts tun - Das Timing zwischen sequentiell und gleichzeitig ist sehr ähnlich, obwohl meine Netzwerkschnittstelle definitiv nicht gesättigt ist und die CPU auch nicht gebunden ist
- Die von der gleichzeitigen Version benötigte Zeit nimmt linear mit der Anzahl der Dateien zu.
Ich bin mir sicher, dass etwas fehlt, aber ich kann meinen Kopf einfach nicht um was wickeln.
Irgendwelche Ideen?
Antworten
Nachdem ich einige Stunden verloren hatte, um zu verstehen, wie man aioboto3richtig verwendet, entschied ich mich, einfach zu meiner Backup-Lösung zu wechseln. Am Ende habe ich meine eigene naive Version multiprocessing.Poolfür die Verwendung in einer AWS-Lambda-Umgebung gerollt.
Wenn jemand in Zukunft über diesen Thread stolpert, ist er hier. Es ist alles andere multiprocessing.Poolals perfekt, aber für meine einfachen Fälle einfach zu ersetzen .
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