Aioboto3 ile Python 3 asyncio sıralı görünüyor
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
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