0Pricing
AI Prompt Engineering · Leçon

Traitement par lots et exécution asynchrone

API Batch d’OpenAI, Python asynchrone et exécution concurrente des prompts.

Traitement par lots et exécution asynchrone est une leçon AI Prompt Engineering gratuite sur CoddyKit. Ceci est la leçon 2 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage AI Prompt Engineering, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours AI Prompt Engineering comprend 4 leçons au total.

Pourquoi utiliser le traitement par lots et l'asynchronisme ?

Traiter séquentiellement des milliers de requêtes de LLM est lent et coûteux. Le traitement par lots regroupe les requêtes et réduit les coûts de 50 %. L'exécution asynchrone parallélise les requêtes afin de maximiser le débit dans les limites de débit autorisées. Ensemble, ces approches réduisent considérablement à la fois les coûts et le temps total d'exécution.

API de traitement par lots d'OpenAI : réduction des coûts de 50 %

L'API de traitement par lots d'OpenAI traite les requêtes de manière asynchrone en arrière-plan, dans un délai maximal de 24 heures, à 50 % du prix normal de l'API. Elle convient particulièrement aux évaluations, au traitement de jeux de données et aux charges de travail qui ne nécessitent pas de résultats en temps réel.

import openai
import json

client = openai.OpenAI(api_key='YOUR_API_KEY')

# Step 1: Create batch input file (JSONL format)
batch_requests = [
    {
        'custom_id': f'request-{i}',
        'method': 'POST',
        'url': '/v1/chat/completions',
        'body': {
            'model': 'gpt-4o-mini',
            'messages': [
                {'role': 'user', 'content': f'Summarize this document: {doc}'}
            ],
            'max_tokens': 200
        }
    }
    for i, doc in enumerate(['Doc A text...', 'Doc B text...', 'Doc C text...'])
]

# Write to JSONL file
with open('batch_input.jsonl', 'w') as f:
    for req in batch_requests:
        f.write(json.dumps(req) + '\n')

# Step 2: Upload the file
batch_file = client.files.create(
    file=open('batch_input.jsonl', 'rb'),
    purpose='batch'
)
print(f'Batch file uploaded: {batch_file.id}')

Soumission et interrogation d'une tâche de traitement par lots

Après avoir téléversé le fichier d'entrée, créez la tâche de traitement par lots et interrogez-la périodiquement jusqu'à son achèvement. L'API de traitement par lots traite les requêtes dans un délai de 24 heures, généralement bien plus rapidement pour les petits lots.

import time

# Step 3: Create batch job
batch = client.batches.create(
    input_file_id=batch_file.id,
    endpoint='/v1/chat/completions',
    completion_window='24h'
)
print(f'Batch created: {batch.id} | Status: {batch.status}')

# Step 4: Poll for completion
def wait_for_batch(batch_id, poll_interval=30, timeout=3600):
    start = time.time()
    while time.time() - start < timeout:
        batch = client.batches.retrieve(batch_id)
        print(f'Status: {batch.status} | '
              f'Completed: {batch.request_counts.completed}/ '
              f'{batch.request_counts.total}')
        if batch.status == 'completed':
            return batch
        if batch.status in ('failed', 'expired', 'cancelling', 'cancelled'):
            raise RuntimeError(f'Batch {batch_id} ended with status: {batch.status}')
        time.sleep(poll_interval)
    raise TimeoutError('Batch polling timed out')

# batch = wait_for_batch(batch.id)

Récupération des résultats d'un traitement par lots

Une fois le lot terminé, téléchargez le fichier de sortie et analysez les résultats JSONL pour les reconvertir dans un format exploitable.

def retrieve_batch_results(batch):
    if not batch.output_file_id:
        raise ValueError('No output file — batch may have failed')

    # Download output file
    content = client.files.content(batch.output_file_id).text

    # Parse JSONL: one result per line
    results = {}
    for line in content.strip().split('\n'):
        if not line:
            continue
        result = json.loads(line)
        custom_id = result['custom_id']
        if result.get('error'):
            results[custom_id] = {'error': result['error']}
        else:
            response_body = result['response']['body']
            text = response_body['choices'][0]['message']['content']
            results[custom_id] = {'text': text}

    # Report error rate
    errors = sum(1 for r in results.values() if 'error' in r)
    print(f'Retrieved {len(results)} results, {errors} errors')
    return results

