Apache NiFi Data Provenance: Trace RAG Chunks to the Source

Apache NiFi Data Provenance: Trace RAG Chunks to the Source

Someone on the support team pastes a screenshot into Slack. The company chatbot told a customer the refund window is 60 days. It's 30, and it has been 30 since March. The question lands on your desk: which document did that answer come from, and why is it still in the index?

With a pile of cron jobs and Python scripts, the honest answer is "give me a day." If it runs through Apache NiFi, the answer is already sitting on disk. NiFi records a provenance event every time a piece of data is received, split, modified, routed or sent. Here's how that record works and how to use it to trace a RAG chunk back to its source file. Everything here was checked against NiFi 2.12.0 (released September 13) and uses core processors only, since NiFi's official Python AI processors haven't shipped a release since November 2024.

You'll need NiFi 2.x and a flow that fetches documents, splits them into chunks and sends them to a vector store.

FlowFiles and Events in One Minute

Every piece of data in NiFi is a FlowFile: content plus attributes, with a UUID. Each time a processor touches one, NiFi writes an event with the FlowFile UUID, the processor ID, the event type, a timestamp, the attributes before and after, and pointers to the content before and after.

The event types you'll see most in an ingestion flow:

Event typeWhen it firesRAG example
RECEIVE / FETCHData enters the flow from outsideFetchS3Object pulls refund-policy.pdf
CONTENT_MODIFIEDA processor rewrites the contentText extraction or cleanup
ATTRIBUTES_MODIFIEDOnly attributes changeHashing, tagging with a version
FORKOne FlowFile becomes many childrenSplitText cuts the doc into chunks
SENDData leaves NiFiInvokeHTTP posts a chunk to the embedding and vector store API
DROPThe FlowFile's journey endsChunk finished, or filtered out

FORK is what makes lineage possible. When SplitText turns one document into 40 chunks, each chunk gets a new UUID, and the FORK event stores the parent UUID next to all 40 children. Follow those links backwards and you get from any chunk to the file it was cut from.

Where the Events Actually Live

Events go to the provenance repository, set in conf/nifi.properties. Here are the 2.12.0 defaults that matter:

nifi.provenance.repository.implementation=org.apache.nifi.provenance.WriteAheadProvenanceRepository
nifi.provenance.repository.directory.default=./provenance_repository
nifi.provenance.repository.max.storage.time=30 days
nifi.provenance.repository.max.storage.size=10 GB
nifi.provenance.repository.indexed.fields=EventType, FlowFileUUID, Filename, ProcessorID, Relationship
nifi.provenance.repository.indexed.attributes=

Events are appended to journal files and indexed with Lucene, so searches don't scan raw logs. Two things to know before you rely on it.

First, retention is a rolling window. Whichever limit hits first (30 days or 10 GB) wins, and older events are gone. A busy flow can burn through 10 GB in a week.

Second, the event only points at content. The bytes live in the content repository, and old content is kept for nifi.content.repository.archive.max.retention.period, which defaults to 3 hours. After that, you still have the metadata and the lineage, but you can't view or replay the actual file.

Make Your Documents Searchable by a Stable ID

Out of the box you can search by filename and FlowFile UUID. Filenames repeat ("policy.pdf" exists in every folder) and UUIDs change at every split, so neither is a great key. What you want is an ID that survives the whole pipeline and also ends up in your vector store.

Add a CryptographicHashContent processor right after the fetch, set to SHA-256. It writes the attribute content_SHA-256. Then rename it with UpdateAttribute (doc.sha256 = ${content_SHA-256}) and tell NiFi to index it:

nifi.provenance.repository.indexed.attributes=doc.sha256

Restart NiFi after changing this. Only events recorded after the restart are indexed on that attribute.

Children inherit their parent's attributes, so every chunk out of SplitText carries doc.sha256, plus fragment.index and segment.original.filename. When you write chunks to your vector store, put doc.sha256 and the chunk's uuid in the metadata. That's the bridge between "this chunk was retrieved" and "this is how it got here."

