A retrieval-augmented assistant is a data product. Its answers are derived from documents that passed through a parser, a chunker, an embedding model and an index, and each of those steps has versions. When a user asks why the assistant quoted last year's leave policy, or legal asks you to remove a customer's documents everywhere, you need to answer questions that ordinary search logs cannot: which source produced this chunk, which index version served it, which answers used it, and where else it lives.
Two disciplines from data engineering answer those questions. Data contracts state what a source promises, such as schema, meaning, freshness and sensitivity, and are enforced at the boundary. Lineage records how every derived artefact was produced from its inputs. This article applies both to AI knowledge systems, with an identity scheme, a contract format, a validator, lineage tables, a worked example of an update and a deletion, and a safe embedding-model migration.
Why knowledge systems need lineage
In a warehouse, lineage connects tables and columns. In a knowledge system the units are smaller and more numerous: a single policy document becomes dozens of chunks, each chunk becomes one embedding per model, and each answer cites several chunks. The pipeline also changes more often than a warehouse: chunkers are tuned, embedding models are replaced, prompts change weekly.
Without lineage, four routine operations become guesswork. Debugging: an answer is wrong, and you cannot tell whether the source was wrong, the chunk was cut badly or a stale index served it. Impact analysis: a source team fixes an error in a document, and you cannot tell which answers relied on the wrong version. Deletion: a data subject or a contract requires removal, and copies remain in an index replica, a cache or an evaluation set. Migration: you change the embedding model and cannot prove the new index covers the same documents as the old one.
The architecture
The diagram shows the shape. Every source passes a contract gate before any processing; failures go to a quarantine with a reason, not into the index. Each pipeline job emits run events to a lineage store and writes identity rows that map document versions to chunks and chunks to embeddings in specific index versions. The answer service logs the chunk ids it cited and the index version it read. Impact analysis, deletion and audit are then queries over those tables.
Nothing here requires a particular product. The run-level events can follow the OpenLineage model, in which each job run emits a START event and a terminal COMPLETE, ABORT or FAIL event naming its input and output datasets, with extra metadata attached as typed facets. That lets a catalogue that already understands OpenLineage show the knowledge pipeline next to your warehouse jobs. The fine-grained identity tables are yours to design.
Identity first: ids that survive reprocessing
Lineage is only as good as its identifiers. Random ids assigned at index time break every link the moment you reprocess. Derive ids from content and configuration instead, so the same input and the same code always produce the same id, and any change produces a new one.
from hashlib import sha256
def h(*parts: str, n: int = 20) -> str:
return sha256("|".join(parts).encode("utf-8")).hexdigest()[:n]
def doc_version(doc_id: str, raw: bytes) -> str:
return h(doc_id, sha256(raw).hexdigest()) # changes when content changes
def chunk_id(doc_ver: str, chunker: str, start: int, end: int) -> str:
return h(doc_ver, chunker, str(start), str(end)) # changes when chunking changes
def embedding_key(chunk: str, model_id: str) -> str:
return h(chunk, model_id) # one vector per chunk per modelWith these ids, re-running ingestion on unchanged documents is a no-op, a changed document produces new chunk ids while the old ones remain traceable, and switching chunker or model is visible in every id. Store the stable doc_id from the source system as well, because deletion requests arrive in terms of source ids, not hashes.
What a knowledge contract should promise
A contract is an agreement between the team that owns a source and the team that consumes it, written as a file both can read and a machine can check. For knowledge systems it needs more than a schema. The example below is this site's illustrative format; the Open Data Contract Standard (ODCS, at major version 3) is a vendor-neutral format you can map these ideas onto if you want tooling to read it.
contract: hr-policies-knowledge
version: 2.1.0
owner: people-ops-data@company.example
consumers: [hr-assistant, onboarding-agent]
source:
system: wiki-space-HR
classification: internal
pii_allowed: false
schema:
required: [doc_id, title, body, effective_date, policy_owner, region]
region: {enum: [global, eu, us, in]}
effective_date: {type: date}
quality:
min_body_chars: 400
max_future_effective_days: 90
freshness:
max_lag_hours: 24
processing:
chunker: heading-aware-v4
embedding_model: embed-model-a-1024
retention:
delete_propagation_hours: 72Each block exists because of a failure mode. Classification and PII stop a sensitive document reaching an index that many users can query. Semantic fields such as effective_date and region let retrieval filter out superseded or irrelevant policies, which is the root cause of many wrong answers. Freshness makes staleness a contract breach with an owner rather than a surprise. Processing pins the chunker and model, so a change is a versioned contract change. Retention turns deletion into a measured promise.
Enforce the contract at the gate
The gate runs before parsing. It is deliberately boring: required fields, allowed values, simple quality rules and a PII scan. A failure quarantines the document with a reason and notifies the owner; it never silently drops it.
import datetime as dt
def check(doc: dict, c: dict, pii_scan) -> list[str]:
errors = [f"missing {f}" for f in c["schema"]["required"] if not doc.get(f)]
if doc.get("region") not in c["schema"]["region"]["enum"]:
errors.append(f"region {doc.get('region')!r} not allowed")
if len(doc.get("body", "")) < c["quality"]["min_body_chars"]:
errors.append("body too short")
eff = doc.get("effective_date")
limit = dt.date.today() + dt.timedelta(days=c["quality"]["max_future_effective_days"])
if eff and eff > limit:
errors.append("effective_date too far in the future")
if not c["source"]["pii_allowed"] and pii_scan(doc.get("body", "")):
errors.append("PII found in a source that must not contain it")
return errors
def gate(doc, contract, pii_scan, quarantine, lineage):
errors = check(doc, contract, pii_scan)
lineage.record_check(doc["doc_id"], contract["version"], passed=not errors, errors=errors)
if errors:
quarantine.put(doc, errors)
return False
return TrueVersion the contract like an API. Adding an optional field is a minor version; removing a field, changing an enum or changing the embedding model is a major version, which consumers must accept before it takes effect. Contract changes go through review by the owning team and the consumers, exactly like a change to a shared schema in a registry; schema registries for streaming describe the same compatibility thinking for events.
Record lineage at two grains
Run-level lineage records jobs: this ingestion run read these sources and wrote this index version. It is small, and it belongs in the same catalogue as the rest of your data platform. Chunk-level lineage records individual artefacts and is far too voluminous for per-item events, so keep it in tables.
CREATE TABLE chunk_lineage (
chunk_id text PRIMARY KEY,
doc_id text NOT NULL,
doc_version text NOT NULL,
source_uri text NOT NULL,
contract_ver text NOT NULL,
chunker text NOT NULL,
span_start int, span_end int,
run_id uuid NOT NULL,
deleted_at timestamptz
);
CREATE TABLE embedding_lineage (
chunk_id text REFERENCES chunk_lineage,
model_id text NOT NULL,
index_version text NOT NULL,
run_id uuid NOT NULL,
PRIMARY KEY (chunk_id, model_id, index_version)
);
CREATE TABLE answer_log (
answer_id uuid PRIMARY KEY,
asked_at timestamptz NOT NULL,
index_version text NOT NULL,
prompt_version text NOT NULL,
model_id text NOT NULL,
cited_chunks text[] NOT NULL
);A run-level event for the embedding job then looks like this, abridged: the specification also requires a schemaURL naming the spec version your client emits, omitted here. Dataset names match the index versions in the tables:
{
"eventType": "COMPLETE",
"eventTime": "2026-09-29T02:14:07Z",
"run": {"runId": "7c1e0f4a-0b7d-4a55-9a61-2f3f0d2c9b10"},
"job": {"namespace": "knowledge", "name": "embed_hr_policies"},
"inputs": [{"namespace": "knowledge", "name": "chunks.hr_policies"}],
"outputs": [{"namespace": "knowledge", "name": "index.hr_policies.v41"}],
"producer": "https://git.company.example/knowledge-pipeline"
}Write the lineage rows in the same transaction as the artefact where you can, or through an outbox when the artefact lives in another system such as a vector database. Lineage that is written best-effort after the fact drifts, and drifting lineage is worse than none because people trust it. The transactional outbox pattern is the standard way to keep the two consistent.
Worked example: a policy fix and a deletion
On Monday the HR team corrects the parental-leave policy: the EU entitlement in the published page was wrong. The new version passes the gate, gets a new doc_version, is chunked and embedded, and index version v42 goes live. Because the answer log records cited chunks, you can find every answer that relied on the wrong text:
SELECT a.answer_id, a.asked_at, a.index_version
FROM answer_log a
JOIN chunk_lineage c ON c.chunk_id = ANY (a.cited_chunks)
WHERE c.doc_id = 'hr-policy-parental-leave'
AND c.doc_version = :wrong_version
ORDER BY a.asked_at;Support can now contact the affected users, and the evaluation team can add those questions to the regression set. On Wednesday a different request arrives: a contractor's personal document was uploaded to a shared space by mistake and must be removed within the contract's 72 hours. The deletion job walks lineage forward from the doc_id: every chunk_id for every version, every embedding in every live index version, plus any cache and evaluation set keyed by chunk id. It deletes in each store, sets deleted_at, and then verifies by querying each index for the chunk ids. The job's completion time against the contract's promise is a metric, not a hope.
Migrating an embedding model under contract
Changing embedding models is the riskiest routine change, because vectors from different models cannot be mixed in one index. Lineage makes it a controlled migration. Build index v50 with the new model beside the live index, and check coverage with a query: every live chunk in v49 must have an embedding in v50. Run the evaluation golden set against both. Shift a slice of traffic, compare answer quality and citation rates, then cut over and keep v49 for a rollback window. The answer log records which index served each answer, so a quality dip can be attributed to the migration rather than argued about. Longer-lived agent memories built on vector stores need the same care; see memory and vector stores in the agentic stack.
Failure modes
- Random ids. Reprocessing breaks every lineage link. Derive ids from content and configuration.
- Best-effort lineage. Rows written after the artefact, outside a transaction, drift. Use the same transaction or an outbox.
- Uncited answers. If the answer service does not log chunk ids, impact analysis stops at the index. Make citation logging part of the response path.
- Forgotten copies. Caches, evaluation sets, fine-tuning exports and replicas hold chunks outside the main index. Register every consumer in the contract and include it in deletion.
- Contract theatre. Contracts that nobody enforces become stale documentation. The gate must block, and breaches must page the owner.
- Silent superseding. Old and new versions of a policy both stay retrievable. Mark superseded versions and filter by effective date at query time.
Trade-offs and operations
Chunk-level lineage costs storage and write latency, roughly a row per chunk and per embedding, which is small next to the vectors themselves. Answer logging raises privacy questions of its own: store chunk ids and versions, not full prompts, unless you have a reason and a retention policy. Contract gates slow ingestion slightly and occasionally block a document someone wanted indexed urgently; that is their purpose, so provide an expedited review path rather than a bypass.
Track contract breaches by source and owner, quarantine age, ingestion lag against the freshness promise, deletion completion time against its promise, the share of answers with logged citations, and index coverage after each migration. For the wider retrieval pipeline see RAG pipeline design at scale and hybrid search. A companion article in this category, on graph RAG versus vector search, applies the same evidence discipline to graph edges.
What to do next
- Switch to content-derived ids for documents, chunks and embeddings, and keep the source system's own id alongside them.
- Write a contract for your most-used source with owner, classification, required fields, freshness, pinned chunker and model, and retention.
- Put a blocking gate with a quarantine in front of the chunker.
- Create the chunk, embedding and answer-log tables, and write them transactionally or through an outbox.
- Emit run-level events for each pipeline job into the catalogue your data platform already uses.
- Write and rehearse the impact query and the deletion job, and measure deletion time against the contract.
- Use coverage queries and the answer log for your next embedding migration.