# results = retrieve_batch_results(batch)
# for req_id, result in results.items():
#     print(req_id, result.get('text', result.get('error', ''))[:50])

Python asynchrone avec asyncio.gather()

Pour le parallélisme en temps réel, hors traitement par lots, asyncio de Python avec asyncio.gather() lance plusieurs appels d'API simultanément et attend qu'ils soient tous terminés. Cela réduit considérablement le temps total d'exécution pour les charges de travail comportant plusieurs requêtes.

import asyncio
import openai

async_client = openai.AsyncOpenAI(api_key='YOUR_API_KEY')

async def async_completion(messages, model='gpt-4o-mini', max_tokens=200):
    response = await async_client.chat.completions.create(
        model=model,
        messages=messages,
        max_tokens=max_tokens
    )
    return response.choices[0].message.content

async def process_parallel(prompts):
    tasks = [
        async_completion([{'role': 'user', 'content': p}])
        for p in prompts
    ]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    return results

# Usage
async def main():
    prompts = ['Explain photosynthesis.', 'Explain gravity.', 'Explain evolution.']
    results = await process_parallel(prompts)
    for prompt, result in zip(prompts, results):
        if isinstance(result, Exception):
            print(f'ERROR: {result}')
        else:
            print(f'{prompt[:30]}... -> {result[:60]}...')

# asyncio.run(main())
print('asyncio.gather: all 3 requests fire simultaneously')

Appels simultanés tenant compte des limites de débit

Lancer trop de requêtes simultanées déclenche des erreurs de limitation de débit. Un sémaphore limite la concurrence pour rester dans les limites autorisées tout en maximisant le débit.

import asyncio

# Rate limits (example for gpt-4o-mini):
# RPM (requests per minute): 500
# TPM (tokens per minute): 200,000

MAX_CONCURRENT = 20  # stay well below rate limit

async def process_with_rate_limit(prompts, max_concurrent=MAX_CONCURRENT):
    semaphore = asyncio.Semaphore(max_concurrent)
    results = [None] * len(prompts)

    async def bounded_completion(i, prompt):
        async with semaphore:
            try:
                result = await async_completion(
                    [{'role': 'user', 'content': prompt}]
                )
                results[i] = result
            except openai.RateLimitError as e:
                print(f'Rate limited on prompt {i}: {e}')
                await asyncio.sleep(60)  # back off and retry
                result = await async_completion(
                    [{'role': 'user', 'content': prompt}]
                )
                results[i] = result

    await asyncio.gather(*[
        bounded_completion(i, p) for i, p in enumerate(prompts)
    ])
    return results

print('Semaphore limits to', MAX_CONCURRENT, 'concurrent requests')

API de traitement par lots d'Anthropic

Anthropic propose également une API de lots de messages, dont le modèle économique est similaire à celui de l'API de traitement par lots d'OpenAI. Les lots sont traités de manière asynchrone et leurs résultats sont interrogés périodiquement ou transmis en flux.

import anthropic

client = anthropic.Anthropic(api_key='YOUR_API_KEY')

# Create a batch of messages
batch = client.messages.batches.create(
    requests=[
        {
            'custom_id': f'doc-{i}',
            'params': {
                'model': 'claude-haiku-4-5',
                'max_tokens': 200,
                'messages': [{
                    'role': 'user',
                    'content': f'Classify the sentiment of: {text}'
                }]
            }
        }
        for i, text in enumerate([
            'Amazing product, exceeded expectations!',
            'Terrible quality, broke after one use.',
            'It works as described.'
        ])
    ]
)
print(f'Batch created: {batch.id} | Status: {batch.processing_status}')

# Poll for completion
# while (batch := client.messages.batches.retrieve(batch.id)).processing_status != 'ended':
#     time.sleep(30)

