Python 3 asyncio mit aioboto3 scheint sequentiell zu sein

Aug 28 2020

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

NewbiZ Aug 28 2020 at 11:50

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