June 20, 2026

Realtime data

Designed a reproducible data pipeline with Kafka ingestion, checkpointed Spark streaming, Parquet bronze storage, silver validation, and hourly Airflow loads into a four-table analytics model.

Data engineering
Role
Data engineer — streaming and analytics systems
Published
June 2026
Focus
Data engineering
Engineer
Saaim Abdullah
E-commerce data platform — from live events to trusted sales reporting system overview

The actual problem: a sales report needs a definition of truth

An e-commerce application can produce orders all day and still fail to answer a basic question: “What was net sales yesterday?” Events arrive more than once, consumers restart, products change names, and an event accepted at 11:59 might not reach the warehouse until after midnight. Analytics fails when those conditions are treated as exceptions rather than a normal property of data systems. I built this pipeline to connect continuous event capture with repeatable analytical transformations. Kafka receives order events, Spark Structured Streaming persists them into date-partitioned Parquet, and Airflow schedules downstream cleaning and warehouse loads. PostgreSQL hosts a sales star schema that reporting tools can query without reconstructing raw events for every chart. The architectural goal was not to list popular technologies. It was to establish where data can be replayed, where errors are rejected, when reporting becomes fresh, and which keys make totals consistent.

Delivery in numbers

DimensionDocumented systemMeaning
Logical data quality layers3 — bronze, silver, goldPreserve source, standardize records, publish reporting model
Warehouse tables4One sales fact table and 3 dimensions
Dimension tables3 — customer, product, dateConsistent grouping axes for sales analysis
Analytics fact tables1 — salesCentralized quantitative event measures
Batch orchestration cadenceHourlyReporting latency includes scheduled batch wait
Continuous ingestionKafka + Spark Structured StreamingInput acceptance is decoupled from batch reporting
Columnar storageParquet, partitioned by dateReplayable intermediate event history
The four-table model and hourly schedule are implementation facts. No claimed event throughput, p99 latency, or cloud cost is being inferred from them.

End-to-end data path

BoundaryWhat is written or processedContract and recovery concern
Producer → KafkaOrder event and associated fieldsSchema stability, unique event identity
Kafka → SparkStream consumption and parsingOffset progress, invalid payload handling
Spark → bronze ParquetAppend source records, partition by dateRestart from checkpoint; retain raw replay data
Bronze → silverValidate, clean, standardize, deduplicateDecide winner for duplicate or conflicting events
Silver → goldResolve dimension keys; write sales fact rowsIdempotent loads, stable grain, referential integrity
Airflow → batch stagesExecute scheduled dependenciesRetry semantics, reruns, observability
PostgreSQL → reportsQuery facts joined to dimensionsBusiness definitions and freshness expectations

Stage 1: capture events without pretending they are perfect

Kafka sits between producers and the streaming processor. That separation allows event capture and downstream computation to progress at different rates. Spark Structured Streaming uses checkpoints to track its stream progress and writes the ingested data to Parquet. The bronze layer preserves source records rather than overwriting them with the latest business interpretation. This matters during change. If the silver deduplication key turns out to be wrong, retained bronze data provides a route to rebuild downstream tables. A stream checkpoint protects processing progress; it does not, on its own, guarantee unique business events or exactly-once inserts into every downstream system.

Stage 2: turn raw events into a defined analytical contract

Silver applies the validation and deduplication rules between raw input and reporting. The crucial design question is the deduplication grain: two records may represent a retransmission of one order event, two legitimate updates to an order, or two separate purchases. Those cases require domain-aware keys and ordering rules. A blanket dropDuplicates without a clear event identity would be an implementation detail, not a data quality strategy. For a production revision, I would publish a schema and expectations for required IDs, event timestamps, monetary precision, and allowed status changes. Invalid records should enter a quarantine or error table rather than silently disappearing from revenue calculations. The current case study claims the implemented validation/deduplication stage, not a measured defect-elimination rate.

Stage 3: build the warehouse for questions, not storage convenience

The gold layer is a four-table PostgreSQL star schema. The fact table represents the sales measures; customer, product, and date dimensions provide reusable reporting attributes. Its value is consistency: revenue grouped by product in one dashboard should share definitions and keys with revenue grouped by customer in another.
TableGrain / intended roleExample analytical question
Sales factEach modeled sales record at a declared business grainWhat revenue or quantity was recorded?
Customer dimensionCustomer identity and reporting attributesWhich customer segments purchased?
Product dimensionProduct identity and descriptive fieldsWhich products contributed to sales?
Date dimensionCalendar date attributesHow do totals vary by day/month?
Star-schema discipline includes stable surrogate keys, clear measures, careful decimal types for currency, and documented handling of corrections. These are the decisions I would review before expanding the dashboard surface.

Two clocks: continuous ingestion and hourly reporting

A useful mental model is that this system has two clocks. Kafka/Spark ingest continuously; Airflow runs the reporting transformations hourly. Therefore, “event written” and “event visible to analytics” are distinct statuses. A report can be consistent and still be behind the live storefront. An illustrative schedule, not a measured service-level agreement: if a clean event lands one minute after an hourly cutoff, it may wait nearly 59 minutes before the next batch begins, plus transform and load time. That arithmetic explains why an hourly pipeline should not be marketed as sub-second end-to-end analytics.

The failure cases worth designing for

Failure scenarioWhy it mattersRecovery principle
Spark restarts mid-streamSome events may be processed againCheckpoints plus idempotent downstream interpretation
Source publishes duplicate order eventsSales might be double-countedStable event/order identity and deterministic deduplication
Airflow load fails halfwayFact/dimension state may disagreeTransactional or repeatable warehouse loads
A late event belongs to yesterdayDate partitions and totals changeEvent-time conventions and backfill policy
Product attributes changeHistorical reports can shift unexpectedlyExplicit dimension history policy
Schema evolves without noticeParsing and warehouse writes can breakVersioned schemas and compatibility checks

Why this architecture instead of a single script?

A direct producer-to-PostgreSQL insert would be simpler for a tiny workload, but would entangle ingestion failure with reporting schema changes and lose a convenient replay source. Kafka + bronze storage creates a durable separation. Airflow provides explicit scheduled dependencies instead of hidden cron chains. Spark handles structured event processing, while PostgreSQL remains a familiar query surface for the gold model. Those choices have costs: more local services, checkpoints to manage, batch/stream coordination, and operational overhead. They are justified by recoverability and a clear transformation contract—not by assuming that every small e-commerce workload needs Kafka.

What I delivered, and what I would benchmark

Delivered: producer-to-Kafka event path, Spark/Parquet ingestion, three logical data layers, hourly Airflow-orchestrated batch transformations, validation/deduplication, and a PostgreSQL warehouse with one fact and three dimensions. The stack is documented for containerized execution and S3-compatible storage; this is not presented as an unverified managed AWS/EMR production deployment. Next validation plan: replay a fixed, published dataset twice and confirm equal gold outputs; inject duplicates and late arrivals; interrupt Spark and Airflow mid-run; compare source/silver/gold counts; measure hourly freshness and job duration. Throughput and cost numbers belong in the final report only after environment, event size, partition counts, and hardware are recorded. My movie recommendation service applies related orchestration ideas to training artifacts and a serving API.

Repository and deep dive

E-commerce data platform — from live events to trusted sales reporting architecture diagram 1E-commerce data platform — from live events to trusted sales reporting architecture diagram 2

More to explore

Let’s talk

I like working through complex problems with people who care about the details. Have a product to build, an engineering role, or an interesting challenge? Let’s start a conversation.

A little note

SaaimOpen to full-time roles, contract work, and conversations about things worth building.

ϟ 1
Contact