Kejuruteraan Arahan AI · Pelajaran

Pemprosesan Kelompok dan Pelaksanaan Tak Segerak

OpenAI Batch API, Python tak segerak dan pelaksanaan gesaan serentak.

Pelajaran 2 daripada 413 langkah

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 results

Membahagikan 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
Percuma untuk bermula

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

  1. Strategi Penyimpanan Cache untuk Gesaan
  2. Pemprosesan Kelompok dan Pelaksanaan Tak Segerak
  3. Pengimbangan Beban Merentas Model
  4. Pemantauan dan Makluman untuk Pipeline Gesaan
← Kembali ke Kejuruteraan Arahan AI