0Pricing
AI Prompt Engineering · Lekcja

Przetwarzanie wsadowe i wykonywanie asynchroniczne

OpenAI Batch API, asynchroniczny Python i równoległe wykonywanie promptów.

Przetwarzanie wsadowe i wykonywanie asynchroniczne to bezpłatna lekcja AI Prompt Engineering na CoddyKit. To lekcja 2 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AI Prompt Engineering, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AI Prompt Engineering zawiera 4 lekcji w sumie.

Dlaczego przetwarzanie wsadowe i asynchroniczne?

Sekwencyjne przetwarzanie tysięcy żądań LLM jest powolne i kosztowne. Przetwarzanie wsadowe grupuje żądania, obniżając koszt o 50%. Wykonywanie asynchroniczne równolegle przetwarza żądania, maksymalizując przepustowość w ramach limitów żądań. Razem te rozwiązania znacznie zmniejszają zarówno koszt, jak i rzeczywisty czas wykonania.

OpenAI Batch API: obniżenie kosztu o 50%

OpenAI Batch API przetwarza żądania asynchronicznie w tle (do 24 godzin), pobierając 50% standardowej ceny API. Jest idealne do przeprowadzania ewaluacji, przetwarzania zbiorów danych i zadań niewymagających wyników w czasie rzeczywistym.

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

Przesyłanie zadania wsadowego i sprawdzanie jego statusu

Po przesłaniu pliku wejściowego należy utworzyć zadanie wsadowe i odpytywać jego status aż do ukończenia. Batch API przetwarza żądania w ciągu 24 godzin (zwykle znacznie szybciej w przypadku małych partii).

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)

Pobieranie wyników zadania wsadowego

Po ukończeniu zadania wsadowego należy pobrać plik wyjściowy i przeanalizować wyniki JSONL, aby ponownie uzyskać format możliwy do wykorzystania.

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

Asynchroniczny Python z asyncio.gather()

W przypadku równoległości w czasie rzeczywistym (poza przetwarzaniem wsadowym) asyncio w Pythonie wraz z asyncio.gather() uruchamia jednocześnie wiele wywołań API i czeka na ich zakończenie. Znacznie skraca to całkowity rzeczywisty czas przetwarzania zadań obejmujących wiele żądań.

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

Równoczesne wywołania uwzględniające limity żądań

Uruchomienie zbyt wielu równoczesnych żądań powoduje błędy przekroczenia limitu żądań. Semafor ogranicza równoczesność, aby zachować limity żądań i jednocześnie maksymalizować przepustowość.

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 również oferuje Message Batches API o podobnej ekonomice co Batch API firmy OpenAI. Partie są przetwarzane asynchronicznie, a wyniki można pobierać przez odpytywanie statusu lub strumieniować.

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

Optymalizacja przepustowości: strategie przetwarzania wsadowego

Aby zmaksymalizować przepustowość, należy wybrać odpowiednią strategię przetwarzania wsadowego na podstawie wymagań dotyczących opóźnień i charakterystyki obciążenia.

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

Obsługa błędów i ponawianie w trybie wsadowym i asynchronicznym

Równoczesne i wsadowe obciążenia wymagają solidnej obsługi błędów. Pojedyncze niepowodzenia żądań nie powinny przerywać całej partii — należy je rejestrować, ponawiać z narastającym opóźnieniem i raportować zbiorcze współczynniki powodzenia.

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

Dzielenie dużych danych wejściowych na fragmenty na potrzeby przetwarzania wsadowego

Dokumenty większe niż okno kontekstu modelu należy podzielić na fragmenty przed przetwarzaniem wsadowym. Każdy fragment staje się osobnym żądaniem wsadowym, a wyniki są później scalane lub podsumowywane.

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

Śledzenie postępu dużych partii

W przypadku dużych zadań wsadowych (tysięcy żądań) należy wyświetlać postęp w czasie rzeczywistym, aby operatorzy mogli monitorować przepustowość i szacować czas ukończenia.

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

Szybkie sprawdzenie

Należy przez noc oznaczyć 10 000 dokumentów pod kątem analizy sentymentu. Chcą Państwo zminimalizować koszt i nie potrzebują wyników w czasie rzeczywistym. Które podejście będzie najlepsze?

Podsumowanie przetwarzania wsadowego i asynchronicznego

Przetwarzanie wsadowe i wykonywanie asynchroniczne są niezbędne w inżynierii promptów na dużą skalę:

  • OpenAI/Anthropic Batch API: obniżenie kosztu o 50%, przetwarzanie w tle, do 50 tys. żądań w jednej partii
  • asyncio.gather(): równoczesne żądania w czasie rzeczywistym, wszystkie uruchamiane jednocześnie i oczekujące na komplet wyników
  • Semafor: kontrola równoczesności uwzględniająca limity żądań (zwykle 10–50 równoczesnych żądań)
  • Wykładnicze zwiększanie opóźnienia: ponawianie z podwajanym opóźnieniem w przypadku błędów limitu żądań lub przekroczenia czasu
  • Odporne gather: parametr return_exceptions=True zapobiega przerwaniu całej partii przez pojedyncze niepowodzenie
  • Śledzenie postępu: raportowanie ukończonych zadań wraz z szybkością i przewidywanym czasem zakończenia dla dużych zadań

Często zadawane pytania

Czy lekcja „Przetwarzanie wsadowe i wykonywanie asynchroniczne” jest bezpłatna?

Tak — pełny tekst „Przetwarzanie wsadowe i wykonywanie asynchroniczne” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AI Prompt Engineering, przejdź na CoddyKit PRO. Kurs AI Prompt Engineering zawiera 4 lekcji w sumie.

Co nauczysz się w „Przetwarzanie wsadowe i wykonywanie asynchroniczne”?

OpenAI Batch API, asynchroniczny Python i równoległe wykonywanie promptów. Ćwiczysz AI Prompt Engineering z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.

Czy potrzebuję doświadczenia, aby zacząć AI Prompt Engineering?

Nie wymagamy żadnego doświadczenia. AI Prompt Engineering w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 2 z 4.

Ile czasu zajmuje lekcja „Przetwarzanie wsadowe i wykonywanie asynchroniczne”?

Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.

Czy mogę pisać i uruchamiać kod w tej lekcji AI Prompt Engineering?

Tak. Każda lekcja AI Prompt Engineering zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.

Wszystkie lekcje w tym kursie

  1. Strategie buforowania promptów
  2. Przetwarzanie wsadowe i wykonywanie asynchroniczne
  3. Równoważenie obciążenia między modelami
  4. Monitorowanie i alerty dla potoków promptów
← Powrót do AI Prompt Engineering