# Retrieve results
# for result in client.messages.batches.results(batch.id):
#     print(result.custom_id, result.result.message.content[0].text[:50])

Optimisation du débit&nbsp;: stratégies de traitement par lots

Maximisez le débit en choisissant la stratégie de traitement par lots adaptée à vos exigences de latence et aux caractéristiques de votre charge de travail.

throughput_strategies = {
    'API Batch (OpenAI/Anthropic)': {
        'cost': '50% of normal price',
        'latency': 'Minutes to hours (background processing)',
        'best_for': 'Offline workloads: eval runs, dataset labeling, report generation',
        'max_batch_size': '50,000 requests per batch'
    },
    'asyncio.gather()': {
        'cost': 'Normal price',
        'latency': 'Same as slowest individual request',
        'best_for': 'Real-time parallel enrichment, multi-step pipelines',
        'max_concurrent': '10-50 depending on rate limits'
    },
    'Streaming + Concurrent': {
        'cost': 'Normal price',
        'latency': 'First token arrives faster, total similar',
        'best_for': 'User-facing applications needing perceived speed',
        'pattern': 'asyncio with stream=True per request'
    },
    'Worker Queue (Celery, RQ)': {
        'cost': 'Normal price',
        'latency': 'Variable (depends on queue depth)',
        'best_for': 'High-volume production with auto-scaling workers',
        'backends': 'Redis, RabbitMQ'
    }
}

for strategy, details in throughput_strategies.items():
    print(f'{strategy}: {details["best_for"][:60]}')

Gestion des erreurs et logique de nouvelle tentative pour le traitement par lots et l'asynchronisme

Les charges de travail simultanées et par lots nécessitent une gestion robuste des erreurs. Les échecs de requêtes individuelles ne doivent pas interrompre tout le lot : journalisez-les, réessayez avec une temporisation croissante et indiquez les taux de réussite globaux.

import asyncio
import random

async def resilient_completion(prompt, max_retries=3, base_delay=1.0):
    for attempt in range(max_retries):
        try:
            return await async_completion(
                [{'role': 'user', 'content': prompt}]
            )
        except openai.RateLimitError:
            wait = base_delay * (2 ** attempt) + random.uniform(0, 1)
            print(f'Rate limited. Waiting {wait:.1f}s (attempt {attempt+1})')
            await asyncio.sleep(wait)
        except openai.APITimeoutError:
            print(f'Timeout on attempt {attempt+1}')
            await asyncio.sleep(base_delay)
        except openai.APIError as e:
            if e.status_code >= 500:
                await asyncio.sleep(base_delay * (attempt + 1))
            else:
                raise  # Don't retry 4xx errors
    raise RuntimeError(f'Failed after {max_retries} attempts')

async def batch_with_error_reporting(prompts):
    tasks = [resilient_completion(p) for p in prompts]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    successes = sum(1 for r in results if not isinstance(r, Exception))
    print(f'Batch complete: {successes}/{len(prompts)} succeeded')
    return results

Découpage des entrées volumineuses pour le traitement par lots

Les documents plus volumineux que la fenêtre de contexte du modèle doivent être découpés avant le traitement par lots. Chaque fragment devient une requête distincte du lot ; les résultats sont ensuite fusionnés ou résumés.

def chunk_document(text, max_tokens=3000, overlap_tokens=200):
    '''
    Split a long document into overlapping chunks for batch processing.
    Approximate: 1 token ~ 4 characters
    '''
    max_chars = max_tokens * 4
    overlap_chars = overlap_tokens * 4
    chunks = []
    start = 0
    while start < len(text):
        end = min(start + max_chars, len(text))
        # Try to break at a sentence boundary
        if end < len(text):
            last_period = text.rfind('.', start, end)
            if last_period > start + max_chars // 2:
                end = last_period + 1
        chunks.append({'text': text[start:end], 'start': start, 'end': end})
        start = end - overlap_chars  # overlap for context continuity
    return chunks