Tracing a Chunk Back to Its Source

Back to the refund chatbot. Your RAG app logs which chunks it retrieved, and the metadata gives you doc.sha256. Now ask NiFi. First, get a token (NiFi 2.x runs HTTPS with single-user credentials by default):

NIFI=https://your-nifi.example.com
TOKEN=$(curl -s -X POST "$NIFI/nifi-api/access/token" \
  --data-urlencode "username=$NIFI_USER" --data-urlencode "password=$NIFI_PASSWORD")

Then submit a provenance query. Search terms use field IDs for built-in fields and the attribute name for indexed attributes:

curl -s -X POST "$NIFI/nifi-api/provenance" \
  -H "Authorization: Bearer $TOKEN" -H "Content-Type: application/json" \
  -d '{"provenance":{"request":{"maxResults":100,
       "searchTerms":{"doc.sha256":{"value":"3f9a...e1c2","inverse":false}}}}}'

Queries are asynchronous. The response contains an id. Poll GET /nifi-api/provenance/{id} until finished is true, read the events, then DELETE /nifi-api/provenance/{id} to free the query on the server. You'll see every event recorded after the hash was added: the FORK into chunks and one SEND per chunk to your vector store. The FETCH happened before the hash existed, so for that step you need lineage.

For the full graph, ask for lineage from the chunk's UUID:

curl -s -X POST "$NIFI/nifi-api/provenance/lineage" \
  -H "Authorization: Bearer $TOKEN" -H "Content-Type: application/json" \
  -d '{"lineage":{"request":{"lineageRequestType":"FLOWFILE","uuid":"<chunk-uuid>"}}}'

Same async pattern: poll GET /nifi-api/provenance/lineage/{id}, then delete it. In the UI, it's the lineage icon next to any event in the Data Provenance view.

For the refund bug, the trail would typically end somewhere boring: the FETCH event shows an S3 key like archive/2025/refund-policy.pdf, because the List processor watches the whole bucket and an old copy got ingested again. That's ten minutes with provenance instead of an afternoon of grepping logs.

Replay: Re-Run a Document Without Re-Fetching It

Once you've fixed the flow (say, excluding archive/), you can push a specific event's content through again: POST /nifi-api/provenance-events/replays with {"eventId": 1234}, or the Replay button in the UI. It re-enters the flow at the same processor, with the same content, and gets its own REPLAY event, so the audit trail stays honest.

The catch is the content archive. Replay needs the bytes, and with a 3-hour default retention, yesterday's document can't be replayed. If replay matters to you, raise nifi.content.repository.archive.max.retention.period and make sure the disk can hold it.

Running It on Elestio

You can deploy Apache NiFi as a fully managed service on Elestio. NiFi is Apache 2.0 licensed, so you only pay for the VM. I'd start at 2 vCPU / 4 GB of RAM (about $16/month on Netcup), since the JVM and Lucene indexing both want memory, and go to 4 vCPU / 8 GB (about $29/month) once you index attributes on a busy flow. Updates, backups and monitoring are handled for you.

Troubleshooting

Search by doc.sha256 returns nothing. The attribute wasn't indexed when those events were written. Check indexed.attributes, restart, and process a new file. Older events stay unsearchable by that attribute.

Lineage stops at the chunk. The split happened outside NiFi (a script that wrote chunks back to disk, for example). Do the chunking inside the flow so FORK links parent and children.

"Content is no longer available" on replay or view. The content archive expired. Raise the retention period, or rely on your source bucket for re-ingestion.

Events disappear after a few days. You're hitting max.storage.size before max.storage.time. Raise the size, or stream events out with SiteToSiteProvenanceReportingTask if you need months of history.

Queries are slow. Keep indexed.attributes to the few you actually search by. Every indexed attribute costs disk and CPU on every event.

When the support screenshot lands next time, you won't need a day. You'll need a hash, two API calls, and a coffee.

Thanks for reading ❤️ See you in the next one 👋