Skip to content
Leivadev
Back to Blog

A resumable CSV inventory pipeline with medallion architecture on Cloudflare

How I moved tens of thousands of retail inventory rows from a CSV into a Postgres catalog reliably, using a Bronze/Silver/Gold pipeline built from Cloudflare Workflows, Queues, Durable Objects, R2, and D1.

· 7 min read · Data EngineeringCloudflare WorkersArchitecture

A retail client needed to sync their inventory into our product catalog from a CSV export: tens of thousands of rows, uploaded by a human, on a schedule that was really “whenever someone remembers.” The naive version of this (“parse the file and insert the rows”) fails in every way you would expect. A worker times out on a big file. One malformed row kills the whole batch. A re-upload of yesterday’s file double-writes everything. And when it half-fails, nobody can tell you which rows made it and which did not.

I built it instead as a medallion pipeline (Bronze, Silver, Gold) on Cloudflare. Medallion architecture is usually talked about in the data-lake world, but the discipline (separate the raw landing from the cleaning from the business shaping, and let each layer fail on its own terms) is exactly what an unreliable, human-driven import needs. Here is how it maps onto Cloudflare primitives.

The shape of the pipeline

Each layer has one job and hands its output to the next through R2, and a Cloudflare Workflow orchestrates the whole thing as durable, individually-retryable steps.

flowchart TB
  U[CSV upload] --> B
  subgraph Bronze
    B[raw CSV in R2 + run record in D1]
  end
  subgraph Silver
    S[validate + normalize -> valid rows in R2]
  end
  subgraph Gold
    G[business rules + enrich -> batches in R2]
  end
  B --> S --> G --> Q[upsert queue]
  Q --> DO[Durable Object - batch coordinator]
  DO --> P[(Postgres catalog - bulk upsert)]
  S -.rejected rows.-> DLQ[R2 JSONL + D1 error table]
  G -.rejected rows.-> DLQ

The orchestration lives in a Cloudflare Workflow. Workflows let you break a long job into named steps (step.do("silver-validate", ...)) where each step’s result is persisted and each step retries independently. If the Gold step fails, the completed Silver step does not re-run. That durability is the backbone of everything below.

Bronze: land the raw file, decide if you should even start

The upload endpoint does as little as possible. It streams the CSV straight into R2, computes a hash of the file, writes an import_run record to D1 with a pending status, and dispatches the Workflow with the run id, the R2 key, and the file hash. Then it returns 202 Accepted. The heavy work happens behind the workflow, not on the request.

One guard matters here: before accepting an upload, it checks whether a run is already in progress and refuses if so. A single-threaded catalog upsert does not want two imports racing each other, and stopping that at the door is far simpler than reconciling it later.

Bronze is deliberately dumb. It does not look at the contents. Its only job is to get the raw bytes somewhere durable and start the machine.

Silver: make the data trustworthy

The Silver step reads the raw CSV back out of R2 and turns it into clean, canonical rows. This is schema and shape work, not business logic:

  • Parse with a real CSV parser (I moved to csv-parse after a sniffer-based approach mishandled edge cases), tolerating BOMs, loose quoting, and ragged columns.
  • Validate the header against the required columns, and reject the whole file early if the delimiter or header is wrong (a single-column “header” almost always means a broken export).
  • Validate and normalize each row: a fixed-length alphanumeric product code, a numeric item number of bounded length (after stripping commas and spaces), length-capped names and brands, dates parsed into ISO, and prices parsed from a formatted currency string into integer minor units so money never rides on a float.
  • Normalize category and product names into URL-safe handles with a Spanish-aware slugifier (so “JAMÓN & QUESO” becomes jamon-y-queso), and drop duplicate item numbers within the same file.

A row that fails any of this is not fatal. It gets pushed to a rejected list with its row index and a human-readable reason. The valid rows are written to R2 as JSON for the next step, and the rejects are written twice: as a JSONL file in R2 and as rows in a D1 import_row_error table, so a human can later ask “why did line 4,213 not import?” and get a real answer.

Gold: apply business rules and shape for the write

Silver guarantees the data is well-formed. Gold decides whether it is acceptable and shapes it for the catalog:

  • Business-rule validation that has nothing to do with types: the price must fall inside an allowed range, and a last-purchase date cannot be in the future. These are policy, so they live here, not in Silver.
  • Enrichment: each row is matched against a table of fiscal category codes (loaded once from D1) and tagged with processing metadata (a timestamp, the run id, the file hash).
  • Splitting: the enriched rows are chunked into batches sized for the write phase, and the batches are written to R2.

Gold rejects, like Silver’s, are recorded to the D1 error table rather than aborting the run. By the time Gold finishes, R2 holds a set of batch files that are known-good and ready to persist, and the run record knows exactly how many rows were accepted and how many were rejected at each layer.

The write phase: a queue, a coordinator, and a bulk upsert

The Workflow enqueues each batch onto a Cloudflare Queue and then hands off. A queue consumer picks up batches and performs the actual write into the Postgres catalog, and a Durable Object acts as the coordinator that tracks how many batches have completed for a run. The Workflow’s final step asks the coordinator to wait until every batch is done, then computes a final status (success, partial, or failed) from the processed and rejected counts.

The write itself is a bulk, multi-phase upsert, not row-by-row inserts. Within a batch it upserts brands, then categories, then products, then their custom properties, options, variants, and prices, each as a set operation. New catalog entities get prefixed ULID identifiers. Doing it in phases respects the foreign-key ordering while still moving each phase in bulk, which is the difference between a batch that lands in a fraction of a second and one that makes hundreds of round trips.

The two details that made it actually work

Two things are worth calling out because they are not obvious until production teaches them to you.

R2 is the hand-off medium between steps, on purpose. A Cloudflare Workflow persists each step’s return value in its own storage, which has size limits. Passing tens of thousands of parsed rows from Silver to Gold as a step return value blows past those limits (you get a SQLITE_TOOBIG-class failure). So each layer writes its output to R2 and passes only the key forward. The steps stay small and durable; the data rides in object storage where size is a non-issue.

Resumability keys off the file hash. Because the raw file and its hash are recorded up front, a re-run can look up a prior run of the same file, copy the batches that already completed, and enqueue only the ones that are still missing. A pipeline that died three-quarters of the way through a fifty-thousand-row file does not start over; it finishes the last quarter. Combined with the per-row error table, that turns “the import broke” from an incident into a resumable, inspectable state.

What medallion bought me

I could have written a single big function that parsed, validated, and inserted. It would have been shorter and it would have been miserable to operate. Splitting the work into Bronze, Silver, and Gold gave me three things that mattered more than brevity:

  • Localized failure. A parsing problem is a Silver problem, a policy violation is a Gold problem, and a database hiccup is a write-phase problem. Each fails in one place with one kind of error, and the durable workflow step retries just that part.
  • Auditability. Every rejected row is recorded with a reason and a line number, at the layer that rejected it. “Which rows didn’t import and why” is a query, not an investigation.
  • Resumability and safety. Raw-first landing plus content hashing made re-runs cheap and idempotent, and the one-run-at-a-time guard kept two imports from fighting over the catalog.

Medallion architecture is usually sold as a data-warehouse pattern, but it is really just a discipline for untrusted data moving through stages. A messy CSV uploaded by a human is exactly that, and the Cloudflare primitives (Workflows for durable orchestration, R2 for inter-stage storage, Queues and a Durable Object for the write phase, D1 for run and error tracking) map onto the pattern cleanly enough that the architecture mostly explains itself.