def batch_summarize_long_document(document_text, summary_prompt):
    chunks = chunk_document(document_text)
    print(f'Document split into {len(chunks)} chunks')
    # Create one batch request per chunk
    batch_inputs = [
        {'custom_id': f'chunk-{i}',
         'content': summary_prompt + '\n\n' + chunk['text']}
        for i, chunk in enumerate(chunks)
    ]
    # Submit all chunks as one batch job
    return batch_inputs

long_doc = 'Lorem ipsum ' * 5000  # ~20K character document
chunks = chunk_document(long_doc)
print(f'Chunks: {len(chunks)}, first chunk length: {len(chunks[0]["text"])} chars')

Suivi de la progression des grands lots

Pour les grandes tâches de traitement par lots, comportant des milliers de requêtes, affichez la progression en temps réel afin que les opérateurs puissent surveiller le débit et estimer le temps restant.

import asyncio
import time

async def batch_with_progress(prompts, max_concurrent=20):
    semaphore = asyncio.Semaphore(max_concurrent)
    completed = 0
    total = len(prompts)
    start_time = time.time()
    results = [None] * total

    async def process_one(i, prompt):
        nonlocal completed
        async with semaphore:
            results[i] = await resilient_completion(prompt)
            completed += 1

        elapsed = time.time() - start_time
        rate = completed / elapsed if elapsed > 0 else 0
        eta = (total - completed) / rate if rate > 0 else float('inf')

        if completed % 10 == 0 or completed == total:
            print(f'Progress: {completed}/{total} '
                  f'({completed/total:.0%}) | '
                  f'{rate:.1f} req/s | '
                  f'ETA: {eta:.0f}s')

    await asyncio.gather(*[
        process_one(i, p) for i, p in enumerate(prompts)
    ])
    return results

print('Progress tracking: reports every 10 completions with ETA.')

Vérification rapide

Vous devez attribuer une étiquette à 10 000 documents pour une analyse des sentiments pendant la nuit. Vous souhaitez réduire les coûts au minimum et n'avez pas besoin de résultats en temps réel. Quelle approche est la meilleure ?

Résumé du traitement par lots et de l'asynchronisme

Le traitement par lots et l'exécution asynchrone sont essentiels à l'ingénierie des instructions à grande échelle :

  • API de traitement par lots d'OpenAI/Anthropic : réduction des coûts de 50 %, traitement en arrière-plan, jusqu'à 50 000 requêtes par lot
  • asyncio.gather() : requêtes simultanées en temps réel, toutes lancées en même temps et attendant l'ensemble des résultats
  • Semaphore : contrôle de la concurrence tenant compte des limites de débit, généralement de 10 à 50 requêtes simultanées
  • Temporisation exponentielle : nouvelle tentative avec un délai doublé en cas d'erreur de limitation de débit ou de délai d'expiration
  • Rassemblement résilient : return_exceptions=True empêche un seul échec d'interrompre tout le lot
  • Suivi de la progression : signalement des achèvements avec le débit et l'ETA pour les grandes tâches

Questions Fréquemment Posées

La leçon « Traitement par lots et exécution asynchrone » est-elle gratuite ?

Oui — le texte complet de « Traitement par lots et exécution asynchrone » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours AI Prompt Engineering, passe à CoddyKit PRO. Le cours AI Prompt Engineering comprend 4 leçons au total.

Qu'est-ce que j'apprendrai dans « Traitement par lots et exécution asynchrone » ?

API Batch d’OpenAI, Python asynchrone et exécution concurrente des prompts. Tu pratiques AI Prompt Engineering avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.

Dois-je avoir de l'expérience pour commencer AI Prompt Engineering ?

Aucune expérience préalable n'est requise. AI Prompt Engineering sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 2 sur 4.

Combien de temps prend la leçon « Traitement par lots et exécution asynchrone » ?

La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.

Peux-tu écrire et exécuter du code dans cette leçon AI Prompt Engineering ?

Oui. Chaque leçon AI Prompt Engineering inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.

Toutes les leçons de ce cours

  1. Stratégies de mise en cache des prompts
  2. Traitement par lots et exécution asynchrone
  3. Équilibrage de charge entre les modèles
  4. Surveillance et alertes pour les pipelines de prompts
← Retour à AI Prompt Engineering