0Pricing
AI Engineering Academy · Lektion

Batch-Verarbeitung mit Async und Queues

Erstellen Sie eine asynchrone Extraktionspipeline mit asyncio und einer Job-Queue, um Tausende Dokumente parallel zu verarbeiten, dabei Rate-Limits einzuhalten und den Fortschritt zu verfolgen.

Batch-Verarbeitung mit Async und Queues ist eine kostenlose AI Engineering Academy-Lektion auf CoddyKit. Dies ist Lektion 3 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des AI Engineering Academy-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der AI Engineering Academy-Kurs umfasst insgesamt 4 Lektionen.

Warum Batch-Verarbeitung wichtig ist

Tausende Dokumente einzeln zu verarbeiten, ist für den Produktionseinsatz zu langsam. Eine synchrone Schleife, die die OpenAI API sequenziell aufruft, verarbeitet möglicherweise ein Dokument pro Sekunde. Das bedeutet, dass 10.000 Dokumente fast 3 Stunden benötigen. Asynchrone Batch-Verarbeitung kann Hunderte Anfragen gleichzeitig parallelisieren und dadurch die gesamte verstrichene Zeit um eine Größenordnung reduzieren.

Die Grundlage von asyncio

Der asyncio-Event-Loop von Python ermöglicht es Ihnen, viele I/O-gebundene Aufgaben gleichzeitig ohne Threads auszuführen. Wenn ein API-Aufruf auf eine Netzwerkantwort wartet, wechselt der Event-Loop zur Verarbeitung einer anderen Aufgabe. Sie schreiben Code mit den Schlüsselwörtern async def und await, während die Laufzeitumgebung die Planung übernimmt. Das ist ideal für LLM-Aufrufe, die den größten Teil ihrer Zeit mit dem Warten auf den Server verbringen.

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

Mehrere Extraktionen mit gather ausführen

asyncio.gather führt eine Liste von Coroutinen gleichzeitig aus und gibt alle Ergebnisse zurück, sobald die letzte abgeschlossen ist. Für eine kleine Gruppe von Dokumenten reicht das aus. Verpacken Sie Ihre Extraktionscoroutinen in einer List Comprehension und übergeben Sie sie an gather. Die gesamte Dauer entspricht ungefähr der langsamsten Einzelanfrage, nicht der Summe aller Anfragen.

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

Gleichzeitigkeit mit Semaphoren steuern

Wenn Sie Tausende Anfragen gleichzeitig senden, stoßen Sie an Rate Limits und verursachen 429-Fehler. Verwenden Sie asyncio.Semaphore, um die Anzahl gleichzeitiger API-Aufrufe zu begrenzen. Ein Semaphor mit dem Wert 50 bedeutet, dass höchstens 50 Aufrufe gleichzeitig aktiv sind. Stimmen Sie diesen Wert auf das Rate Limit Ihres OpenAI-Tarifs für das Zielmodell ab.

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)

Exponentielles Backoff bei Rate-Limit-Fehlern

Auch mit einem Semaphor können bei Verkehrsspitzen Rate Limits erreicht werden. Implementieren Sie ein exponentielles Backoff: Warten Sie vor einem erneuten Versuch 1 Sekunde, dann 2, dann 4 und anschließend 8 Sekunden. Fügen Sie Jitter hinzu, also eine kleine zufällige Abweichung, damit nicht alle gleichzeitigen Aufrufer exakt zur selben Zeit einen neuen Versuch starten und dadurch eine weitere Spitze verursachen. Die Bibliothek tenacity macht dies einfach.

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)

Eine Job-Warteschlange für große Batches verwenden

Für Batches mit mehr als einigen Tausend Elementen ist eine persistente Job-Warteschlange besser geeignet als asyncio.gather. Warteschlangen wie Redis Queue (RQ), Celery oder Dramatiq speichern Jobs über Neustarts hinweg, ermöglichen die horizontale Skalierung mit mehreren Workern und bieten Einblick in den Status und die Fehler von Jobs. Worker holen Jobs aus der Warteschlange und rufen die API unabhängig voneinander auf.

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

Fortschritt mit einer Datenbank verfolgen

Lang laufende Batch-Jobs benötigen eine Fortschrittsverfolgung, damit Sie den Status überwachen, festgefahrene Jobs erkennen und die Verarbeitung nach Fehlern fortsetzen können. Verwenden Sie in Ihrer Datenbank eine Statustabelle mit Feldern für Dokument-ID, Status (pending, processing, completed, failed), Zeitstempel und Fehlermeldungen. Aktualisieren Sie den Status rund um jeden Extraktionsaufruf atomar.

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
            )

Fehlgeschlagene Jobs fortsetzen

