AI Engineering Academy · Lektion

Batchbearbetning med async och köer

Bygg en asynkron extraktionspipeline med asyncio och en job queue för att parallellt bearbeta tusentals dokument samtidigt som ni respekterar hastighetsbegränsningar och följer upp förloppet.

Lektion 3 av 413 steg

Batchbearbetning med async och köer är en gratis lektion i AI Engineering Academy på CoddyKit. Detta är lektion 3 av 4. Ni kan läsa hela lektionen gratis nedan och sedan öva praktiskt i webbläsaren med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt. Den ingår i lärvägen för AI Engineering Academy, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i AI Engineering Academy innehåller totalt 4 lektioner.

Varför batchbearbetning är viktig

Att bearbeta tusentals dokument ett i taget går för långsamt i produktion. En synkron loop som anropar OpenAI-API:t sekventiellt kanske bearbetar ett dokument per sekund, vilket innebär att 10 000 dokument tar nästan 3 timmar. Asynkron batchbearbetning kan parallellisera hundratals förfrågningar samtidigt och minska den totala faktiska körtiden med en storleksordning.

Grunderna i asyncio

Pythons händelseloop asyncio låter er köra många I/O-bundna uppgifter samtidigt utan trådar. När ett API-anrop väntar på ett nätverkssvar växlar händelseloopen till att bearbeta ett annat. Ni skriver kod med nyckelorden async def och await, och körmiljön hanterar schemaläggningen. Detta är idealiskt för LLM-anrop, som tillbringar större delen av sin tid med att vänta på servern.

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ör flera extraheringar med gather

asyncio.gather kör en lista med coroutines samtidigt och returnerar alla resultat när den sista är klar. För en liten grupp dokument räcker detta. Lägg era extraktions-coroutines i en list comprehension och skicka dem till gather. Den totala tiden motsvarar ungefär tiden för det långsammaste enskilda anropet, inte summan av alla anrop.

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

Styra samtidighet med semaforer

Om ni skickar tusentals anrop samtidigt kommer ni att nå hastighetsbegränsningar och få 429-fel. Använd asyncio.Semaphore för att begränsa antalet samtidiga API-anrop. En semafor med värdet 50 innebär att högst 50 anrop kan vara aktiva samtidigt. Anpassa detta antal utifrån hastighetsgränsen för er OpenAI-nivå och målmodell.

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)

Exponentiell backoff vid hastighetsbegränsningsfel

Även med en semafor kan ni nå hastighetsbegränsningar under trafiktoppar. Implementera exponentiell backoff: vänta 1 sekund, sedan 2, 4 och 8 sekunder innan ni försöker igen. Lägg till jitter (en liten slumpmässig förskjutning) för att förhindra att alla samtidiga anrop försöker igen exakt samtidigt, vilket annars skulle orsaka ännu en trafikökning. Biblioteket tenacity gör detta 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)

Använda en jobbköt för stora batchar

För batchar som är större än några tusen objekt är en beständig jobbkö bättre än asyncio.gather. Köer som Redis Queue (RQ), Celery och Dramatiq bevarar jobb vid omstarter, möjliggör horisontell skalning med flera workers och ger insyn i jobbstatus och fel. Workers hämtar jobb från kön och anropar API:et oberoende av varandra.

# 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ölja förloppet med en databas

Långkörande batchjobb behöver förloppsspårning så att ni kan övervaka status, identifiera jobb som har fastnat och återuppta körningen efter fel. Använd en statustabell i databasen med fält för dokument-ID, status (pending, processing, completed, failed), tidsstämplar och felmeddelanden. Uppdatera statusen atomiskt runt varje extraktionsanrop.

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
            )

Återuppta misslyckade jobb

Ett batchjobb bör kunna startas om på ett säkert sätt. Fråga vid uppstart databasen efter dokument med status pending eller failed och försök igen med dem. Använd en idempotensnyckel per dokument, så att ett andra försök upptäcker det färdiga resultatet och hoppar över ny bearbetning om samma dokument råkat läggas i kön två gånger. Det förhindrar dubbla skrivningar till nedströmsystem.

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)

Batchbearbetning med OpenAI Batch API

OpenAI:s Batch API låter er skicka upp till 50 000 anrop i en enda fil och ta emot resultaten inom 24 timmar till 50 % lägre kostnad. Detta passar utmärkt för icke-brådskande extraktionspipelines där kostnaden är viktigare än fördröjningen. Ni laddar upp en JSONL-fil med anrop, frågar regelbundet efter status och laddar ner 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)

Övervaka genomströmning och kostnad

Följ extraktionens genomströmning (dokument per minut) och kostnaden per dokument under batchkörningar. Dela de totala API-utgifterna med antalet behandlade dokument för att få en kostnadsbaslinje. När ni skalar bör ni se upp för linjär kostnadsökning — superlinjär ökning tyder på att tokens slösas på onödigt långa prompts. En enkel mätpanel hjälper er att upptäcka ineffektivitet innan den växer.

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

Följa hastighetsgränser för tokens per minut

OpenAI:s hastighetsgränser gäller både anrop per minut (RPM) och tokens per minut (TPM). Att skicka 50 samtidiga anrop fungerar bra för RPM, men om varje anrop använder 2 000 tokens motsvarar 50 anrop 100 000 tokens per minut — vilket enkelt överskrider gränserna för nivå 1. Räkna ut det förväntade antalet tokens innan ni skickar anropen med hjälp av tiktoken och implementera en tokenbudget parallellt med samtidighetssemaforen.

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_limit

Snabbtest

Testa er förståelse av asynkron batchbearbetning för dokumentextraktion.

Sammanfattning av lektionen

I den här lektionen lärde ni er att asyncio och Semaphore möjliggör samtidiga API-anrop samtidigt som hastighetsgränser respekteras, att job queues och status tables gör stora batchjobb möjliga att återuppta och övervaka, samt att OpenAI Batch API ger 50 % lägre kostnad för icke-brådskande arbetsbelastningar på bekostnad av 24 timmars fördröjning. Härnäst hanterar vi schemautveckling i långkörande extraktionspipelines.

Gratis att börja

Lär dig Python med en AI-lärare – gratis

Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.

Kurser
30
Lektioner
120

Vanliga frågor

Är lektionen ”Batchbearbetning med async och köer” gratis?

Ja – hela texten till ”Batchbearbetning med async och köer” kan läsas gratis här på webben. Om Ni vill öva interaktivt med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt och låsa upp resten av kursen i AI Engineering Academy, kan Ni uppgradera till CoddyKit PRO. Kursen i AI Engineering Academy innehåller totalt 4 lektioner.

Vad lär jag mig i ”Batchbearbetning med async och köer”?

Bygg en asynkron extraktionspipeline med asyncio och en job queue för att parallellt bearbeta tusentals dokument samtidigt som ni respekterar hastighetsbegränsningar och följer upp förloppet. Ni övar på AI Engineering Academy med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.

Behöver jag någon erfarenhet för att börja lära mig AI Engineering Academy?

Du behöver inga förkunskaper. Utbildningen i AI Engineering Academy på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 3 av 4.

Hur lång tid tar lektionen ”Batchbearbetning med async och köer”?

De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.

Kan jag skriva och köra kod i den här AI Engineering Academy-lektionen?

Ja. Varje AI Engineering Academy-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.

Alla lektioner i den här kursen

  1. Instructor: typad extraktion med Pydantic
  2. Hantera partiella och saknade data
  3. Batchbearbetning med async och köer
  4. Schemautveckling och bakåtkompatibilitet
← Tillbaka till AI Engineering Academy