Skip to content
Notifications
Clear all

How do I batch-process 1000 support emails with ChatGPT without hitting limits?

1 Posts
1 Users
0 Reactions
25 Views
(@alexg)
Honorable Member
Joined: 3 months ago
Posts: 564
Topic starter   [#15627]

The core challenge here isn't the LLM's capability, but designing a resilient pipeline that respects token economics and API rate limits. Simply looping through a CSV will failβ€”spectacularly and expensively. You need a system that handles batching, retries, state persistence, and cost tracking.

First, architect your workflow. Don't send raw emails; pre-process and chunk intelligently.
1. **Ingestion & Chunking:** Extract email bodies, strip HTML, and remove signatures/threads. Use a character count heuristic (e.g., 3000 chars) to split long emails, preserving semantic boundaries.
2. **Batch Queue:** Place processed chunks into a queue (like SQS, Redis, or even a PostgreSQL `SKIP LOCKED` table). This decouples processing from ingestion.
3. **Orchestrator:** A controller service that pulls from the queue, manages concurrent API calls, and respects RPM/TPM limits. You *must* implement exponential backoff on `429` errors.

Here's a minimal Python sketch for the orchestrator using tenacity for retries and `asyncio` for concurrency control. Assume your queue client is abstracted.

```python
import asyncio
import backoff
from openai import AsyncOpenAI, RateLimitError

client = AsyncOpenAI(api_key="your_key")
SEMAPHORE = asyncio.Semaphore(10) # Control concurrent requests

@backoff.on_exception(backoff.expo, RateLimitError, max_tries=5)
async def process_email_chunk(chunk_id, text, system_prompt):
async with SEMAPHORE:
try:
response = await client.chat.completions.create(
model="gpt-4-turbo-preview",
messages=[
{"role": "system", "content": system_prompt},
{"role": "user", "content": f"Analyze this support email: {text}"}
],
temperature=0.2
)
# Store result with chunk_id in your DB
return chunk_id, response.choices[0].message.content
except RateLimitError as e:
# Log for FinOps tracking
raise e

async def main(email_chunks):
tasks = [process_email_chunk(c['id'], c['text'], sys_prompt) for c in email_chunks]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Handle exceptions and consolidate results per original email
```

**Critical Considerations:**
* **Cost:** With 1000 emails, even at ~$0.01 per 1K input tokens, costs can balloon if emails are long. Pre-calculate estimated tokens using `tiktoken` and set a budget alert.
* **State Management:** You must track which chunks have been processed. Use a database with a status field (`pending`, `processed`, `failed`).
* **Idempotency:** Make retries safe. If a request fails after API processing but before your storage, you could be charged for a duplicate. Use idempotency keys if your provider supports them.
* **Observability:** Log token counts, latency, and cost per request. This is non-negotiable for optimization.

The naive approach will hit rate limits and lose data. Treat this as a distributed systems problem, not a simple script. The code to call the API is trivial; the infrastructure to do it reliably at scale is not.

-- alex



   
Quote