Pemprosesan Kelompok dan Pelaksanaan Tak Segerak
OpenAI Batch API, Python tak segerak dan pelaksanaan gesaan serentak.
Pemprosesan Kelompok dan Pelaksanaan Tak Segerak ialah pelajaran Kejuruteraan Arahan AI percuma di CoddyKit. Ini ialah pelajaran 2 daripada 4. Sebanyak 3 pelajaran dalam laluan pembelajaran ini boleh dibaca sepenuhnya secara percuma — selepas itu, CoddyKit PRO membuka akses kepada semua pelajaran, serta latihan praktikal dengan penyunting kod terbina dalam dan tutor kecerdasan buatan yang tersedia 24/7. Pelajaran ini merupakan sebahagian daripada laluan pembelajaran Kejuruteraan Arahan AI, dan kemajuan anda disegerakkan merentas web serta aplikasi CoddyKit. Kursus Kejuruteraan Arahan AI merangkumi sejumlah 4 pelajaran.
Mengapa Kelompok dan Tak Segerak?
Memproses ribuan permintaan LLM secara berjujukan adalah lambat dan mahal. Pemprosesan kelompok mengumpulkan permintaan untuk mengurangkan kos sebanyak 50%. Pelaksanaan tak segerak menjalankan permintaan secara selari bagi memaksimumkan daya pemprosesan dalam had kadar. Bersama-sama, kedua-duanya mengurangkan kos dan masa jam sebenar dengan ketara.
API Kelompok OpenAI: Pengurangan Kos 50%
API Kelompok OpenAI memproses permintaan secara tak segerak di latar belakang (sehingga 24 jam) pada 50% daripada harga API biasa. Pendekatan ini sesuai untuk pelaksanaan penilaian, pemprosesan set data dan beban kerja bukan masa nyata.
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}')Menghantar dan Meninjau Kerja Kelompok
Selepas memuat naik fail masukan, gunakan create untuk menghasilkan kerja kelompok dan lakukan tinjauan sehingga selesai. API Kelompok memproses permintaan dalam tempoh 24 jam (biasanya lebih pantas untuk kelompok kecil).
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)Mendapatkan Hasil Kelompok
Setelah kelompok selesai, muat turun fail output dan huraikan results JSONL kembali kepada format yang boleh digunakan.
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 Tak Segerak dengan asyncio.gather()
Untuk pemprosesan selari masa nyata (bukan kelompok), asyncio Python dengan asyncio.gather() melancarkan berbilang panggilan API secara serentak dan menunggu semuanya selesai. Pendekatan ini mengurangkan masa jam sebenar keseluruhan dengan ketara untuk beban kerja yang melibatkan banyak permintaan.
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')Panggilan Serentak yang Peka terhadap Had Kadar
Melancarkan terlalu banyak permintaan serentak akan mencetuskan ralat had kadar. Semaphore mengehadkan keserentakan supaya kekal dalam had kadar sambil memaksimumkan daya pemprosesan.
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 Kelompok Anthropic
Anthropic juga menawarkan API Message Batches dengan ekonomi yang serupa dengan API Kelompok OpenAI. Kelompok diproses secara tak segerak dan results ditinjau atau distrim.
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])Pengoptimuman Daya Pemprosesan: Strategi Pengelompokan
Maksimumkan daya pemprosesan dengan memilih strategi pengelompokan yang sesuai berdasarkan keperluan kependaman dan ciri-ciri beban kerja anda.
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]}')Pengendalian Ralat dan Logik Percubaan Semula dalam Kelompok/Tak Segerak
Beban kerja serentak dan kelompok memerlukan pengendalian ralat yang kukuh. Kegagalan permintaan secara individu tidak sepatutnya menyebabkan seluruh kelompok terhenti — catat kegagalan tersebut, cuba semula dengan penangguhan beransur-ansur dan laporkan kadar kejayaan keseluruhan.
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 resultsMembahagikan Masukan Besar untuk Pemprosesan Kelompok
Dokumen yang lebih besar daripada tetingkap konteks model mesti dibahagikan sebelum dikelompokkan. Setiap bahagian menjadi permintaan kelompok yang berasingan; hasilnya kemudian digabungkan atau diringkaskan.
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')Menjejaki Kemajuan untuk Kelompok Besar
Untuk kerja kelompok besar (ribuan permintaan), paparkan kemajuan masa nyata supaya pengendali boleh memantau daya pemprosesan dan estimate masa penyelesaian.
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.')Semakan Ringkas
Anda perlu melabel 10,000 dokumen untuk analisis sentimen sepanjang malam. Anda mahu meminimumkan kos dan tidak memerlukan hasil masa nyata. Pendekatan manakah yang paling sesuai?
Ringkasan Kelompok dan Tak Segerak
Pemprosesan kelompok dan pelaksanaan tak segerak penting untuk kejuruteraan gesaan pada skala besar:
- API Kelompok OpenAI/Anthropic: pengurangan kos 50%, pemprosesan latar belakang, sehingga 50K permintaan bagi setiap kelompok
- asyncio.gather(): permintaan masa nyata serentak, semuanya dilancarkan serentak dan menunggu semua hasil
- Semaphore: kawalan keserentakan yang peka terhadap had kadar (kebiasaannya 10-50 permintaan serentak)
- Undur balik eksponen: cuba semula dengan kelewatan berganda apabila berlaku ralat had kadar atau tamat masa
- Pengumpulan yang tahan ralat: return_exceptions=True menghalang satu kegagalan daripada menyebabkan seluruh kelompok terhenti
- Penjejakan kemajuan: laporkan penyelesaian bersama kadar dan ETA untuk kerja besar
Pelajari Kejuruteraan Arahan AI dengan tutor kecerdasan buatan — percuma
Tulis dan jalankan kod sebenar dalam pelayar anda, dapatkan bantuan segera daripada tutor kecerdasan buatan yang tersedia 24/7, dan sambung semula dari tempat anda berhenti di web atau dalam aplikasi.
- Kursus
- 53
- Pelajaran
- 199
Soalan Lazim
Adakah pelajaran “Pemprosesan Kelompok dan Pelaksanaan Tak Segerak” percuma?
Ya — sebanyak 3 pelajaran dalam laluan pembelajaran Kejuruteraan Arahan AI, termasuk “Pemprosesan Kelompok dan Pelaksanaan Tak Segerak”, boleh dibaca sepenuhnya secara percuma di web ini. Selepas itu, CoddyKit PRO membuka akses kepada semua pelajaran, serta latihan interaktif dengan penyunting kod terbina dalam dan tutor kecerdasan buatan yang tersedia 24/7. Kursus Kejuruteraan Arahan AI merangkumi sejumlah 4 pelajaran.
Apakah yang akan saya pelajari dalam “Pemprosesan Kelompok dan Pelaksanaan Tak Segerak”?
OpenAI Batch API, Python tak segerak dan pelaksanaan gesaan serentak. Anda berlatih Kejuruteraan Arahan AI menggunakan kod praktikal yang dijalankan terus dalam pelayar, manakala tutor kecerdasan buatan 24/7 menjawab soalan anda semasa anda mengikuti pelajaran.
Adakah saya memerlukan pengalaman untuk memulakan Kejuruteraan Arahan AI?
Tiada pengalaman terdahulu diperlukan. Pembelajaran Kejuruteraan Arahan AI di CoddyKit disusun untuk pelajar daripada peringkat pemula hingga lanjutan, jadi anda boleh bermula di sini atau dari awal dan belajar mengikut kadar anda sendiri. Ini ialah pelajaran 2 daripada 4.
Berapa lamakah pelajaran “Pemprosesan Kelompok dan Pelaksanaan Tak Segerak” diambil?
Kebanyakan pelajaran CoddyKit mengambil masa kira-kira 5–10 minit. Setiap pelajaran ringkas dan interaktif, jadi anda boleh membuat kemajuan secara berterusan dan menyambung tepat dari tempat anda berhenti di web atau aplikasi.
Bolehkah saya menulis dan menjalankan kod dalam pelajaran Kejuruteraan Arahan AI ini?
Ya. Setiap pelajaran Kejuruteraan Arahan AI menyertakan penyunting kod terbina dalam, jadi anda boleh menulis dan menjalankan kod sebenar terus dalam pelayar serta menerima maklum balas kecerdasan buatan serta-merta — tanpa memerlukan persediaan setempat.
Semua pelajaran dalam kursus ini
- Strategi Penyimpanan Cache untuk Gesaan
- Pemprosesan Kelompok dan Pelaksanaan Tak Segerak
- Pengimbangan Beban Merentas Model
- Pemantauan dan Makluman untuk Pipeline Gesaan