Pemrosesan Kelompok dengan Asinkron dan Antrean
Bangun alur pemrosesan ekstraksi asinkron menggunakan asyncio dan antrean tugas untuk memproses ribuan dokumen secara paralel, dengan tetap mematuhi batas laju dan melacak kemajuan.
Pemrosesan Kelompok dengan Asinkron dan Antrean adalah pelajaran AI Engineering Academy gratis di CoddyKit. Ini adalah pelajaran 3 dari 4. Kamu bisa membaca pelajaran lengkapnya di bawah secara gratis — lalu praktikkan langsung di browser dengan editor kode bawaan dan tutor AI 24/7. Ini adalah bagian dari jalur belajar AI Engineering Academy, dan progresmu tersinkronisasi di web dan aplikasi CoddyKit. Kursus AI Engineering Academy mencakup 4 pelajaran total.
Mengapa Pemrosesan Batch Penting
Memproses ribuan dokumen satu per satu terlalu lambat untuk produksi. Perulangan sinkron yang memanggil API OpenAI secara berurutan mungkin hanya memproses 1 dokumen per detik, sehingga 10.000 dokumen memerlukan hampir 3 jam. Pemrosesan batch asinkron dapat memparalelkan ratusan permintaan secara bersamaan dan mengurangi total waktu nyata hingga sepuluh kali lipat.
Dasar asyncio
Perulangan peristiwa asyncio Python memungkinkan Anda menjalankan banyak tugas yang terikat I/O secara bersamaan tanpa thread. Saat satu pemanggilan API menunggu respons jaringan, perulangan peristiwa beralih untuk memproses pemanggilan lainnya. Anda menulis kode dengan kata kunci async def dan await, lalu runtime menangani penjadwalannya. Hal ini ideal untuk pemanggilan LLM yang sebagian besar waktunya dihabiskan untuk menunggu server.
import asyncio
import instructor
from openai import AsyncOpenAI
async_client = instructor.from_openai(AsyncOpenAI())
async def extract_one(text: str) -> PersonExtract:
return await async_client.chat.completions.create(
model='gpt-4o-mini',
response_model=PersonExtract,
messages=[{'role': 'user', 'content': text}]
)Menjalankan Beberapa Ekstraksi dengan gather
asyncio.gather menjalankan daftar coroutine secara bersamaan dan mengembalikan semua hasil setelah yang terakhir selesai. Untuk kumpulan kecil dokumen, cara ini sudah memadai. Bungkus coroutine ekstraksi Anda dalam pemahaman daftar, lalu teruskan ke gather. Total waktunya kira-kira sama dengan satu panggilan yang paling lambat, bukan jumlah semua panggilan.
async def batch_extract(texts: list) -> list:
tasks = [extract_one(text) for text in texts]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Filter out exceptions
successes = [r for r in results if not isinstance(r, Exception)]
failures = [r for r in results if isinstance(r, Exception)]
print(f'Success: {len(successes)}, Failures: {len(failures)}')
return successes
results = asyncio.run(batch_extract(documents))Mengendalikan Konkurensi dengan Semaphore
Mengirim ribuan permintaan secara bersamaan akan mencapai batas laju dan menyebabkan error 429. Gunakan asyncio.Semaphore untuk membatasi jumlah panggilan API secara bersamaan. Semaphore dengan nilai 50 berarti paling banyak 50 panggilan yang sedang berlangsung pada satu waktu. Sesuaikan angka ini berdasarkan batas laju tingkat akun OpenAI Anda untuk model target.
import asyncio
sem = asyncio.Semaphore(50) # max 50 concurrent calls
async def extract_with_limit(text: str, semaphore: asyncio.Semaphore):
async with semaphore:
return await extract_one(text)
async def batch_extract_limited(texts: list):
tasks = [extract_with_limit(t, sem) for t in texts]
return await asyncio.gather(*tasks, return_exceptions=True)Backoff Eksponensial saat Terjadi Error Batas Laju
Bahkan dengan semaphore, Anda mungkin mencapai batas laju saat terjadi lonjakan lalu lintas. Terapkan backoff eksponensial: tunggu 1 detik, lalu 2, 4, dan 8 detik sebelum mencoba lagi. Tambahkan jitter (offset acak kecil) agar semua pemanggil yang berjalan bersamaan tidak mencoba lagi tepat pada waktu yang sama, yang dapat menyebabkan lonjakan lain. Pustaka tenacity memudahkan penerapan ini.
from tenacity import retry, wait_exponential, stop_after_attempt, retry_if_exception_type
from openai import RateLimitError
@retry(
wait=wait_exponential(multiplier=1, min=1, max=60),
stop=stop_after_attempt(5),
retry=retry_if_exception_type(RateLimitError)
)
async def extract_with_retry(text: str):
return await extract_one(text)Menggunakan Antrean Tugas untuk Kumpulan Besar
Untuk kumpulan yang lebih besar daripada beberapa ribu item, antrean tugas persisten lebih baik daripada asyncio.gather. Antrean seperti Redis Queue (RQ), Celery, atau Dramatiq mempertahankan tugas saat aplikasi dimulai ulang, memungkinkan penskalaan horizontal dengan beberapa pekerja, serta memberi Anda visibilitas terhadap status dan kegagalan tugas. Pekerja mengambil tugas dari antrean dan memanggil API secara mandiri.
# With Redis Queue (RQ)
from rq import Queue
from redis import Redis
redis_conn = Redis()
q = Queue('extractions', connection=redis_conn)
def enqueue_documents(doc_ids: list):
for doc_id in doc_ids:
q.enqueue(
'workers.extract_document',
doc_id,
job_timeout=120,
result_ttl=3600
)
enqueue_documents(all_doc_ids)Melacak Kemajuan dengan Basis Data
Tugas kumpulan yang berjalan lama memerlukan pelacakan kemajuan agar Anda dapat memantau status, mengidentifikasi tugas yang macet, dan melanjutkan setelah kegagalan. Gunakan tabel status di basis data Anda dengan kolom untuk ID dokumen, status (menunggu, diproses, selesai, gagal), stempel waktu, dan pesan error. Perbarui status secara atomik di sekitar setiap panggilan ekstraksi.
import asyncpg
async def process_document(pool, doc_id: str, text: str):
async with pool.acquire() as conn:
await conn.execute(
'UPDATE extractions SET status=$1, started_at=NOW() WHERE doc_id=$2',
'processing', doc_id
)
try:
result = await extract_one(text)
await conn.execute(
'UPDATE extractions SET status=$1, result=$2, completed_at=NOW() WHERE doc_id=$3',
'completed', result.model_dump_json(), doc_id
)
except Exception as e:
await conn.execute(
'UPDATE extractions SET status=$1, error=$2 WHERE doc_id=$3',
'failed', str(e), doc_id
)Melanjutkan Tugas yang Gagal
Tugas kumpulan harus dapat dimulai ulang dengan aman. Saat memulai, kueri basis data untuk mencari dokumen dengan status pending atau failed, lalu coba lagi. Gunakan kunci idempotensi untuk setiap dokumen. Dengan begitu, jika dokumen yang sama tidak sengaja dimasukkan ke antrean dua kali, percobaan kedua akan mendeteksi hasil yang sudah selesai dan melewati pemrosesan ulang. Ini mencegah penulisan duplikat ke sistem hilir.
async def get_pending_docs(pool) -> list:
async with pool.acquire() as conn:
rows = await conn.fetch(
'SELECT doc_id, raw_text FROM extractions WHERE status IN ($1, $2)',
'pending', 'failed'
)
return [dict(row) for row in rows]
async def resume_batch(pool):
docs = await get_pending_docs(pool)
print(f'Resuming {len(docs)} unprocessed documents')
tasks = [process_document(pool, d['doc_id'], d['raw_text']) for d in docs]
await asyncio.gather(*tasks, return_exceptions=True)Mengelompokkan Permintaan dengan Batch API OpenAI
Batch API OpenAI memungkinkan Anda mengirim hingga 50.000 permintaan dalam satu file dan menerima hasilnya dalam 24 jam dengan diskon 50%. Fitur ini ideal untuk pipeline ekstraksi yang tidak mendesak dan lebih mengutamakan biaya daripada latensi. Anda mengunggah file JSONL berisi permintaan, memeriksa status penyelesaian secara berkala, lalu mengunduh file hasilnya.
from openai import OpenAI
import json
client = OpenAI()
# Build JSONL batch file
with open('/tmp/batch_requests.jsonl', 'w') as f:
for i, text in enumerate(documents):
request = {
'custom_id': f'doc_{i}',
'method': 'POST',
'url': '/v1/chat/completions',
'body': {
'model': 'gpt-4o-mini',
'messages': [{'role': 'user', 'content': text}]
}
}
f.write(json.dumps(request) + '\n')
# Upload and submit
batch_file = client.files.create(file=open('/tmp/batch_requests.jsonl', 'rb'), purpose='batch')
batch = client.batches.create(input_file_id=batch_file.id, endpoint='/v1/chat/completions', completion_window='24h')
print(batch.id)Memantau Throughput dan Biaya
Lacak throughput ekstraksi (dokumen per menit) dan biaya per dokumen selama pemrosesan kumpulan. Bagi total pengeluaran API dengan jumlah dokumen yang diproses untuk mendapatkan tolok ukur biaya. Saat melakukan penskalaan, perhatikan pertumbuhan biaya linear — pertumbuhan superlinear menunjukkan bahwa Anda membuang token untuk prompt yang terlalu panjang dan tidak perlu. Dasbor metrik sederhana membantu Anda menemukan ketidakefisienan sebelum dampaknya semakin besar.
import time
class BatchMetrics:
def __init__(self):
self.start_time = time.time()
self.processed = 0
self.total_tokens = 0
self.cost = 0.0
def record(self, usage):
self.processed += 1
self.total_tokens += usage.total_tokens
self.cost += usage.prompt_tokens * 0.00000015 + usage.completion_tokens * 0.0000006
def report(self):
elapsed = time.time() - self.start_time
print(f'{self.processed} docs in {elapsed:.1f}s = {self.processed/elapsed:.1f} docs/sec')
print(f'Cost: ${self.cost:.4f} = ${self.cost/self.processed:.6f} per doc')Mematuhi Batas Laju Token per Menit
Batas laju OpenAI berlaku untuk permintaan per menit (RPM) dan token per menit (TPM). Mengirim 50 panggilan secara bersamaan tidak masalah untuk RPM, tetapi jika setiap panggilan menggunakan 2.000 token, 50 panggilan sama dengan 100.000 token per menit — dengan mudah melampaui batas Tingkat 1. Hitung perkiraan token sebelum mengirim permintaan menggunakan tiktoken, lalu terapkan anggaran token bersama semaphore konkurensi Anda.
import tiktoken
enc = tiktoken.encoding_for_model('gpt-4o-mini')
def estimate_tokens(text: str) -> int:
return len(enc.encode(text)) + 300 # +300 for schema + response
# TPM_LIMIT = 200_000 # Tier 2 limit
# Only submit a batch if estimated total tokens fits within budget
def fits_in_budget(texts: list, tpm_limit: int = 200_000) -> bool:
total = sum(estimate_tokens(t) for t in texts)
return total <= tpm_limitPemeriksaan Singkat
Uji pemahaman Anda tentang pemrosesan kumpulan asinkron untuk ekstraksi dokumen.
Ringkasan Pelajaran
Dalam pelajaran ini Anda mempelajari bahwa asyncio dan Semaphore memungkinkan panggilan API secara bersamaan sekaligus mematuhi batas laju, antrean tugas dan tabel status membuat tugas kumpulan besar dapat dilanjutkan dan dipantau, serta Batch API OpenAI menawarkan penghematan biaya 50% untuk beban kerja yang tidak mendesak dengan konsekuensi latensi 24 jam. Selanjutnya, kita akan menangani evolusi skema dalam pipeline ekstraksi yang berjalan lama.
Pertanyaan yang Sering Diajukan
Apakah pelajaran “Pemrosesan Kelompok dengan Asinkron dan Antrean” gratis?
Ya — teks lengkap “Pemrosesan Kelompok dengan Asinkron dan Antrean” gratis dibaca di sini di web. Untuk praktiknya secara interaktif (editor kode bawaan dan tutor AI 24/7) dan buka sisa kursus AI Engineering Academy, upgrade ke CoddyKit PRO. Kursus AI Engineering Academy mencakup 4 pelajaran total.
Apa yang akan aku pelajari di “Pemrosesan Kelompok dengan Asinkron dan Antrean”?
Bangun alur pemrosesan ekstraksi asinkron menggunakan asyncio dan antrean tugas untuk memproses ribuan dokumen secara paralel, dengan tetap mematuhi batas laju dan melacak kemajuan. Kamu berlatih AI Engineering Academy dengan kode praktik yang langsung kamu jalankan di browser, dan tutor AI 24/7 menjawab pertanyaanmu saat kamu mengerjakan pelajaran ini.
Apakah aku perlu pengalaman untuk memulai AI Engineering Academy?
Tidak diperlukan pengalaman sebelumnya. AI Engineering Academy di CoddyKit dirancang untuk pemula hingga pelajar tingkat lanjut, jadi kamu bisa memulai di sini atau dari awal dan belajar sesuai kecepatan kamu sendiri. Ini adalah pelajaran 3 dari 4.
Berapa lama pelajaran “Pemrosesan Kelompok dengan Asinkron dan Antrean” memakan waktu?
Sebagian besar pelajaran CoddyKit memakan waktu sekitar 5–10 menit. Setiap pelajaran ringkas dan interaktif, jadi kamu membuat kemajuan stabil dan melanjutkan dari tempat kamu tinggalkan di web dan aplikasi.
Bisakah aku menulis dan menjalankan kode dalam pelajaran AI Engineering Academy ini?
Ya. Setiap pelajaran AI Engineering Academy menyertakan editor kode bawaan, jadi kamu menulis dan menjalankan kode nyata langsung di browser dan mendapatkan umpan balik AI instan — tidak diperlukan penyiapan lokal.
Semua pelajaran dalam kursus ini
- Instructor: Ekstraksi Bertipe dengan Pydantic
- Menangani Data Parsial dan Hilang
- Pemrosesan Kelompok dengan Asinkron dan Antrean
- Evolusi Skema dan Kompatibilitas Mundur