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

The actual problem: a sales report needs a definition of truth
Delivery in numbers
| Dimension | Documented system | Meaning |
|---|---|---|
| Logical data quality layers | 3 — bronze, silver, gold | Preserve source, standardize records, publish reporting model |
| Warehouse tables | 4 | One sales fact table and 3 dimensions |
| Dimension tables | 3 — customer, product, date | Consistent grouping axes for sales analysis |
| Analytics fact tables | 1 — sales | Centralized quantitative event measures |
| Batch orchestration cadence | Hourly | Reporting latency includes scheduled batch wait |
| Continuous ingestion | Kafka + Spark Structured Streaming | Input acceptance is decoupled from batch reporting |
| Columnar storage | Parquet, partitioned by date | Replayable intermediate event history |
End-to-end data path
| Boundary | What is written or processed | Contract and recovery concern |
|---|---|---|
| Producer → Kafka | Order event and associated fields | Schema stability, unique event identity |
| Kafka → Spark | Stream consumption and parsing | Offset progress, invalid payload handling |
| Spark → bronze Parquet | Append source records, partition by date | Restart from checkpoint; retain raw replay data |
| Bronze → silver | Validate, clean, standardize, deduplicate | Decide winner for duplicate or conflicting events |
| Silver → gold | Resolve dimension keys; write sales fact rows | Idempotent loads, stable grain, referential integrity |
| Airflow → batch stages | Execute scheduled dependencies | Retry semantics, reruns, observability |
| PostgreSQL → reports | Query facts joined to dimensions | Business definitions and freshness expectations |
Stage 1: capture events without pretending they are perfect
Stage 2: turn raw events into a defined analytical contract
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
| Table | Grain / intended role | Example analytical question |
|---|---|---|
| Sales fact | Each modeled sales record at a declared business grain | What revenue or quantity was recorded? |
| Customer dimension | Customer identity and reporting attributes | Which customer segments purchased? |
| Product dimension | Product identity and descriptive fields | Which products contributed to sales? |
| Date dimension | Calendar date attributes | How do totals vary by day/month? |
Two clocks: continuous ingestion and hourly reporting
The failure cases worth designing for
| Failure scenario | Why it matters | Recovery principle |
|---|---|---|
| Spark restarts mid-stream | Some events may be processed again | Checkpoints plus idempotent downstream interpretation |
| Source publishes duplicate order events | Sales might be double-counted | Stable event/order identity and deterministic deduplication |
| Airflow load fails halfway | Fact/dimension state may disagree | Transactional or repeatable warehouse loads |
| A late event belongs to yesterday | Date partitions and totals change | Event-time conventions and backfill policy |
| Product attributes change | Historical reports can shift unexpectedly | Explicit dimension history policy |
| Schema evolves without notice | Parsing and warehouse writes can break | Versioned schemas and compatibility checks |
Why this architecture instead of a single script?
What I delivered, and what I would benchmark
Repository and deep dive




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