Queue grain, not raw HTML
After I match a signal to a parcel, I may queue that grain. I do not queue raw HTML or PDFs. The transform and the lists stay a separate track.
Problem
A grain here is a signal after it sits on a parcel. Raw HTML is the markup I fetched. A PDF is the file. Both are bytes I collected, not the origin. A source is that origin. One job per source.
Once land and match no longer fit in one process, the fork is what goes on a bus: those bytes, nothing, or only the grain after match.
Portal-shaped sources still need a browser. Government apps, vendor listing screens, JavaScript tables with no JSON. File-shaped sources do not. Rolls, tax files, documented downloads. Both have to sit on the property grain before they are list-eligible.
Some tools transform before they store. Others load raw first, then transform in the warehouse. That later step is my transform. Loaders such as Airbyte, dlt, or Singer taps are for documented APIs and files. They are not a crawler. They are not a PDF tagger.
I build property lists today. Shared landing: the same landed record could match another grain later. Not this post.
Options considered
- Stay in one process — collect, extract, and match write tables in one source job forever. Simple. The database becomes a file system. Match stays a search-index guess, with no table for the misses.
- Queue raw HTML and PDFs — fetched bytes on a bus, workers parse, a “ready” queue loads tables. I know that pattern. It scales crawl speed. I cap crawl speed on purpose. Unmatched events become traffic.
- Queue matched grains — land bytes (crawlers for portals, loaders for files) into object storage. Extract under the same signal contracts. Match as its own step. Send a small message only when a signal sits on a parcel. Loaders write source tables. The transform stays a separate track.
Decision
I went with option 3.
Option 1 holds until raw files, gold lists, and orders fight for the same disk, and until a fuzzy match cannot be its own table. Option 2 is the wrong first pattern: a pipeline of every attachment. I still schedule sources. The queue only separates “this source finished match” from “rows are ready to transform.”
Loaders do not replace portal collectors. A document parser does not replace the signal contracts. The search index is for candidates. Probabilistic match is for the link. The transform is for exact keys and for gold. Unmatched is stored. It is not queued. It is not a list row. The qualification gate stays the match, not the download.
Work before that grain message stays boring. “Unprocessed row, then an extract worker” is a status column in the database (pending, processing, done, claim the row so two workers do not take it). I already do that shape for a re-collect. A broker is an option when many enrichments need their own retries and counts. I do not need one because the data can be a week old. An outbox is only worth it if an API must start enrich and cannot lose the job. A nightly cron does not need one. Run ids, a count of failed rows, and a dead-letter status beat adding a streaming cluster.
How it works
Same filing. A crawler lands the PDF in object storage. A parser plus my tagger emit parties and identifiers. Match tries the parcel. If it hits, a short event goes on the queue: signal, parcel, confidence, a pointer to the file. Not the file. A loader writes the source row. The transform builds gold as it does today. If it misses, the file stays. The unmatched row stays. The queue does not see it. That event never puts the property on the absentee-and-filing list.
flowchart TB
portals[Portal sources] --> crawlers[Collectors]
files[File sources] --> elt[Loaders for files]
crawlers --> lake[Object storage]
elt --> lake
lake --> extract[Extract then match]
extract --> q[Matched grain queue]
extract --> out[Unmatched stays out]
q --> load[Load source tables]
load --> compose[Transform]
Hard parts
- A bus of raw HTML looks like progress because messages move. It burns sources and skips the gate.
- A loader pointed at a portal is still a crawler with extra parts. File loads and portal lands have to share a raw prefix, or the transform cares how a file arrived.
- The message has to hold a pointer to the bytes, not the PDF. Otherwise the queue is a second file store.
- Loaders have to be safe to run twice. Bad messages go aside. Otherwise one bad match blocks gold.
- It is easy to name the match step and then move the gate into collection. List eligibility stays the match.
- A streaming cluster is more ops than this needs when week-old data is fine. I would reach for a broker when parallelism and per-job failures hurt. Not because a brochure said so.
What we’d change
I would land bytes in object storage before I introduced a queue. I would turn extract-status polling into an explicit work queue in the database sooner, on the same path as a re-collect. I would not stream every PDF through workers just to use the machines. Splitting the warehouse, and changing how lists are served, wait until a night refresh or a slow query actually hurts. That is a separate track. It is not a dependency of the collectors.
References
- dbt — ETL vs ELT — the letters are load versus transform. My transform is the in-warehouse step.
- Airbyte — Connector Builder — JSON APIs, not a browser.
- Splink — probabilistic linkage. Unmatched is a table, not a failed join.
- Apache Iceberg — a lake as the raw layer, not as the list store.
- PostgreSQL — locking clause (
SKIP LOCKED) — claim a work row without a separate broker.