Aioboto3 ile Python 3 asyncio sıralı görünüyor

Aug 28 2020

AWS Lambda'ya basit bir python 3 betiği taşıyorum. Komut dosyası basittir: bir düzine S3 nesnesinden bilgi toplar ve sonuçları döndürür.

multiprocessing.PoolTüm dosyaları paralel olarak toplamak için kullanılan komut dosyası . Gerçi multiprocessingbu yana AWS Lambda ortamında kullanılamaz /dev/shmeksik. Bu yüzden bir kirli multiprocessing.Process/ multiprocessing.Queuedeğiştirme yazmak yerine, asynciobunun yerine deneyeceğimi düşündüm .

aioboto3Python 3.8'de (8.0.5) ' in en son sürümünü kullanıyorum .

Benim sorunum, dosyaların saf bir sıralı indirilmesi ile indirmeleri çoğullayan bir asyncio olay döngüsü arasında herhangi bir gelişme elde edemememdir.

İşte kodumun iki versiyonu.

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())

Her ikisi için zamanlama run_sequential()ve run_concurrent()oldukça benzer (bir düzine 10MB dosya için ~ 3 saniye). Eşzamanlı sürümün birden fazla nedenden dolayı olmadığına ikna oldum:

  • 'E geçmeyi denedim Process/ThreadPoolExecutorve hiçbir şey yapmıyor olsalar da, işlevin süresi boyunca süreçler / iş parçacıkları ortaya çıktı
  • Sıralı ve eşzamanlı arasındaki zamanlama birbirine çok yakın, ancak ağ arayüzüm kesinlikle doymamış ve CPU da bağlı değil
  • Eşzamanlı sürümün harcadığı süre, dosya sayısı ile doğrusal olarak artar.

Eminim bir şeyler eksik, ama kafamı neyin etrafına dolayamıyorum.

Herhangi bir fikir?

Yanıtlar

NewbiZ Aug 28 2020 at 11:50

Nasıl doğru kullanılacağını anlamaya çalışırken birkaç saat kaybettikten sonra aioboto3, yedekleme çözümüme geçmeye karar verdim. multiprocessing.PoolBir AWS lambda ortamında kullanım için kendi saf sürümümü kullanıma sunmaya son verdim .

Gelecekte birisi bu konuya rastlarsa, işte burada. Mükemmel olmaktan uzak, ancak multiprocessing.Poolbasit vakalarım için olduğu gibi değiştirilebilecek kadar kolay.

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