Пакетная обработка и асинхронное выполнение
Пакетный API OpenAI, асинхронный Python и параллельное выполнение запросов.
«Пакетная обработка и асинхронное выполнение» — бесплатный урок AI Prompt Engineering на CoddyKit. Это урок 2 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Prompt Engineering, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Prompt Engineering содержит 4 уроков всего.
Зачем нужны пакетная обработка и асинхронное выполнение
Последовательная обработка тысяч запросов к LLM занимает много времени и стоит дорого. Пакетная обработка объединяет запросы и снижает стоимость на 50%. Асинхронное выполнение параллелизирует запросы, чтобы максимально увеличить пропускную способность в пределах ограничений частоты запросов. Вместе эти подходы значительно сокращают и стоимость, и общее фактическое время выполнения.
Пакетный интерфейс OpenAI: снижение стоимости на 50%
Пакетный интерфейс OpenAI асинхронно обрабатывает запросы в фоновом режиме (до 24 часов) по цене, составляющей 50% от обычной стоимости интерфейса. Он идеально подходит для оценочных запусков, обработки наборов данных и задач, не требующих результатов в реальном времени.
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}')Отправка пакетного задания и проверка его состояния
После загрузки входного файла создайте пакетное задание и проверяйте его состояние до завершения. Пакетный интерфейс обрабатывает запросы в течение 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()
Для параллельной обработки в реальном времени (не пакетной) Python использует asyncio и asyncio.gather(), одновременно запуская несколько вызовов программного интерфейса и ожидая их завершения. Это значительно сокращает общее фактическое время выполнения задач с несколькими запросами.
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')Параллельные вызовы с учётом ограничения частоты
Слишком большое количество одновременных запросов вызывает ошибки ограничения частоты. Семафор ограничивает параллелизм, позволяя оставаться в пределах лимитов и одновременно максимально увеличивать пропускную способность.
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
Anthropic также предоставляет программный интерфейс Message Batches с экономикой, аналогичной пакетному интерфейсу 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: снижение стоимости на 50%, фоновая обработка, до 50 тыс. запросов в пакете
- asyncio.gather(): одновременные запросы в реальном времени, все запускаются сразу, после чего выполняется ожидание всех результатов
- Semaphore: управление параллелизмом с учётом ограничений частоты (обычно 10–50 одновременных запросов)
- Экспоненциальная задержка: повторная попытка с удваивающейся задержкой при ошибках ограничения частоты или тайм-аутах
- Устойчивый сбор результатов: возврат исключений предотвращает остановку всего пакета из-за одного сбоя
- Отслеживание прогресса: для больших заданий сообщайте о завершённых операциях, пропускной способности и ETA
Часто задаваемые вопросы
Урок «Пакетная обработка и асинхронное выполнение» бесплатный?
Да — полный текст урока «Пакетная обработка и асинхронное выполнение» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Prompt Engineering, подпишись на CoddyKit PRO. Курс AI Prompt Engineering содержит 4 уроков всего.
Чему я научусь в уроке «Пакетная обработка и асинхронное выполнение»?
Пакетный API OpenAI, асинхронный Python и параллельное выполнение запросов. Ты практикуешь AI Prompt Engineering с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AI Prompt Engineering?
Предыдущий опыт не требуется. AI Prompt Engineering на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 2 из 4.
Сколько времени занимает урок «Пакетная обработка и асинхронное выполнение»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AI Prompt Engineering?
Да. Каждый урок AI Prompt Engineering включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Стратегии кэширования запросов
- Пакетная обработка и асинхронное выполнение
- Балансировка нагрузки между моделями
- Мониторинг и оповещения для конвейеров запросов