05 / PROJECTStreaming Data Engineering
Real-Time Retail Data Pipeline
Streaming ingestion with a replayable staging layer
Retail transactions streamed through Kafka into Spark Structured Streaming, staged on S3 so batches can be replayed, then loaded into Redshift under Airflow orchestration.
Architecture
Ingestion
Kafka
Durable event log for retail transactions, with retention that permits replay.
Retail events → Kafka → Spark Structured Streaming → Amazon S3 (staging) → Redshift → Airflow orchestration
The problem
Streaming pipelines that write straight to a warehouse have no recovery story. When a transformation is wrong, the source events are already gone.
Engineering challenge
Handling continuous transaction volume while keeping a path back to the source data when a transformation turns out to be wrong.
Approach
Kafka retains the event log, Spark Structured Streaming processes it in micro-batches, and every processed batch lands on S3 before Redshift. That staging layer makes reprocessing possible without re-ingesting. Airflow orchestrates the downstream steps with task-level retries.
- Apache Kafka
- Durable, replayable event log rather than a fire-and-forget queue.
- Spark Structured Streaming
- Micro-batch semantics with checkpointing and fault tolerance.
- Amazon S3
- Cheap, immutable staging that makes reprocessing a normal operation.
- Apache Airflow
- DAG orchestration with retries and failure visibility at task granularity.
Engineering detail
- 01Kafka → Spark Structured Streaming → S3 → Redshift
- 02Micro-batch processing through Spark Structured Streaming
- 03Replayable S3 staging layer — batches can be reprocessed after a logic fix
- 04Airflow orchestration across the load and transform steps
- 05Task-level retries, so one transient failure does not fail the run
Engineering practices
- Staging before warehouse load, so batches remain replayable.
- Stream checkpointing for fault tolerance.
- Retries scoped to the task, not the whole DAG.
Features
- Kafka ingestion of transaction events.
- Micro-batch Spark transformation.
- Replayable S3 staging.
- Airflow-orchestrated Redshift loads.