Ein Batch-Job sollte sicher neu gestartet werden können. Fragen Sie beim Start in der Datenbank Dokumente mit dem Status pending oder failed ab und versuchen Sie diese erneut. Verwenden Sie pro Dokument einen Idempotenzschlüssel. Wenn dasselbe Dokument versehentlich zweimal in die Warteschlange gestellt wurde, erkennt der zweite Versuch das bereits vorhandene Ergebnis und überspringt die erneute Verarbeitung. So verhindern Sie doppelte Schreibvorgänge in nachgelagerten Systemen.

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)

Batch-Verarbeitung mit der OpenAI Batch API

Mit der Batch API von OpenAI können Sie bis zu 50.000 Anfragen in einer einzelnen Datei übermitteln und innerhalb von 24 Stunden Ergebnisse mit einem Rabatt von 50 % erhalten. Das ist ideal für nicht dringende Extraktionspipelines, bei denen die Kosten wichtiger sind als die Latenz. Sie laden eine JSONL-Datei mit Anfragen hoch, fragen den Abschluss regelmäßig ab und laden die Ergebnisdatei herunter.

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)

Durchsatz und Kosten überwachen

Verfolgen Sie während der Batch-Verarbeitung den Extraktionsdurchsatz (Dokumente pro Minute) und die Kosten pro Dokument. Teilen Sie die gesamten API-Ausgaben durch die Anzahl der verarbeiteten Dokumente, um eine Kostenbasis zu erhalten. Achten Sie bei der Skalierung auf ein lineares Kostenwachstum – überproportionales Wachstum deutet darauf hin, dass Sie Tokens für unnötig lange Prompts verschwenden. Ein einfaches Metrik-Dashboard hilft Ihnen, Ineffizienzen zu erkennen, bevor sie sich verstärken.

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

Token-pro-Minute-Rate-Limits einhalten

Die Rate Limits von OpenAI gelten sowohl für Anfragen pro Minute (RPM) als auch für Tokens pro Minute (TPM). 50 gleichzeitige Aufrufe sind im Hinblick auf RPM unproblematisch. Wenn jeder Aufruf jedoch 2.000 Tokens verwendet, entsprechen 50 Aufrufe 100.000 Tokens pro Minute – damit werden die Limits von Tier 1 leicht überschritten. Zählen Sie die erwartete Token-Anzahl vor dem Absenden mithilfe von tiktoken und implementieren Sie zusätzlich zu Ihrem Gleichzeitigkeitsemaphor ein Token-Budget.

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

Kurzer Test

Testen Sie Ihr Verständnis der asynchronen Batch-Verarbeitung für die Dokumentextraktion.

Zusammenfassung der Lektion

In dieser Lektion haben Sie Folgendes gelernt: asyncio und Semaphore ermöglichen gleichzeitige API-Aufrufe unter Einhaltung der Rate Limits, Job-Warteschlangen und Statustabellen machen große Batch-Jobs fortsetzbar und transparent, und die OpenAI Batch API bietet bei nicht dringenden Workloads eine Kostenersparnis von 50 %, allerdings auf Kosten einer Latenz von 24 Stunden. Als Nächstes behandeln wir die Schemaentwicklung in langfristig betriebenen Extraktionspipelines.

Häufig gestellte Fragen

Ist die Lektion „Batch-Verarbeitung mit Async und Queues“ kostenlos?

Ja — der vollständige Text von „Batch-Verarbeitung mit Async und Queues“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des AI Engineering Academy-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der AI Engineering Academy-Kurs umfasst insgesamt 4 Lektionen.

Was lerne ich in „Batch-Verarbeitung mit Async und Queues“?

Erstellen Sie eine asynchrone Extraktionspipeline mit asyncio und einer Job-Queue, um Tausende Dokumente parallel zu verarbeiten, dabei Rate-Limits einzuhalten und den Fortschritt zu verfolgen. Du übst AI Engineering Academy mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.

Brauche ich Erfahrung, um AI Engineering Academy zu starten?

Keine Vorkenntnisse erforderlich. AI Engineering Academy auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 3 von 4.

Wie lange dauert die Lektion „Batch-Verarbeitung mit Async und Queues“?

Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.

Kann ich in dieser AI Engineering Academy-Lektion Code schreiben und ausführen?

Ja. Jede AI Engineering Academy-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.

Alle Lektionen in diesem Kurs

  1. Instructor: Typisierte Extraktion mit Pydantic
  2. Unvollständige und fehlende Daten verarbeiten
  3. Batch-Verarbeitung mit Async und Queues
  4. Schema-Weiterentwicklung und Abwärtskompatibilität
← Zurück zu AI Engineering Academy