Feature and Data Pipelines in Production

Document ingestion pipelines for RAG

A RAG ingestion pipeline receives new documents, splits them into chunks, embeds them, and stores them safely, so it can be re-run without ever double-processing a file.

On this page 7
  1. The short answer
  2. The analogy you have already lived
  3. Why it exists
  4. How it works
  5. A real example you have seen
  6. Remember this
  7. What to learn next

One lesson, three depths. Pick the one that fits you today — you can switch any time.

Beginner — No maths. Plain English.

The short answer

A RAG ingestion pipeline receives documents, splits them into pieces, and stores them so an AI system can look things up later.

The analogy you have already lived

A librarian receiving a fresh stack of new books does not toss them straight onto a shelf. Each book is checked in against an accession register, given a catalogue entry, and shelved in the right place. Only then can it actually be found later.

RAG stands for retrieval-augmented generation. It means an AI system that looks things up before answering, like an open-book exam instead of one from memory. The ingestion pipeline is the librarian's check-in process for the documents it will look things up in.

Why it exists

An AI assistant answering questions about a company's policies needs those policies stored somewhere it can search. A raw folder of PDFs and text files is not searchable in the way the assistant needs.

The ingestion pipeline turns raw documents into small, embedded, indexed pieces. They sit ready for the low-latency lookup that happens the instant a user asks a question.

How it works

new document arrives
      |
      v
  split into small chunks (a paragraph or two each)
      |
      v
  turn each chunk into an embedding (a list of numbers)
      |
      v
  store chunk + embedding + where it came from
      |
      v
  record "this file is now ingested" so it is never redone by accident

That last step matters as much as the others. A pipeline that re-runs safely, without creating duplicates, is far easier to trust. Compare that to one that has to be run exactly once, perfectly, forever.

A real example you have seen

A customer support chatbot that can answer "what is your refund policy" correctly. That works because someone's ingestion pipeline already chopped up the refund policy document. It was embedded and stored before a single user ever asked the question.

Remember this

  • Ingestion turns raw documents into small, searchable pieces, ready before anyone asks a question.
  • Re-running the pipeline on files already ingested should do nothing, not duplicate the work.
  • "Chunk, embed, store" is the same three steps whether it is ten documents or ten million.

What to learn next

Developer — Code and libraries.

Setup

No installation needed — this uses only Python's standard library.

An idempotent ingestion pipeline

ingest.py
import hashlib
import sqlite3
from pathlib import Path
from datetime import datetime

# A small docs/ folder to ingest, created here so the script is fully
# self-contained -- in a real pipeline these files already exist on disk.
Path("docs").mkdir(exist_ok=True)
Path("docs/refund-policy.txt").write_text(
    "Refunds are processed within 7 business days of the return being received.\n"
    "Store credit is issued immediately if the customer chooses that option instead of a bank refund.\n"
)
Path("docs/shipping-policy.txt").write_text(
    "Standard shipping takes 5 to 8 business days within the country for most locations.\n"
    "Express shipping is available at checkout and typically arrives within 2 business days.\n"
    "International shipping times vary by destination and customs processing at the border.\n"
)

conn = sqlite3.connect(":memory:")
conn.execute("""
CREATE TABLE chunks (
    chunk_id TEXT PRIMARY KEY, source TEXT, chunk_index INTEGER,
    text TEXT, ingested_at TEXT
)
""")
# The ledger: one row per file already ingested, keyed by a hash of its
# contents. This is what makes re-running the pipeline safe.
conn.execute("""
CREATE TABLE ingested_files (path TEXT PRIMARY KEY, content_hash TEXT, ingested_at TEXT)
""")
conn.commit()


def chunk_text(text, max_chars=120):
    paras = [p.strip() for p in text.split("\n") if p.strip()]
    chunks, current = [], ""
    for p in paras:
        if len(current) + len(p) > max_chars and current:
            chunks.append(current)
            current = p
        else:
            current = (current + " " + p).strip()
    if current:
        chunks.append(current)
    return chunks


def ingest_folder(folder):
    added, skipped = 0, 0
    for path in sorted(Path(folder).glob("*.txt")):
        content = path.read_text()
        content_hash = hashlib.sha256(content.encode()).hexdigest()

        existing = conn.execute(
            "SELECT content_hash FROM ingested_files WHERE path = ?", (str(path),)
        ).fetchone()
        if existing and existing[0] == content_hash:
            skipped += 1
            continue  # already ingested, and the file has not changed

        for i, chunk in enumerate(chunk_text(content)):
            chunk_id = hashlib.sha256(f"{path}:{i}".encode()).hexdigest()[:12]
            conn.execute(
                "INSERT OR REPLACE INTO chunks VALUES (?,?,?,?,?)",
                (chunk_id, path.name, i, chunk, datetime.now().isoformat()),
            )
        conn.execute(
            "INSERT OR REPLACE INTO ingested_files VALUES (?,?,?)",
            (str(path), content_hash, datetime.now().isoformat()),
        )
        added += 1
    conn.commit()
    return added, skipped


added, skipped = ingest_folder("docs")
print(f"first run  -> files newly ingested: {added}, files skipped: {skipped}")
n_chunks = conn.execute("SELECT COUNT(*) FROM chunks").fetchone()[0]
print(f"chunks stored: {n_chunks}")

