Toplu İşleme ve Eşzamansız Yürütme
OpenAI Batch API, eşzamansız Python ve istemlerin eşzamanlı yürütülmesi.
Toplu İşleme ve Eşzamansız Yürütme, CoddyKit'te ücretsiz bir AI Prompt Engineering dersidir. Bu, 4 dersinin 2. dersidir. Aşağıdan dersin tamamını ücretsiz okuyabilir, sonra tarayıcıda yerleşik kod editörü ve 7/24 yapay zeka koçu ile uygulamalı olarak pratik yapabilirsin. Bu, AI Prompt Engineering öğrenme yolunun bir parçasıdır ve ilerlemeniz web ve CoddyKit uygulaması arasında senkronize olur. AI Prompt Engineering kursu toplamda 4 dersten oluşur.
Neden Toplu İş ve Eşzamansızlık?
Binlerce LLM isteğini ardışık olarak işlemek yavaş ve pahalıdır. Toplu işleme, maliyeti %50 azaltmak için istekleri gruplandırır. Eşzamansız yürütme, hız sınırları içinde iş hacmini en üst düzeye çıkarmak için istekleri paralel hale getirir. Birlikte kullanıldıklarında hem maliyeti hem de geçen süreyi büyük ölçüde azaltırlar.
OpenAI Toplu İş API'si: %50 Maliyet Azaltma
OpenAI Toplu İş API'si, istekleri arka planda eşzamansız olarak (24 saate kadar) normal API ücretinin %50'si karşılığında işler. Değerlendirme çalıştırmaları, veri kümesi işleme ve gerçek zamanlı olmayan iş yükleri için idealdir.
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}')Toplu İş Gönderme ve Durumunu Sorgulama
Girdi dosyasını yükledikten sonra toplu işi oluşturun ve tamamlanana kadar durumunu sorgulayın. Toplu İş API'si istekleri 24 saat içinde işler (küçük toplu işler genellikle çok daha hızlı tamamlanır).
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)Toplu İş Sonuçlarını Alma
Toplu iş tamamlandığında çıktı dosyasını indirin ve JSONL sonuçlarını kullanılabilir bir biçime ayrıştırın.
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])asyncio.gather() ile Eşzamansız Python
Gerçek zamanlı (toplu iş olmayan) paralellik için Python'daki asyncio, birden çok API çağrısını eşzamanlı olarak başlatır ve tümünün tamamlanmasını bekler. Bu, birden çok istek içeren iş yüklerinde toplam geçen süreyi büyük ölçüde azaltır.
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')Hız Sınırlarını Dikkate Alan Eşzamanlı Çağrılar
Çok fazla eşzamanlı istek başlatmak, hız sınırı hatalarını tetikler. Bir semafor, iş hacmini en üst düzeye çıkarırken hız sınırları içinde kalmak için eşzamanlılığı sınırlar.
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 Toplu İş API'si
Anthropic de OpenAI'nin Toplu İş API'sine benzer ekonomik koşullara sahip bir Mesaj Topluları API'si sunar. Toplu işler eşzamansız olarak işlenir ve sonuçların durumu sorgulanabilir veya sonuçlar akış olarak alınabilir.
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])İş Hacmi Optimizasyonu: Toplu İşleme Stratejileri
Gecikme gereksinimlerinize ve iş yükünüzün özelliklerine göre doğru toplu işleme stratejisini seçerek iş hacmini en üst düzeye çıkarın.
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]}')Toplu İşlerde ve Eşzamansız İşlemlerde Hata İşleme ve Yeniden Deneme Mantığı
Eşzamanlı ve toplu iş yükleri sağlam hata işleme gerektirir. Tek tek isteklerin başarısız olması tüm toplu işi çökertmemelidir — bunları günlüğe kaydedin, geri çekilmeyle yeniden deneyin ve toplam başarı oranlarını bildirin.
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 resultsToplu İşleme İçin Büyük Girdileri Parçalara Ayırma
Modelin bağlam penceresinden daha büyük belgeler, toplu işlemden önce parçalara ayrılmalıdır. Her parça ayrı bir toplu iş isteğine dönüşür; sonuçlar daha sonra birleştirilir veya özetlenir.
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')Büyük Toplu İşlerde İlerleme Takibi
Büyük toplu işlerde (binlerce istek) işleçlerin iş hacmini izlemesi ve tamamlanma süresini tahmin edebilmesi için gerçek zamanlı ilerlemeyi gösterin.
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.')Kısa Kontrol
Gece boyunca duygu analizi için 10.000 belgeyi etiketlemeniz gerekiyor. Maliyeti en aza indirmek istiyor ve gerçek zamanlı sonuçlara ihtiyaç duymuyorsunuz. En iyi yaklaşım hangisidir?
Toplu İş ve Eşzamansızlık Özeti
Toplu işleme ve eşzamansız yürütme, büyük ölçekte istem mühendisliği için gereklidir:
- OpenAI/Anthropic Toplu İş API'si: %50 maliyet azaltma, arka planda işleme, toplu iş başına 50 bine kadar istek
- asyncio.gather(): eşzamanlı gerçek zamanlı istekler; tümü aynı anda başlatılır ve tüm sonuçlar beklenir
- Semaphore: hız sınırlarını dikkate alan eşzamanlılık denetimi (genellikle 10-50 eşzamanlı istek)
- Üstel geri çekilme: hız sınırı veya zaman aşımı hatalarında gecikmeyi iki katına çıkararak yeniden deneme
- Dayanıklı toplama: return_exceptions=True, tek bir başarısızlığın tüm toplu işi çökertmesini önler
- İlerleme takibi: büyük işler için tamamlanmaları hız ve ETA ile bildirme
Sıkça Sorulan Sorular
“Toplu İşleme ve Eşzamansız Yürütme” dersi ücretsiz mi?
Evet — “Toplu İşleme ve Eşzamansız Yürütme” dersin tüm metni burada web'de ücretsiz olarak okunabilir. Etkileşimli olarak pratik yapmak (yerleşik kod editörü ve 7/24 yapay zeka koçu) ve AI Prompt Engineering kursunun geri kalanını açmak için CoddyKit PRO'ya yükselt. AI Prompt Engineering kursu toplamda 4 dersten oluşur.
“Toplu İşleme ve Eşzamansız Yürütme” dersinde ne öğreneceğim?
OpenAI Batch API, eşzamansız Python ve istemlerin eşzamanlı yürütülmesi. AI Prompt Engineering ile uygulamalı kodu tarayıcıda doğrudan çalıştırarak pratik yaparsın ve 7/24 yapay zeka koçu dersi çalışırken sorularını yanıtlar.
AI Prompt Engineering öğrenmeye başlamak için deneyim gerekli mi?
Önceden deneyim gerekmez. CoddyKit'te AI Prompt Engineering, başlangıçtan ileri seviyeye kadar yapılandırıldığı için buradan başlayabilir veya başından başlayıp kendi hızında ilerleme yapabilirsin. Bu, 4 dersinin 2. dersidir.
“Toplu İşleme ve Eşzamansız Yürütme” dersi ne kadar sürer?
Çoğu CoddyKit dersi yaklaşık 5–10 dakika sürer. Her biri kısa ve etkileşimli olduğu için sabit ilerleme yaparsın ve web ile uygulama arasında tam olarak bıraktığın yerden devam edebilirsin.
Bu AI Prompt Engineering dersinde kod yazıp çalıştırabilir miyim?
Evet. Her AI Prompt Engineering dersi yerleşik bir kod editörü içerir, bu sayede tarayıcıda gerçek kod yazıp çalıştırabilir ve anlık yapay zeka geri bildirimi alırsın — yerel kurulum gerekli değildir.
Bu kursun tüm dersleri
- İstemler İçin Önbellekleme Stratejileri
- Toplu İşleme ve Eşzamansız Yürütme
- Modeller Arasında Yük Dengeleme
- İstem İş Akışlarını İzleme ve Uyarı Verme