Przetwarzanie wsadowe z async i kolejkami
Zbuduj asynchroniczny potok ekstrakcji z użyciem asyncio i kolejki zadań, aby równolegle przetwarzać tysiące dokumentów z zachowaniem limitów szybkości i śledzeniem postępu.
Przetwarzanie wsadowe z async i kolejkami to bezpłatna lekcja AI Engineering Academy na CoddyKit. To lekcja 3 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AI Engineering Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AI Engineering Academy zawiera 4 lekcji w sumie.
Dlaczego przetwarzanie wsadowe ma znaczenie
Przetwarzanie tysięcy dokumentów pojedynczo jest zbyt wolne w środowisku produkcyjnym. Synchroniczna pętla, która sekwencyjnie wywołuje API OpenAI, może przetwarzać 1 dokument na sekundę, co oznacza, że obsługa 10 000 dokumentów zajmie niemal 3 godziny. Asynchroniczne przetwarzanie wsadowe może równolegle obsługiwać setki żądań, skracając całkowity czas zegarowy nawet o rząd wielkości.
Podstawy asyncio
Pętla zdarzeń asyncio w Pythonie umożliwia równoczesne uruchamianie wielu zadań związanych z operacjami wejścia-wyjścia bez użycia wątków. Gdy jedno wywołanie API oczekuje na odpowiedź sieciową, pętla zdarzeń przełącza się na przetwarzanie innego. Kod jest pisany za pomocą słów kluczowych async def i await, a środowisko wykonawcze zajmuje się planowaniem. Jest to idealne rozwiązanie dla wywołań LLM, które przez większość czasu oczekują na odpowiedź serwera.
import asyncio
import instructor
from openai import AsyncOpenAI
async_client = instructor.from_openai(AsyncOpenAI())
async def extract_one(text: str) -> PersonExtract:
return await async_client.chat.completions.create(
model='gpt-4o-mini',
response_model=PersonExtract,
messages=[{'role': 'user', 'content': text}]
)Uruchamianie wielu operacji wydobywania danych za pomocą gather
asyncio.gather uruchamia listę korutyn współbieżnie i zwraca wszystkie wyniki po zakończeniu ostatniej z nich. W przypadku niewielkiej partii dokumentów to wystarczy. Umieść korutyny ekstrakcji w wyrażeniu listowym i przekaż je do gather. Łączny czas jest w przybliżeniu równy czasowi najwolniejszego pojedynczego wywołania, a nie sumie czasów wszystkich wywołań.
async def batch_extract(texts: list) -> list:
tasks = [extract_one(text) for text in texts]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Filter out exceptions
successes = [r for r in results if not isinstance(r, Exception)]
failures = [r for r in results if isinstance(r, Exception)]
print(f'Success: {len(successes)}, Failures: {len(failures)}')
return successes
results = asyncio.run(batch_extract(documents))Kontrolowanie współbieżności za pomocą semaforów
Jednoczesne wysłanie tysięcy żądań spowoduje przekroczenie limitów zapytań i błędy 429. Użyj asyncio.Semaphore, aby ograniczyć liczbę współbieżnych wywołań API. Semafor o wartości 50 oznacza, że jednocześnie może być aktywnych najwyżej 50 wywołań. Dobierz tę wartość na podstawie limitu zapytań w Państwa warstwie OpenAI dla docelowego modelu.
import asyncio
sem = asyncio.Semaphore(50) # max 50 concurrent calls
async def extract_with_limit(text: str, semaphore: asyncio.Semaphore):
async with semaphore:
return await extract_one(text)
async def batch_extract_limited(texts: list):
tasks = [extract_with_limit(t, sem) for t in texts]
return await asyncio.gather(*tasks, return_exceptions=True)Wykładnicze zwiększanie opóźnienia po błędach limitu zapytań
Nawet w przypadku użycia semafora podczas nagłych wzrostów ruchu może dojść do przekroczenia limitów zapytań. Należy zaimplementować wykładnicze zwiększanie opóźnienia: przed ponowieniem próby odczekać 1 sekundę, następnie 2, 4 i 8 sekund. Należy dodać jitter (niewielkie losowe przesunięcie), aby zapobiec ponownym próbom wszystkich współbieżnych wywołań dokładnie w tym samym momencie, co spowodowałoby kolejny nagły wzrost ruchu. Biblioteka tenacity ułatwia to zadanie.
from tenacity import retry, wait_exponential, stop_after_attempt, retry_if_exception_type
from openai import RateLimitError
@retry(
wait=wait_exponential(multiplier=1, min=1, max=60),
stop=stop_after_attempt(5),
retry=retry_if_exception_type(RateLimitError)
)
async def extract_with_retry(text: str):
return await extract_one(text)Korzystanie z kolejki zadań dla dużych partii
W przypadku partii liczących więcej niż kilka tysięcy elementów trwała kolejka zadań sprawdza się lepiej niż asyncio.gather. Kolejki takie jak Redis Queue (RQ), Celery czy Dramatiq zachowują zadania po ponownym uruchomieniu, umożliwiają horyzontalne skalowanie za pomocą wielu workerów oraz zapewniają wgląd w stan i błędy zadań. Workery pobierają zadania z kolejki i niezależnie wywołują API.
# With Redis Queue (RQ)
from rq import Queue
from redis import Redis
redis_conn = Redis()
q = Queue('extractions', connection=redis_conn)
def enqueue_documents(doc_ids: list):
for doc_id in doc_ids:
q.enqueue(
'workers.extract_document',
doc_id,
job_timeout=120,
result_ttl=3600
)
enqueue_documents(all_doc_ids)Śledzenie postępu za pomocą bazy danych
Długotrwałe zadania przetwarzania partii wymagają śledzenia postępu, aby można było monitorować stan, identyfikować zablokowane zadania i wznawiać pracę po awariach. Należy użyć w bazie danych tabeli stanów zawierającej pola na identyfikator dokumentu, stan (pending, processing, completed, failed), znaczniki czasu i komunikaty o błędach. Stan należy aktualizować atomowo wokół każdego wywołania ekstrakcji.
import asyncpg
async def process_document(pool, doc_id: str, text: str):
async with pool.acquire() as conn:
await conn.execute(
'UPDATE extractions SET status=$1, started_at=NOW() WHERE doc_id=$2',
'processing', doc_id
)
try:
result = await extract_one(text)
await conn.execute(
'UPDATE extractions SET status=$1, result=$2, completed_at=NOW() WHERE doc_id=$3',
'completed', result.model_dump_json(), doc_id
)
except Exception as e:
await conn.execute(
'UPDATE extractions SET status=$1, error=$2 WHERE doc_id=$3',
'failed', str(e), doc_id
)Wznawianie nieudanych zadań
Zadanie przetwarzania partii powinno umożliwiać bezpieczne ponowne uruchomienie. Przy starcie należy wyszukać w bazie dokumenty ze stanem pending lub failed i ponowić ich przetwarzanie. Należy użyć klucza idempotencji dla każdego dokumentu, aby w przypadku przypadkowego dwukrotnego dodania tego samego dokumentu do kolejki druga próba wykryła ukończony wynik i pominęła ponowne przetwarzanie. Zapobiega to zduplikowanym zapisom w systemach downstream.
async def get_pending_docs(pool) -> list:
async with pool.acquire() as conn:
rows = await conn.fetch(
'SELECT doc_id, raw_text FROM extractions WHERE status IN ($1, $2)',
'pending', 'failed'
)
return [dict(row) for row in rows]
async def resume_batch(pool):
docs = await get_pending_docs(pool)
print(f'Resuming {len(docs)} unprocessed documents')
tasks = [process_document(pool, d['doc_id'], d['raw_text']) for d in docs]
await asyncio.gather(*tasks, return_exceptions=True)Przetwarzanie wsadowe za pomocą OpenAI Batch API
OpenAI Batch API umożliwia przesłanie do 50 000 żądań w jednym pliku i otrzymanie wyników w ciągu 24 godzin przy 50% rabacie. To idealne rozwiązanie dla niepilnych potoków ekstrakcji, w których koszt jest ważniejszy niż opóźnienie. Należy przesłać plik JSONL z żądaniami, odpytywać usługę o ukończenie przetwarzania i pobrać plik wyników.
from openai import OpenAI
import json
client = OpenAI()
# Build JSONL batch file
with open('/tmp/batch_requests.jsonl', 'w') as f:
for i, text in enumerate(documents):
request = {
'custom_id': f'doc_{i}',
'method': 'POST',
'url': '/v1/chat/completions',
'body': {
'model': 'gpt-4o-mini',
'messages': [{'role': 'user', 'content': text}]
}
}
f.write(json.dumps(request) + '\n')
# Upload and submit
batch_file = client.files.create(file=open('/tmp/batch_requests.jsonl', 'rb'), purpose='batch')
batch = client.batches.create(input_file_id=batch_file.id, endpoint='/v1/chat/completions', completion_window='24h')
print(batch.id)Monitorowanie przepustowości i kosztów
Podczas przetwarzania partii należy śledzić przepustowość ekstrakcji (dokumenty na minutę) oraz koszt przypadający na dokument. Podzielenie całkowitych wydatków na API przez liczbę przetworzonych dokumentów pozwala wyznaczyć koszt bazowy. W miarę skalowania należy szukać liniowego wzrostu kosztów — wzrost ponadliniowy sugeruje marnowanie tokenów na niepotrzebnie długie prompty. Prosty dashboard metryk pomaga wykryć nieefektywności, zanim zaczną się kumulować.
import time
class BatchMetrics:
def __init__(self):
self.start_time = time.time()
self.processed = 0
self.total_tokens = 0
self.cost = 0.0
def record(self, usage):
self.processed += 1
self.total_tokens += usage.total_tokens
self.cost += usage.prompt_tokens * 0.00000015 + usage.completion_tokens * 0.0000006
def report(self):
elapsed = time.time() - self.start_time
print(f'{self.processed} docs in {elapsed:.1f}s = {self.processed/elapsed:.1f} docs/sec')
print(f'Cost: ${self.cost:.4f} = ${self.cost/self.processed:.6f} per doc')Przestrzeganie limitów tokenów na minutę
Limity zapytań OpenAI dotyczą zarówno zapytań na minutę (RPM), jak i tokenów na minutę (TPM). Wysłanie 50 współbieżnych wywołań jest zgodne z limitem RPM, ale jeśli każde wywołanie zużywa 2000 tokenów, 50 wywołań oznacza 100 000 tokenów na minutę — co z łatwością przekracza limity poziomu 1. Przed wysłaniem należy oszacować liczbę tokenów za pomocą tiktoken i wdrożyć budżet tokenów obok semafora współbieżności.
import tiktoken
enc = tiktoken.encoding_for_model('gpt-4o-mini')
def estimate_tokens(text: str) -> int:
return len(enc.encode(text)) + 300 # +300 for schema + response
# TPM_LIMIT = 200_000 # Tier 2 limit
# Only submit a batch if estimated total tokens fits within budget
def fits_in_budget(texts: list, tpm_limit: int = 200_000) -> bool:
total = sum(estimate_tokens(t) for t in texts)
return total <= tpm_limitSzybkie sprawdzenie
Sprawdź swoją wiedzę na temat asynchronicznego przetwarzania wsadowego na potrzeby ekstrakcji dokumentów.
Podsumowanie lekcji
W tej lekcji dowiedzieli się Państwo, że: asyncio i Semaphore umożliwiają współbieżne wywołania API przy przestrzeganiu limitów zapytań, kolejki zadań i tabele stanów sprawiają, że duże zadania wsadowe można wznawiać i monitorować, a OpenAI Batch API zapewnia 50% oszczędności kosztów dla niepilnych obciążeń kosztem 24-godzinnego opóźnienia. W następnej części zajmiemy się ewolucją schematów w długotrwałych potokach ekstrakcji.
Często zadawane pytania
Czy lekcja „Przetwarzanie wsadowe z async i kolejkami” jest bezpłatna?
Tak — pełny tekst „Przetwarzanie wsadowe z async i kolejkami” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AI Engineering Academy, przejdź na CoddyKit PRO. Kurs AI Engineering Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „Przetwarzanie wsadowe z async i kolejkami”?
Zbuduj asynchroniczny potok ekstrakcji z użyciem asyncio i kolejki zadań, aby równolegle przetwarzać tysiące dokumentów z zachowaniem limitów szybkości i śledzeniem postępu. Ćwiczysz AI Engineering Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć AI Engineering Academy?
Nie wymagamy żadnego doświadczenia. AI Engineering Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 3 z 4.
Ile czasu zajmuje lekcja „Przetwarzanie wsadowe z async i kolejkami”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji AI Engineering Academy?
Tak. Każda lekcja AI Engineering Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Instructor: ekstrakcja typowana z Pydantic
- Obsługa częściowych i brakujących danych
- Przetwarzanie wsadowe z async i kolejkami
- Ewolucja schematów i zgodność wsteczna