المعالجة الدفعية والتنفيذ غير المتزامن
واجهة OpenAI Batch API وPython غير المتزامنة والتنفيذ المتزامن للمطالبات
المعالجة الدفعية والتنفيذ غير المتزامن درس مجاني في AI Prompt Engineering على CoddyKit. هذا هو الدرس 2 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7. هذا الدرس جزء من مسار التعلم في AI Prompt Engineering، وتقدمك يتزامن عبر الويب وتطبيق CoddyKit. تتضمن دورة AI Prompt Engineering 4 دروس في المجموع.
لماذا المعالجة الدفعية والتنفيذ غير المتزامن؟
تُعد معالجة آلاف طلبات LLM بالتتابع بطيئة ومكلفة. تجمع المعالجة الدفعية الطلبات لتقليل التكلفة بنسبة 50%. ويوازي التنفيذ غير المتزامن الطلبات لتحقيق أقصى إنتاجية ضمن حدود المعدل. ويؤدي الجمع بينهما إلى تقليل التكلفة والوقت الفعلي المنقضي بدرجة كبيرة.
واجهة OpenAI Batch API: خفض التكلفة بنسبة 50%
تعالج OpenAI Batch API الطلبات بشكل غير متزامن في الخلفية (حتى 24 ساعة) بسعر يعادل 50% من السعر المعتاد لواجهة API. وهي مثالية لعمليات التقييم، ومعالجة مجموعات البيانات، وأحمال العمل غير الآنية.
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}')إرسال مهمة دفعية والاستطلاع الدوري لحالتها
بعد تحميل ملف الإدخال، أنشئ المهمة الدفعية واستطلع حالتها حتى تكتمل. تعالج Batch API الطلبات خلال 24 ساعة (وعادةً ما تكون أسرع بكثير للدفعات الصغيرة).
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)استرداد نتائج المعالجة الدفعية
بعد اكتمال الدفعة، نزّل ملف الإخراج وحلّل نتائج JSONL مجددًا إلى تنسيق قابل للاستخدام.
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 غير المتزامن باستخدام asyncio.gather()
للتوازي الآني (غير الدفعي)، تطلق asyncio في Python باستخدام asyncio.gather() عدة استدعاءات لواجهة API بشكل متزامن، مع الانتظار حتى اكتمالها جميعًا. ويقلل ذلك بدرجة كبيرة الوقت الفعلي المنقضي إجمالًا لأحمال العمل التي تتضمن طلبات متعددة.
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')استدعاءات متزامنة تراعي حدود المعدل
يؤدي إطلاق عدد كبير جدًا من الطلبات المتزامنة إلى أخطاء تجاوز حدود المعدل. ويحدّ كائن semaphore من التزامن للبقاء ضمن حدود المعدل مع زيادة الإنتاجية إلى أقصى حد.
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 أيضًا Message Batches API باقتصاديات مشابهة لـ Batch API من OpenAI. وتُعالج الدفعات بشكل غير متزامن، ثم تُستطلع النتائج أو تُبث.
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])تحسين الإنتاجية: استراتيجيات المعالجة الدفعية
حقّق أقصى إنتاجية باختيار استراتيجية المعالجة الدفعية المناسبة بناءً على متطلبات زمن الاستجابة وخصائص حمل العمل.
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]}')معالجة الأخطاء ومنطق إعادة المحاولة في المعالجة الدفعية وغير المتزامنة
تحتاج أحمال العمل المتزامنة والدفعية إلى معالجة قوية للأخطاء. يجب ألا تؤدي إخفاقات الطلبات الفردية إلى تعطيل الدفعة بأكملها — سجّل هذه الإخفاقات، وأعد المحاولة مع التراجع التدريجي، وأبلغ عن معدلات النجاح الإجمالية.
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تقسيم المدخلات الكبيرة للمعالجة الدفعية
يجب تقسيم المستندات التي يتجاوز حجمها نافذة سياق النموذج إلى أجزاء قبل إجراء المعالجة الدفعية. ويصبح كل جزء طلبًا دفعيًا منفصلًا، ثم تُدمج النتائج أو تُلخّص لاحقًا.
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')تتبّع التقدم للدفعات الكبيرة
بالنسبة إلى المهام الدفعية الكبيرة (التي تتضمن آلاف الطلبات)، اعرض التقدم في الوقت الفعلي حتى يتمكن المشغّلون من مراقبة الإنتاجية وتقدير وقت الاكتمال.
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.')تحقق سريع
تحتاج إلى تصنيف 10,000 مستند لتحليل المشاعر خلال الليل. تريد تقليل التكلفة إلى أدنى حد، ولا تحتاج إلى نتائج آنية. ما النهج الأفضل؟
ملخص المعالجة الدفعية والتنفيذ غير المتزامن
تُعد المعالجة الدفعية والتنفيذ غير المتزامن أساسيين لهندسة المطالبات على نطاق واسع:
- OpenAI/Anthropic Batch API: خفض التكلفة بنسبة 50%، ومعالجة في الخلفية، وما يصل إلى 50 ألف طلب لكل دفعة
- asyncio.gather(): طلبات آنية متزامنة، تُطلق جميعها في الوقت نفسه وتنتظر جميع النتائج
- Semaphore: التحكم في التزامن مع مراعاة حدود المعدل (عادةً من 10 إلى 50 طلبًا متزامنًا)
- التراجع الأُسّي: إعادة المحاولة مع مضاعفة التأخير عند حدوث أخطاء حدود المعدل أو انتهاء المهلة
- التجميع المرن: تمنع return_exceptions=True فشلًا واحدًا من تعطيل الدفعة
- تتبّع التقدم: الإبلاغ عن حالات الاكتمال مع معدل الإنجاز والوقت المقدّر للوصول إلى الاكتمال للمهام الكبيرة
الأسئلة الشائعة
هل درس «المعالجة الدفعية والتنفيذ غير المتزامن» مجاني؟
نعم — نص درس «المعالجة الدفعية والتنفيذ غير المتزامن» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7) وفتح باقي دورة AI Prompt Engineering، انتقل إلى CoddyKit PRO. تتضمن دورة AI Prompt Engineering 4 دروس في المجموع.
ماذا ستتعلم في «المعالجة الدفعية والتنفيذ غير المتزامن»؟
واجهة OpenAI Batch API وPython غير المتزامنة والتنفيذ المتزامن للمطالبات تتمرن على AI Prompt Engineering مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.
هل أحتاج إلى خبرة سابقة لأبدأ AI Prompt Engineering؟
لا تُشترط خبرة سابقة. AI Prompt Engineering على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 2 من أصل 4.
كم من الوقت يستغرق درس «المعالجة الدفعية والتنفيذ غير المتزامن»؟
معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.
هل يمكنني كتابة وتشغيل أكواد في درس AI Prompt Engineering هذا؟
نعم. كل درس في AI Prompt Engineering يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.
جميع الدروس في هذه الدورة
- استراتيجيات التخزين المؤقت للمطالبات
- المعالجة الدفعية والتنفيذ غير المتزامن
- موازنة الحمل بين النماذج
- مراقبة مسارات المطالبات والتنبيه بشأنها