Hands-on: build a queue-based document pipeline with pipes and filters and claim check
A PDF pushed straight into a queue, a processing step that runs twice, a client that resends a POST: all three familiar failures can be handled in a few dozen lines of Python, provided the design is right from the start.
In brief
- The pipe only passes messages along. All logic lives in the filters, and each filter needs to know only its own input and output schemas.
- Large files are written to the store first, and only then is the reference published. The two operations are not atomic, so filters must be idempotent.
- When the customer needs the result within the same request, an asynchronous pipeline is the wrong choice.
- 1API receives the documentChecks the Idempotency-Key, then immediately returns 202 with Location and Retry-After
- 2Write the file to storageOnly once the write succeeds is the reference (claim check) published to the first queue
- 3Validate filterReads the file via the reference; repeated failures go to the dead-letter queue
- 4Extract filterUsually the slowest step, so run several instances in parallel
- 5Store filterIdempotent on job_id, so a duplicate message does not cause a double write
- 6Status endpointReturns 303 See Other once the document has been processed
Only the reference to the file travels through the queues; the file itself stays in storage from start to finish.
Graphic: FDE Times
The detail most often overlooked in Microsoft’s reference documentation on pipes and filters is that the messages passing between steps do not contain the image being processed. Each message holds only a reference (a claim check) pointing to the image sitting in storage.
This keeps messages small, but it does not remove a separate risk: a queue can deliver the same message more than once. So, alongside claim check, every filter must also be idempotent, meaning that receiving a message twice still produces exactly one result. Step 3 deals with this.
On FDE engagements this problem turns up under many names: loan applications, scanned contracts, insurance claims. The customer uploads a heavy file, the file has to pass through several validation and extraction steps, and someone has to know when processing is finished. This guide builds such a pipeline on a laptop using only the Python standard library.
What you will build, and what you need
The pipeline has three steps, validate, extract and store, connected by queues. This is the approach suggested in Microsoft’s best-practice guidance on background jobs: a job that passes through several stages should be split into separate filters connected by queues. Outside the pipeline sits an API that accepts documents and reports processing status.
All you need is Python 3 and an empty directory. Everything below is a simplified simulation: queue.Queue stands in for the message broker, a local directory stands in for blob storage, and the API is written as plain functions with no framework. On a customer site you swap each piece for their real services; the logic stays the same. (In the code, docs is Vietnamese for “application file”.)
Step 1: write the file first, send the ticket second
As Microsoft describes it, claim check means storing a large payload in an external store and sending only a reference through the messaging system. The consumer receives the reference and fetches the payload itself. The order of these two operations matters: the documentation requires that the reference be published only after the payload has been written successfully.
import hashlib, queue, uuid
from pathlib import Path
STORE = Path("blobs"); STORE.mkdir(exist_ok=True)
q_validate, q_extract, q_store, q_dead = (queue.Queue() for _ in range(4))
def put_payload(data: bytes) -> str:
key = hashlib.sha256(data).hexdigest()
(STORE / key).write_bytes(data)
return key
def submit(data: bytes) -> str:
job_id = str(uuid.uuid4())
ref = put_payload(data) # 1. write payload
q_validate.put({"job_id": job_id, "blob": ref, "schema": "docs.v1"}) # 2. publish
return job_id
Check: call submit(b"%PDF-1.4 ...") and look in the blobs/ directory. It should contain exactly one file, while the message in q_validate is only a few dozen bytes.
Using a hash of the content as the key is deliberate: the same file always produces the same key. So if the write is retried, no extra copy is created.
Microsoft notes that the write and the publish do not happen atomically, so you have to account for both duplicate messages and orphaned payloads, meaning files that were written but whose message was never sent.
Step 2: write a common filter template, so the pipe only carries messages
The Azure documentation is explicit: a pipe does no routing and holds no logic; it simply turns the output of one filter into the input of the next. All processing, including error handling, therefore lives in the filters.
def run_filter(inbox, outbox, work, max_attempts=3):
msg = inbox.get()
try:
result = work(msg)
if outbox is not None:
outbox.put({**msg, **result})
except Exception:
msg["attempts"] = msg.get("attempts", 0) + 1
(q_dead if msg["attempts"] >= max_attempts else inbox).put(msg)
The except branch handles poison messages. A message that fails repeatedly is not requeued forever but moved to a dead-letter queue, exactly as the background-job guidance recommends.
The schema field attached to the message has a purpose too: each filter knows only its own input and output schemas, so once schemas are standardised you can reorder the steps without rewriting them.
Step 3: three filters, and the last one must be idempotent
DONE = {}
def validate(msg):
data = (STORE / msg["blob"]).read_bytes()
if not data.startswith(b"%PDF"):
raise ValueError("not a PDF")
return {}
def extract(msg):
data = (STORE / msg["blob"]).read_bytes()
return {"size_bytes": len(data)} # simplified: replace with real OCR/extraction
def store(msg):
if msg["job_id"] in DONE: # already processed, skip
return {}
DONE[msg["job_id"]] = {"blob": msg["blob"], "size_bytes": msg["size_bytes"]}
return {}
(The comments read, in order: “simplified: replace with real OCR/extraction” and “already processed, skip”; the error message means “not a PDF”.)
The reason for the if line in store: a queue guarantees only at-least-once delivery, which means the same message can arrive several times.
Microsoft describes a specific scenario: a filter posts its result and then fails, the message is rerun on another instance, and the result is duplicated.
Any filter that writes to the outside world therefore needs a key to recognise work already done.
Check: run the three filters in turn with run_filter(q_validate, q_extract, validate), run_filter(q_extract, q_store, extract) and run_filter(q_store, None, store). Then manually put the same message into q_store once more and run it again. len(DONE) must still equal 1.
Step 4: find the slowest filter before optimising
The Azure documentation says that the time taken to process a request depends on the slowest filters in the pipeline. Imagine validate takes 1 second, extract 8 seconds and store 1 second.
Optimising validate changes almost nothing. The right move is to run several instances of extract in parallel, all reading from q_extract, and that is a direct benefit of separating the steps with queues.
In the simulation you can create a few threading.Thread instances that all call run_filter on q_extract. Because store is already idempotent, even if two workers happen to process the same message, the result is still correct.
Step 5: the API returns 202, blocks duplicates, and returns 303 when done
The customer should not have to hold a connection open while the pipeline runs. Under the async request-reply pattern, the API validates the request and immediately returns HTTP 202 with Location and Retry-After headers; the client then polls the status endpoint.
JOBS_BY_KEY = {}
def post_docs(body: bytes, idem_key: str):
job_id = JOBS_BY_KEY.get(idem_key)
if job_id is None:
job_id = submit(body)
JOBS_BY_KEY[idem_key] = job_id
return 202, {"Location": f"/docs/status/{job_id}", "Retry-After": "5"}
def get_status(job_id: str):
if job_id in DONE:
return 303, {"Location": f"/docs/{job_id}"}
return 200, {"status": "processing"}
Two details in this code prevent two different failures. The Idempotency-Key header lets the backend return the existing status resource when a client resubmits, instead of putting another work item on the queue.
And when the job is finished, Microsoft recommends 303 See Other rather than 302, because with some clients a 302 can cause the original POST to be resent.
Check: call post_docs twice with the same key. Both calls must return the same Location, and q_validate must grow by only one message.
Final exercise: break validate on purpose
This exercise ties all five steps together. Call post_docs(b"hello", "key-bad"): the file is still written to blobs/ and the API still returns 202, but the content does not start with %PDF, so validate will raise an error.
Now run run_filter(q_validate, q_extract, validate) three times. On the first two runs the message returns to q_validate with attempts incrementing; on the third it lands in q_dead. Confirm that q_dead.qsize() equals 1, q_extract is still empty, and get_status for this job still returns processing.
That last result exposes a gap in the simulation: the client will poll forever without learning that the document failed. The natural next step is to have the status endpoint also read the dead-letter queue and return a failed status.
Common mistakes
The first mistake is publishing the message before the file has been fully written, reversing the order the claim check documentation requires. The first filter may then read a reference that points at nothing.
The second is that nobody owns payload deletion. The claim check documentation requires you to state clearly who owns deletion and retention, and to align the payload’s retention period with the message’s lifetime, so that a still-valid reference never points to data that has already been deleted.
The third concerns security: putting a security token into the claim check for convenience. Microsoft advises against this. The last mistake is applying the pattern in the wrong place. If processing must complete within the original request, request/response style, pipes and filters is not a fit.
How this skill shows up on a customer site
On a customer site the first task is not writing code but asking one question: will users accept receiving the result later, through a status endpoint, or must they have it immediately?
The answer decides whether you use a pipeline at all. Next, map the current steps into filters, mark the slowest one, and ask where files are stored and who is allowed to delete them.
On your CV, do not just write “used a message queue”. Spell out how you handled duplicate delivery, where you placed the dead-letter queue and who was responsible for cleaning up payloads.
The book Enterprise Integration Patterns frames this pattern around a question: how do you process a message through several complex steps while keeping those steps independent and flexible? The exercise above shows the answer: each filter knows only its own input and output schemas, so you can reorder steps or add instances without rewriting any filter.
Was this article useful?
Thanks for the feedback!
5 sources
- Pipes and Filters pattern - Azure Architecture Center | Microsoft Learn · 2024-04-10
- Claim Check pattern - Azure Architecture Center | Microsoft Learn · 2026-08-31
- Best Practices for Background Jobs - Azure Architecture Center | Microsoft Learn · 2026-03-30
- Asynchronous Request-Reply Pattern - Azure Architecture Center | Microsoft Learn · 2026-03-30
- Pipes and Filters - Enterprise Integration Patterns