# Run it again with nothing changed on disk -- an idempotent pipeline
# should do no duplicate work.
added2, skipped2 = ingest_folder("docs")
print(f"\nsecond run -> files newly ingested: {added2}, files skipped: {skipped2}")
n_chunks2 = conn.execute("SELECT COUNT(*) FROM chunks").fetchone()[0]
print(f"chunks stored: {n_chunks2}  (unchanged)")

sample = conn.execute("SELECT source, chunk_index, text FROM chunks LIMIT 2").fetchall()
print("\nsample chunks:")
for s in sample:
    print(" ", s)

The script writes two small sample files into docs/ before ingesting them, so this runs standalone with no files to prepare by hand:

Output
first run  -> files newly ingested: 2, files skipped: 0
chunks stored: 5

second run -> files newly ingested: 0, files skipped: 2
chunks stored: 5  (unchanged)

sample chunks:
  ('refund-policy.txt', 0, 'Refunds are processed within 7 business days of the return being received.')
  ('refund-policy.txt', 1, 'Store credit is issued immediately if the customer chooses that option instead of a bank refund.')

This lesson's chunk_text is a simple paragraph-length splitter, deliberately kept small so the ledger logic stays the main character. Real pipelines usually reach for LangChain or LlamaIndex's chunking utilities, and a real embedding model instead of storing raw text alone.

Line-by-line walkthrough

content_hash is the key idea. It is not enough to check "have I seen this filename before" — a file can change without being renamed. Hashing the actual content means an edited file is correctly re-ingested, and an untouched one is correctly skipped.

INSERT OR REPLACE on chunk_id (itself derived from path and chunk index) means re-ingesting a changed file overwrites its old chunks cleanly, rather than piling new ones on top of stale ones.

Common mistakes

No ledger at all — only re-running the whole pipeline on a schedule. Every run re-embeds every document, even unchanged ones. Embedding is often the most expensive step; skipping unchanged files is the single biggest cost saving available here.

Keying the ledger on filename instead of content. A file that gets edited in place, keeping the same name, would then be silently skipped forever, since the filename never changed.

Chunking without recording where a chunk came from. When the AI system cites an answer, it needs to point back to a real source document, not return text with no origin.

No re-embedding plan when the embedding model changes. This ledger, as written, only detects document changes — it has no idea the embedding model itself changed underneath it. See re-embedding and reindexing for that separate, larger problem.

Try it yourself

Edit refund-policy.txt, changing one sentence, and run ingest_folder a third time. Confirm it reports one file newly ingested (not skipped), and that the old chunks for that file were replaced, not duplicated.

What to learn next

Researcher — Mathematics and papers.

Idempotency as the correctness property

An ingestion pipeline should satisfy $\text{ingest}(\text{ingest}(D)) = \text{ingest}(D)$ for any document set $D$ — running it twice on unchanged input produces the same stored state as running it once. This is the general idempotency requirement behind the content-hash ledger above, and it is what allows a pipeline to be safely re-run after a partial failure, rather than requiring exactly-once execution guarantees from the infrastructure itself.

Chunking strategy affects retrieval quality directly

Chunk size trades off two competing failure modes: chunks too large dilute a specific answer among irrelevant surrounding text, hurting retrieval precision; chunks too small lose context a correct answer depends on, hurting recall of the right meaning. Overlapping chunks (a sliding window with, for example, 20% overlap between consecutive chunks) is a common mitigation for information split awkwardly across a chunk boundary, at the cost of storing and embedding somewhat more text than the source contains.

Semantic chunking — splitting at detected topic boundaries rather than a fixed character count — improves this trade-off further but adds a model call per document at ingestion time, a real cost at scale that a fixed-size splitter avoids.

Incremental ingestion at scale

Beyond a single-file content hash, production pipelines commonly track ingestion as a stream: a change-data-capture feed from the source system (a wiki's edit log, a document store's update stream) rather than a periodic full-folder scan. This converts ingestion from an $O(n)$-per-run scan of every document into an $O(\Delta)$ process reacting only to what changed — the same shift from batch rescanning to event-driven processing argued for in streaming feature aggregation.

Deletion and staleness

An ingestion pipeline that only ever adds is incomplete: a document deleted or retracted at the source needs its chunks removed from the index too, or the retrieval system will keep confidently citing content that no longer exists or is known to be wrong. This requires the ledger to track source deletions, not only source changes — commonly implemented as a periodic reconciliation pass comparing the ledger against the live source list, removing chunks for anything no longer present.

Papers and systems

  • Lewis et al., Retrieval-Augmented Generation for Knowledge-Intensive NLP Tasks, NeurIPS 2020 — the paper that introduced the RAG architecture this pipeline feeds.
  • Karpukhin et al., Dense Passage Retrieval for Open-Domain Question Answering, EMNLP 2020 — foundational work on how passage (chunk) granularity affects retrieval quality.
  • LangChain and LlamaIndex's document-loader and text-splitter documentation are the most current practical reference for chunking strategy trade-offs in production.

What to learn next

What to learn next

These follow on from what you just read.

  • Serving Models in Production

    Serving a model with FastAPI

    Taking a FastAPI model server from working to production-ready means understanding how it handles many requests arriving at once, not repeating the basics of loading and validation.

  • Serving Models in Production

    NVIDIA Triton inference server

    Triton is a dedicated inference server that handles dynamic batching, multiple frameworks and GPU scheduling for you, instead of you hand-rolling that logic inside a web framework.

  • Serving Models in Production

    BentoML

    BentoML is a Python framework that turns a trained model into a packaged, servable API with far less boilerplate than writing every route and Dockerfile by hand.