Batchbehandling med asynkronitet og køer
Bygg en asynkron uttrekkingspipeline med asyncio og en jobbkø for å behandle tusenvis av dokumenter parallelt, samtidig som du overholder hastighetsgrenser og sporer fremdriften.
Batchbehandling med asynkronitet og køer er en gratis leksjon i AI Engineering Academy på CoddyKit. Dette er leksjon 3 av 4. Du kan lese hele leksjonen gratis nedenfor – og deretter øve praktisk i nettleseren med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt. Den er en del av læringsløpet i AI Engineering Academy, og fremdriften din synkroniseres mellom nettet og CoddyKit-appen. Kurset i AI Engineering Academy inneholder totalt 4 leksjoner.
Hvorfor satsvis behandling er viktig
Det går for sakte å behandle tusenvis av dokumenter ett om gangen i produksjon. En synkron løkke som kaller OpenAI-API-et sekvensielt, kan kanskje behandle ett dokument per sekund, noe som betyr at 10 000 dokumenter tar nesten tre timer. Asynkron satsvis behandling kan parallellisere hundrevis av forespørsler samtidig og redusere den totale klokketiden med en størrelsesorden.
Grunnlaget i asyncio
Pythons asyncio-hendelsesløkke lar Dem kjøre mange I/O-bundne oppgaver samtidig uten tråder. Når ett API-kall venter på et nettverkssvar, bytter hendelsesløkken til å behandle et annet. De skriver kode med nøkkelordene async def og await, mens kjøretidsmiljøet håndterer planleggingen. Dette passer godt for LLM-kall, som bruker mesteparten av tiden på å vente på serveren.
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}]
)Kjøre flere uthentinger med gather
asyncio.gather kjører en liste med coroutines samtidig og returnerer alle resultatene når den siste er ferdig. For en liten gruppe dokumenter er dette tilstrekkelig. Pakk extraction-coroutinene inn i en list comprehension, og send dem til gather. Den totale tiden tilsvarer omtrent tiden for det tregeste enkeltkallet, ikke summen av alle kallene.
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))Kontrollere samtidighet med semaforer
Hvis De sender tusenvis av forespørsler samtidig, vil De treffe rate limits og få 429-feil. Bruk asyncio.Semaphore til å begrense antallet samtidige API-kall. En semafor med verdien 50 betyr at maksimalt 50 kall kan være under behandling samtidig. Juster dette tallet basert på rategrensen for OpenAI-nivået Deres for målmodellen.
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)Eksponentiell backoff ved rate limit-feil
Selv med en semafor kan De treffe rate limits under trafikktopper. Implementer eksponentiell backoff: vent 1 sekund, deretter 2, så 4 og til slutt 8 sekunder før nytt forsøk. Legg til jitter (et lite tilfeldig avvik) for å hindre at alle samtidige kallere prøver på nytt nøyaktig samtidig, noe som ellers ville skapt en ny trafikktopp. Biblioteket tenacity gjør dette enkelt.
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)Bruke en jobbkø for store grupper
For grupper på mer enn noen få tusen elementer er en persistent job queue bedre enn asyncio.gather. Køer som Redis Queue (RQ), Celery eller Dramatiq bevarer jobber gjennom omstarter, gjør horisontal skalering med flere arbeidere mulig og gir innsyn i jobbstatus og feil. Arbeiderne henter jobber fra køen og kaller API-et uavhengig av hverandre.
# 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)Følge fremdriften med en database
Langvarige gruppejobber trenger fremdriftssporing slik at De kan overvåke status, finne jobber som har stoppet opp, og fortsette etter feil. Bruk en status table i databasen med felt for dokument-ID, status (pending, processing, completed, failed), tidsstempler og feilmeldinger. Oppdater statusen atomisk rundt hvert extraction-kall.
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
)Fortsette mislykkede jobber
En gruppejobb bør kunne startes på nytt på en trygg måte. Ved oppstart spør De databasen etter dokumenter med status pending eller failed, og prøver dem på nytt. Bruk en idempotency key per dokument. Hvis det samme dokumentet ved et uhell legges i kø to ganger, oppdager det andre forsøket at resultatet allerede er fullført, og hopper over ny behandling. Dette hindrer duplikatskriving til nedstrømssystemer.
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)Gruppering med OpenAI Batch API
OpenAI's Batch API lar Dem sende inn opptil 50 000 forespørsler i én fil og motta resultater innen 24 timer til 50 % rabatt. Dette er ideelt for ikke-hastende extraction-pipelines der kostnad er viktigere enn forsinkelse. De laster opp en JSONL-fil med forespørsler, spør jevnlig om status og laster ned resultatfilen.
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)Overvåke gjennomstrømning og kostnader
Følg med på extraction-gjennomstrømningen (dokumenter per minutt) og kostnaden per dokument under kjøring av grupper. Del de totale API-utgiftene på antallet behandlede dokumenter for å finne et kostnadsgrunnlag. Når De skalerer opp, bør De se etter lineær kostnadsvekst — superlineær vekst tyder på at De sløser tokens på unødvendig lange prompts. Et enkelt metrics-dashboard hjelper Dem med å oppdage ineffektivitet før den forsterkes.
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')Overholde rategrenser for tokens per minutt
OpenAI-rategrenser gjelder både requests per minute (RPM) og tokens per minute (TPM). Det er greit å sende 50 samtidige kall med tanke på RPM, men hvis hvert kall bruker 2 000 tokens, tilsvarer 50 kall 100 000 tokens per minutt — noe som lett overskrider grensene for Tier 1. Tell forventede tokens før innsending ved hjelp av tiktoken, og implementer et tokenbudsjett i tillegg til semaforen for samtidighet.
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_limitKort kontroll
Test forståelsen Deres av asynkron gruppebehandling for dokumentekstraksjon.
Oppsummering av leksjonen
I denne leksjonen lærte De at asyncio og Semaphore muliggjør samtidige API-kall samtidig som rategrenser overholdes, at job queues og status tables gjør store gruppejobber mulige å gjenoppta og overvåke, og at OpenAI Batch API gir 50 % lavere kostnader for ikke-hastende arbeidsbelastninger, på bekostning av 24 timers forsinkelse. Neste tema er håndtering av skjemautvikling i langvarige extraction-pipelines.
Lær deg Python med en AI-veileder – gratis
Skriv og kjør ekte kode i nettleseren, få umiddelbar hjelp fra en AI-veileder som er tilgjengelig døgnet rundt, og fortsett der du slapp – på nettet eller i appen.
- Kurs
- 30
- Leksjoner
- 120
Ofte stilte spørsmål
Er leksjonen «Batchbehandling med asynkronitet og køer» gratis?
Ja – hele teksten i «Batchbehandling med asynkronitet og køer» er gratis å lese her på nettet. For å øve interaktivt med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt, og for å låse opp resten av AI Engineering Academy-kurset, kan du oppgradere til CoddyKit PRO. Kurset i AI Engineering Academy inneholder totalt 4 leksjoner.
Hva lærer jeg i «Batchbehandling med asynkronitet og køer»?
Bygg en asynkron uttrekkingspipeline med asyncio og en jobbkø for å behandle tusenvis av dokumenter parallelt, samtidig som du overholder hastighetsgrenser og sporer fremdriften. Du øver på AI Engineering Academy med praktisk kode som du kjører direkte i nettleseren, mens en AI-veileder som er tilgjengelig døgnet rundt, svarer på spørsmålene dine mens du jobber deg gjennom leksjonen.
Trenger jeg erfaring for å begynne med AI Engineering Academy?
Ingen tidligere erfaring er nødvendig. AI Engineering Academy på CoddyKit er lagt opp for både nybegynnere og viderekomne, så De kan begynne her eller helt fra start og lære i Deres eget tempo. Dette er leksjon 3 av 4.
Hvor lang tid tar leksjonen «Batchbehandling med asynkronitet og køer»?
De fleste CoddyKit-leksjoner tar omtrent 5–10 minutter. Hver leksjon er kort og interaktiv, slik at De gjør jevne fremskritt og kan fortsette akkurat der De slapp – både på nettet og i appen.
Kan jeg skrive og kjøre kode i denne AI Engineering Academy-leksjonen?
Ja. Alle AI Engineering Academy-leksjoner har en innebygd kodeeditor, slik at De kan skrive og kjøre ekte kode direkte i nettleseren og få umiddelbar tilbakemelding fra AI – uten lokal konfigurering.
Alle leksjonene i dette kurset
- Instructor: Typet uttrekking med Pydantic
- Håndtere delvise og manglende data
- Batchbehandling med asynkronitet og køer
- Skjemautvikling og bakoverkompatibilitet