Project · Batch Document Pipeline
The capstone for this track: a resumable pipeline that classifies and summarizes a large folder of documents with the OpenAI Batch API — at half the cost of realtime. You'll tie together files, batches, polling, result-matching by custom_id, per-item error handling, and persistence into one production-shaped job.
Learning objectives
- Turn a folder of documents into a JSONL batch of structured-extraction requests.
- Submit, poll, and collect a batch end to end.
- Match results by
custom_idand handle per-item failures. - Persist state so the pipeline is resumable, not all-or-nothing.
OPENAI_API_KEY) + pip install openai pydantic. The structure runs offline; the batch submission needs the key. Builds directly on OP1 · Batch API.1 · The brief essential
The scenario: you have thousands of support tickets (or contracts, or research PDFs) sitting in a folder, and you need a structured record for each — category, priority, a one-line summary. No human is waiting; it just needs to be done by tomorrow, cheaply. This is the textbook Batch API job, and building it end to end is how the recipes from OP1 become a real pipeline.
The deliverable is a script you can run, kill, and re-run without redoing finished work — because a batch over thousands of items that falls over at item 2,000 must resume, not restart. That resumability is what separates a demo loop from a pipeline.
custom_id, submitted and polled, then collected and matched back into structured records — with state persisted so the job survives a restart.
2 · Step 1 — the schema & the JSONL builder intermediate
Define the structured record you want per document, then turn the folder into a JSONL file — one request per document, custom_id set to the file path so results map straight back.
build.pyimport json, pathlib
# the structured record we want per document
SCHEMA = {"type": "object", "properties": {
"category": {"type": "string"},
"priority": {"type": "string", "enum": ["low","medium","high"]},
"summary": {"type": "string"},
}, "required": ["category","priority","summary"]}
SYSTEM = "Extract category, priority, and a one-line summary from the document."
def build_jsonl(folder, out="requests.jsonl"):
with open(out, "w") as f:
for path in pathlib.Path(folder).glob("*.txt"):
f.write(json.dumps({
"custom_id": str(path), # path = the key back
"method": "POST", "url": "/v1/responses",
"body": {
"model": "gpt-5.5",
"instructions": SYSTEM,
"input": path.read_text()[:8000],
},
}) + "\n")
return out
- Each document becomes one JSONL line;
custom_id = str(path)means the result carries the file path, so matching back is trivial. - The
bodyis a normal Responses request — here extraction viainstructions+ the document asinput(truncated to a sane length). - You could use the full structured-output schema here too; this version keeps the body simple and parses the summary text on collection.
Try this: swap the glob to *.pdf and upload each via the Files API (OP3), referencing file_id in the body — the pipeline shape is identical.
3 · Step 2 — submit & persist the batch id intermediate
Upload, create the batch, and immediately write the batch id to disk. That one line is what makes the job resumable — a fresh process can pick up the id and collect results without re-submitting.
submit.pyimport json
from openai import OpenAI
client = OpenAI()
def submit(jsonl_path, state="state.json"):
up = client.files.create(file=open(jsonl_path,"rb"), purpose="batch")
batch = client.batches.create(
input_file_id=up.id, endpoint="/v1/responses",
completion_window="24h")
json.dump({"batch_id": batch.id}, open(state,"w")) # PERSIST — resumable
return batch.id
batch_id (and later the custom_id→input map) to a small state file is what turns a blocking script into a resumable pipeline. Kill the process after submit; a new run reads state.json and collects — no re-submission, no double spend.4 · Step 3 — poll, collect, match, handle errors advanced
Read the id from state, poll to completion, download the output, and fold each line back by custom_id — successes into records, failures into an errors bucket for re-submission.
collect.pyimport json, time
from openai import OpenAI
client = OpenAI()
def collect(state="state.json"):
batch_id = json.load(open(state))["batch_id"]
while True: # resume-safe: id came from disk
b = client.batches.retrieve(batch_id)
if b.status in {"completed","failed","expired","cancelled"}:
break
time.sleep(30)
records, errors = {}, {}
for line in client.files.content(b.output_file_id).text.splitlines():
row = json.loads(line); cid = row["custom_id"]
if row.get("error"):
errors[cid] = row["error"] # bucket for re-submit
else:
records[cid] = row["response"]["body"]["output_text"]
print(f"{len(records)} ok, {len(errors)} errors")
return records, errors
- The
batch_idcomes fromstate.json, so this runs in a fresh process — the pipeline resumed, not restarted. - Polling stops at any terminal status; then the output JSONL is streamed and each line keyed by
custom_id. - Successes land in
records(keyed by file path), failures inerrors— so one bad document never sinks the run, and you can re-batch just the failures.
Try this: after collection, assert len(records) + len(errors) == input_count — a missing custom_id is a silent data gap, and the assertion makes it loud.
errors bucket to disk, and feed those custom_ids into a fresh JSONL for a second, smaller batch. A production pipeline loops submit→collect→re-batch until the error set is empty or explainable.5 · Tech-lead — making it production-grade tech-lead
The three steps above are the skeleton; production adds the hardening. Chunk huge datasets into multiple batches under the file-size/request-count limits, persisting a row per batch. Decouple submit from collect — submit from one job, poll from a scheduled job, so no process blocks for hours. Idempotency: key everything by custom_id so re-running never double-processes. Observability: log per-batch counts and costs (OP4). The result is a data pipeline that happens to call an LLM — robust, resumable, and cheap by design.
🪜 Practice ladder beginner → industry
- Beginner: run Step 1 on a folder of 3 text files and inspect the JSONL.
- Easy: submit with Step 2 and confirm
state.jsonholds the batch id. - Core: collect with Step 3 in a fresh process and print ok/error counts.
- Stretch: add the reconciliation assertion and deliberately include one bad document.
- Hard: re-batch the errors bucket automatically and merge the second run's results.
- Industry: chunk a 50k-doc dataset into multiple batches with per-batch state and a scheduled collector.
✓ Checkpoint — you've built it when you can…
- Turn a folder into a JSONL batch keyed by
custom_id. - Submit and persist the batch id for resumability.
- Poll, collect, match by
custom_id, and bucket errors. - Reconcile counts and re-batch the failures.
Knowledge check check yourself
What single design choice makes this pipeline resumable rather than all-or-nothing?
Show answer
batch_id (and the custom_id→input map) to a small state file right after submission. Because the batch id lives on disk, a fresh process can read it and collect results without re-submitting, so killing the job mid-run and restarting never redoes finished work or double-spends. Keying everything by custom_id makes the whole thing idempotent.How does the pipeline handle a batch where some documents fail?
Show answer
response or an error, keyed by custom_id, so successes go into a records map and failures into an errors bucket — one bad document never sinks the run. You then reconcile counts against the input (asserting none are silently missing), write the errors bucket to disk, and feed those custom_ids into a fresh, smaller batch, looping until the error set is empty or explained.