0Pricing
AI Prompt Engineering · Lección

Procesamiento por lotes y ejecución asíncrona

OpenAI Batch API, Python asíncrono y ejecución simultánea de prompts.

Procesamiento por lotes y ejecución asíncrona es una lección gratuita de AI Prompt Engineering en CoddyKit. Esta es la lección 2 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de AI Prompt Engineering, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de AI Prompt Engineering incluye 4 lecciones en total.

Por qué usar procesamiento por lotes y ejecución asíncrona

Procesar secuencialmente miles de solicitudes a LLM es lento y costoso. El procesamiento por lotes agrupa las solicitudes para reducir los costes un 50 %. La ejecución asíncrona paraleliza las solicitudes para maximizar el rendimiento dentro de los límites de tasa. Juntos, reducen drásticamente tanto los costes como el tiempo total transcurrido.

OpenAI Batch API: reducción de costes del 50 %

OpenAI Batch API procesa las solicitudes de forma asíncrona en segundo plano (hasta 24 horas) por un precio equivalente al 50 % del precio normal de la API. Es ideal para ejecuciones de evaluación, procesamiento de datasets y cargas de trabajo que no requieren resultados en tiempo real.

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

Envío y consulta periódica de un trabajo por lotes

Después de cargar el archivo de entrada, cree el trabajo por lotes y consulte su estado periódicamente hasta que finalice. Batch API procesa las solicitudes en un plazo de 24 horas (normalmente mucho más rápido en lotes pequeños).

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)

Recuperación de resultados por lotes

Cuando finalice el lote, descargue el archivo de salida y analice los resultados JSONL para convertirlos de nuevo a un formato utilizable.

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 asíncrono con asyncio.gather()

Para el paralelismo en tiempo real (sin lotes), el asyncio de Python con asyncio.gather() ejecuta varias llamadas a la API de forma simultánea y espera a que todas finalicen. Esto reduce drásticamente el tiempo total transcurrido en cargas de trabajo con varias solicitudes.

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

Llamadas simultáneas con control de límites de tasa

Enviar demasiadas solicitudes simultáneas provoca errores por superar los límites de tasa. Un semáforo limita la simultaneidad para mantenerse dentro de esos límites y, al mismo tiempo, maximizar el rendimiento.

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

Anthropic Batch API

Anthropic también ofrece una Message Batches API con una economía similar a la Batch API de OpenAI. Los lotes se procesan de forma asíncrona y sus resultados se consultan periódicamente o se transmiten mediante streaming.

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

Optimización del rendimiento: estrategias de procesamiento por lotes

Maximice el rendimiento eligiendo la estrategia de procesamiento por lotes adecuada según sus requisitos de latencia y las características de su carga de trabajo.

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

Manejo de errores y lógica de reintentos en procesos por lotes y asíncronos

Las cargas de trabajo simultáneas y por lotes necesitan un manejo de errores sólido. Los fallos de solicitudes individuales no deben detener todo el lote: regístrelos, vuelva a intentarlo con backoff y notifique las tasas de éxito agregadas.

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

División de entradas grandes en fragmentos para el procesamiento por lotes

Los documentos que superan la ventana de contexto del modelo deben dividirse en fragmentos antes de procesarlos por lotes. Cada fragmento se convierte en una solicitud por lotes independiente; posteriormente, los resultados se combinan o se resumen.

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

Seguimiento del progreso de lotes grandes

En trabajos por lotes grandes (miles de solicitudes), muestre el progreso en tiempo real para que los operadores puedan supervisar el rendimiento y estimar el tiempo de finalización.

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

Comprobación rápida

Necesita etiquetar 10 000 documentos para analizar su sentimiento durante la noche. Quiere minimizar los costes y no necesita resultados en tiempo real. ¿Qué enfoque es el más adecuado?

Resumen del procesamiento por lotes y la ejecución asíncrona

El procesamiento por lotes y la ejecución asíncrona son esenciales para la ingeniería de prompts a gran escala:

  • OpenAI/Anthropic Batch API: reducción de costes del 50 %, procesamiento en segundo plano y hasta 50K solicitudes por lote
  • asyncio.gather(): solicitudes simultáneas en tiempo real; todas se ejecutan a la vez y se espera a que lleguen todos los resultados
  • Semáforo: control de la simultaneidad teniendo en cuenta los límites de tasa (normalmente, entre 10 y 50 solicitudes simultáneas)
  • Backoff exponencial: reintentos con un retraso que se duplica ante errores de límite de tasa o tiempo de espera
  • gather resistente: return_exceptions=True evita que un fallo detenga todo el lote
  • Seguimiento del progreso: informa de las solicitudes completadas junto con la tasa y la ETA en trabajos grandes

Preguntas frecuentes

¿La lección «Procesamiento por lotes y ejecución asíncrona» es gratis?

Sí — el texto completo de «Procesamiento por lotes y ejecución asíncrona» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de AI Prompt Engineering, actualiza a CoddyKit PRO. El curso de AI Prompt Engineering incluye 4 lecciones en total.

¿Qué aprenderé en «Procesamiento por lotes y ejecución asíncrona»?

OpenAI Batch API, Python asíncrono y ejecución simultánea de prompts. Practicas AI Prompt Engineering con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar AI Prompt Engineering?

No se requiere experiencia previa. AI Prompt Engineering en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 2 de 4.

¿Cuánto tiempo toma la lección «Procesamiento por lotes y ejecución asíncrona»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de AI Prompt Engineering?

Sí. Cada lección de AI Prompt Engineering incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. Estrategias de caché para prompts
  2. Procesamiento por lotes y ejecución asíncrona
  3. Balanceo de carga entre modelos
  4. Monitorización y alertas para pipelines de prompts
← Volver a AI Prompt Engineering