Batchbehandling med async og køer
Opbyg en asynkron udtrækspipeline ved hjælp af asyncio og en jobkø, så tusindvis af dokumenter kan behandles parallelt, samtidig med at hastighedsgrænser overholdes, og fremskridt spores.
Batchbehandling med async og køer er en gratis AI Engineering Academy-lektion på CoddyKit. Dette er lektion 3 af 4. Du kan læse hele lektionen gratis nedenfor — og derefter øve dig praktisk i browseren med en indbygget kodeeditor og en AI-vejleder, der er tilgængelig døgnet rundt. Den er en del af læringsforløbet i AI Engineering Academy, og dine fremskridt synkroniseres på tværs af nettet og CoddyKit-appen. AI Engineering Academy-kurset indeholder 4 lektioner i alt.
Hvorfor batchbehandling er vigtig
Det er for langsomt til produktion at behandle tusindvis af dokumenter ét ad gangen. En synkron løkke, der kalder OpenAI API'et sekventielt, behandler måske 1 dokument i sekundet, hvilket betyder, at 10.000 dokumenter tager næsten 3 timer. Asynkron batchbehandling kan parallelisere hundredvis af forespørgsler samtidigt og dermed reducere den samlede klokketid med en størrelsesorden.
Grundlaget i asyncio
Pythons asyncio-hændelsesløkke lader dig køre mange I/O-bundne opgaver samtidigt uden tråde. Når et API-kald venter på et netværkssvar, skifter hændelsesløkken til at behandle et andet. Du skriver kode med nøgleordene async def og await, mens runtime-miljøet håndterer planlægningen. Det er ideelt til LLM-kald, som bruger det meste af tiden på at 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}]
)Kørsel af flere udtrækninger med gather
asyncio.gather kører en liste af coroutiner samtidigt og returnerer alle resultater, når den sidste er færdig. Til en lille gruppe dokumenter er dette tilstrækkeligt. Pak dine ekstraktionscoroutiner ind i en listeforståelse, og send dem til gather. Den samlede tid svarer omtrent til tiden for det langsomste enkeltkald, ikke summen af alle kald.
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))Styring af samtidighed med semaforer
Hvis du sender tusindvis af forespørgsler samtidigt, rammer du hastighedsbegrænsninger og får 429-fejl. Brug asyncio.Semaphore til at begrænse antallet af samtidige API-kald. En semafor med værdien 50 betyder, at højst 50 kald er i gang på samme tid. Juster dette tal efter din OpenAI-tiers hastighedsbegrænsning 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)Eksponentiel backoff ved fejl på grund af hastighedsbegrænsninger
Selv med en semafor kan du ramme hastighedsbegrænsninger under trafikspidser. Implementer eksponentiel backoff: vent 1 sekund, derefter 2, 4 og 8 sekunder, før du prøver igen. Tilføj jitter (en lille tilfældig forskydning) for at forhindre, at alle samtidige kaldere prøver igen præcis samtidig, hvilket ellers ville skabe endnu en trafikspids. Biblioteket tenacity gør dette nemt.
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)Brug af en jobkø til store grupper
Til grupper på mere end nogle få tusinde elementer er en vedvarende jobkø bedre end asyncio.gather. Køer som Redis Queue (RQ), Celery eller Dramatiq bevarer jobs efter genstarter, gør det muligt at skalere horisontalt med flere workers og giver dig overblik over jobstatus og fejl. Workers henter jobs fra køen og kalder API'et uafhængigt af hinanden.
# 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)Sporing af fremdrift med en database
Langvarige batchjobs har brug for fremdriftssporing, så du kan overvåge status, identificere jobs, der sidder fast, og genoptage efter fejl. Brug en statustabel i din database med felter til dokument-id, status (pending, processing, completed, failed), tidsstempler og fejlmeddelelser. Opdater statussen atomisk omkring hvert ekstraktionskald.
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
)Genoptagelse af fejlede jobs
Et batchjob skal kunne genstartes sikkert. Ved opstart skal du forespørge databasen efter dokumenter med statussen pending eller failed og prøve dem igen. Brug en idempotensnøgle pr. dokument, så det andet forsøg opdager det færdige resultat og springer genbehandlingen over, hvis det samme dokument ved en fejl sættes i kø to gange. Det forhindrer dublerede skrivninger til downstream-systemer.
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)Batchbehandling med OpenAI Batch API
OpenAI's Batch API lader dig indsende op til 50.000 forespørgsler i en enkelt fil og modtage resultater inden for 24 timer med 50 % rabat. Det er ideelt til ikke-akutte ekstraktionspipelines, hvor omkostninger betyder mere end latenstid. Du uploader en JSONL-fil med forespørgsler, forespørger løbende om færdiggørelse og downloader 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ågning af gennemløb og omkostninger
Spor ekstraktionsgennemløb (dokumenter pr. minut) og omkostningen pr. dokument under batchkørsler. Divider det samlede API-forbrug med antallet af behandlede dokumenter for at få et omkostningsgrundlag. Når du skalerer, skal du se efter lineær omkostningsvækst — superlineær vækst tyder på, at du spilder tokens på unødigt lange prompts. Et simpelt metrikdashboard hjælper dig med at opdage ineffektivitet, før den vokser.
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')Overholdelse af hastighedsbegrænsninger for tokens pr. minut
OpenAI's hastighedsbegrænsninger gælder både forespørgsler pr. minut (RPM) og tokens pr. minut (TPM). Det er fint at sende 50 samtidige kald i forhold til RPM, men hvis hvert kald bruger 2.000 tokens, svarer 50 kald til 100.000 tokens pr. minut — hvilket nemt overskrider begrænsningerne for Tier 1. Optæl de forventede tokens med tiktoken, før du indsender, og implementer et tokenbudget sammen med din semafor for samtidighed.
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_limitHurtigt tjek
Test din forståelse af asynkron batchbehandling til dokumentekstraktion.
Opsummering af lektionen
I denne lektion har du lært, at asyncio og Semaphore muliggør samtidige API-kald, samtidig med at hastighedsbegrænsninger overholdes, at jobkøer og statustabeller gør store batchjobs genoptagelige og overvågelige, og at OpenAI Batch API giver 50 % lavere omkostninger for ikke-akutte arbejdsbelastninger på bekostning af 24 timers latenstid. Næste gang håndterer vi skemaudvikling i langvarige ekstraktionspipelines.
Lær Python med en AI-underviser — gratis
Skriv og kør rigtig kode i din browser, få øjeblikkelig hjælp fra en AI-underviser døgnet rundt, og fortsæt, hvor du slap, på web eller i appen.
- Kurser
- 30
- Lektioner
- 120
Ofte stillede spørgsmål
Er lektionen “Batchbehandling med async og køer” gratis?
Ja — hele teksten til “Batchbehandling med async og køer” kan læses gratis her på nettet. Hvis du vil øve dig interaktivt med en indbygget kodeeditor og en AI-vejleder døgnet rundt og få adgang til resten af AI Engineering Academy-kurset, skal du opgradere til CoddyKit PRO. AI Engineering Academy-kurset indeholder 4 lektioner i alt.
Hvad lærer jeg i “Batchbehandling med async og køer”?
Opbyg en asynkron udtrækspipeline ved hjælp af asyncio og en jobkø, så tusindvis af dokumenter kan behandles parallelt, samtidig med at hastighedsgrænser overholdes, og fremskridt spores. Du øver dig i AI Engineering Academy med praktisk kode, som du kører direkte i browseren, og en AI-vejleder døgnet rundt besvarer dine spørgsmål, mens du arbejder dig gennem lektionen.
Skal jeg have erfaring for at begynde på AI Engineering Academy?
Der kræves ingen tidligere erfaring. AI Engineering Academy på CoddyKit er tilrettelagt for både begyndere og øvede, så du kan starte her eller fra begyndelsen og lære i dit eget tempo. Dette er lektion 3 af 4.
Hvor lang tid tager lektionen “Batchbehandling med async og køer”?
De fleste CoddyKit-lektioner tager cirka 5–10 minutter. Hver lektion er kort og interaktiv, så du gør løbende fremskridt og kan fortsætte, hvor du slap – på både web og app.
Kan jeg skrive og køre kode i denne AI Engineering Academy-lektion?
Ja. Alle AI Engineering Academy-lektioner har en indbygget kodeeditor, så du kan skrive og køre rigtig kode direkte i din browser og få øjeblikkelig feedback fra AI – uden lokal opsætning.
Alle lektioner i dette kursus
- Instructor: Typet udtræk med Pydantic
- Håndtering af delvise og manglende data
- Batchbehandling med async og køer
- Skemaudvikling og bagudkompatibilitet