AI EngineeringZero to ProductionHome·About·What’s new·Contact
OpenAI API in Practice · Project

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.

⏱️ ~2 hours🏗️ Build-along🎯 Intermediate→Tech-lead

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_id and handle per-item failures.
  • Persist state so the pipeline is resumable, not all-or-nothing.
⚙️ To run this for realNeeds an OpenAI API key (OPENAI_API_KEY) + pip install openai pydantic. The structure runs offline; the batch submission needs the key. Builds directly on OP1 · Batch API.
This pipeline uploads your documents to the OpenAI API and writes state files locally. Run it on data you're permitted to send, and point it at a scratch directory first.

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.

folder → JSONL → batch → poll → collect → structured records (resumable) docs/a folder build JSONLcustom_id=path submit+pollbatches API collect+matchby custom_id records state persisted at each step → kill & resume without redoing work Folder to records, resumably. Documents become a JSONL batch keyed by 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.

Step 1
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
▶ How this works
  1. Each document becomes one JSONL line; custom_id = str(path) means the result carries the file path, so matching back is trivial.
  2. The body is a normal Responses request — here extraction via instructions + the document as input (truncated to a sane length).
  3. 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.

Step 2
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
Persistence is the whole trickWriting the 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.

Step 3
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
▶ How this works
  1. The batch_id comes from state.json, so this runs in a fresh process — the pipeline resumed, not restarted.
  2. Polling stops at any terminal status; then the output JSONL is streamed and each line keyed by custom_id.
  3. Successes land in records (keyed by file path), failures in errors — 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.

Reconcile, then re-batch failuresNever assume a batch came back complete. Reconcile counts against your input, write the 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

  1. Beginner: run Step 1 on a folder of 3 text files and inspect the JSONL.
  2. Easy: submit with Step 2 and confirm state.json holds the batch id.
  3. Core: collect with Step 3 in a fresh process and print ok/error counts.
  4. Stretch: add the reconciliation assertion and deliberately include one bad document.
  5. Hard: re-batch the errors bucket automatically and merge the second run's results.
  6. 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

✓ Knowledge check

What single design choice makes this pipeline resumable rather than all-or-nothing?

Show answer
Persisting state to disk — writing the 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.
✓ Knowledge check

How does the pipeline handle a batch where some documents fail?

Show answer
Each output line carries either a 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.
© 2026 studybydoing.in · AI Engineering: Zero to Production · All rights reserved. · About · Privacy Policy · Terms · Contact
Educational content, provided as-is and without warranty. Code samples are examples — review, test, and adapt them before using in production. See the Terms of Use & Disclaimer. Use at your own risk.
© studybydoing.in