An end-to-end data platform for ingesting, processing, governing, and serving e-commerce data in real time. The platform combines application events with PostgreSQL change data capture and publishes trusted analytics products for finance, marketing, fulfillment, and inventory teams.
PostgreSQL changes are captured by Debezium and delivered through Kafka alongside application events. Apache Flink writes the immutable Bronze layer with exactly-once checkpointing and event-time semantics. Apache Spark performs CDC merges, deduplication, SCD2 processing, and Iceberg maintenance. dbt and Trino build and serve the Gold analytics layer.
- PostgreSQL CDC using Debezium and logical replication
- Eight versioned application-event streams governed by Avro contracts
- Exactly-once Flink ingestion with watermarks, keyed deduplication, late-data handling, and DLQs
- Apache Iceberg Bronze, Silver, and Gold layers backed by MinIO locally or Amazon S3 in AWS
- Idempotent Spark CDC merges, SCD2 history, backfills, and table maintenance
- Finance, marketing, operations, and inventory data products built with dbt on Trino
- One pipeline specification deployable through either Dagster or Airflow
- Layered quality gates using schema validation, dbt tests, and Great Expectations
- OpenLineage metadata published to Marquez
- Prometheus alerts and Grafana dashboards for freshness, lag, quality, and checkpoint health
- Terraform modules for the AWS development environment
flowchart LR
apps[Application events] --> registry[Schema Registry]
registry --> kafka[(Kafka)]
postgres[(PostgreSQL)] --> debezium[Debezium CDC]
debezium --> kafka
kafka --> flink[Apache Flink]
flink --> bronze[(Iceberg Bronze)]
bronze --> spark[Apache Spark]
spark --> silver[(Iceberg Silver)]
silver --> dbt[dbt on Trino]
dbt --> gold[(Iceberg Gold)]
gold --> trino[Trino]
orchestrator[Dagster or Airflow] -. orchestrates .-> spark
orchestrator -. orchestrates .-> dbt
flink -. metrics .-> grafana[Prometheus and Grafana]
orchestrator -. lineage .-> marquez[OpenLineage and Marquez]
More detail is available in docs/architecture.md.
- Docker Engine with Docker Compose v2
- Python 3.10 or later
- GNU Make
- At least 16 GB RAM for the core services
- Approximately 24 GB RAM for the complete platform
Create the local configuration and install the Python tooling:
cp .env.example .env
make install
make testStart the core infrastructure, initialize PostgreSQL, create Kafka topics, and register the CDC connector:
make up PROFILE=core
make migrate
make seed
make topics
make connectorsStart stream processing, batch processing, observability, and Dagster:
make up PROFILE=core,streaming,batch,observability,orchestrator-dagster
make run-flinkGenerate a deterministic workload and materialize the Silver and Gold layers:
make generate SCENARIO=baseline RATE=20 SEED=42
make run-pipeline TARGET=allOpen a Trino session:
make queryExample query:
SELECT
revenue_date,
order_count,
gross_revenue,
refunds,
net_revenue
FROM iceberg.gold.fct_revenue_daily
ORDER BY revenue_date DESC
LIMIT 30;Stop all platform services with make down.
| Profile | Services |
|---|---|
core |
PostgreSQL, Kafka, Schema Registry, Debezium, MinIO, Iceberg REST, Trino |
streaming |
Flink JobManager and TaskManager |
batch |
Spark with Iceberg runtime |
observability |
Prometheus, Grafana, Marquez, and Marquez Web |
orchestrator-dagster |
Dagster orchestration service |
orchestrator-airflow |
Airflow orchestration service |
Profiles can be combined as a comma-separated PROFILE value.
| Service | Address |
|---|---|
| Trino | http://localhost:8080 |
| Schema Registry | http://localhost:8081 |
| Flink | http://localhost:8082 |
| Kafka Connect | http://localhost:8083 |
| Airflow | http://localhost:8084 |
| MinIO Console | http://localhost:9001 |
| Grafana | http://localhost:3000 |
| Marquez | http://localhost:3001 |
| Dagster | http://localhost:3002 |
Local credentials and connection settings are defined in .env. Development defaults are
provided in .env.example; replace them before using a shared environment.
| Domain | Gold datasets | Primary consumers |
|---|---|---|
| Finance | fct_revenue_daily, fct_payments |
Revenue and payment reporting |
| Marketing | dim_customer, fct_conversion_funnel |
Conversion and customer analysis |
| Operations | fct_fulfillment, fct_shipment_sla |
Fulfillment and carrier performance |
| Inventory | fct_stock_position_daily, stockout_risk |
Stock health and reorder planning |
Every data product includes ownership metadata, descriptions, tests, and downstream exposures in the dbt project.
The generator provides repeatable scenarios for development, validation, and incident exercises:
make generate SCENARIO=baseline RATE=100 SEED=42
make generate SCENARIO=flash_sale RATE=500 SEED=42
make generate SCENARIO=fraud_burst RATE=200 SEED=42baseline represents normal traffic, flash_sale increases ordering pressure, and fraud_burst
raises payment failures and late-event frequency.
make contract-check # validate event contracts
make demo-dlq # inject an invalid event
make demo-timetravel # exercise Iceberg snapshot queries
make terraform-validate # validate the AWS environment
make status # inspect local service healthOperational procedures are documented in docs/runbooks, including DLQ replay, backfills, late data, schema changes, freshness incidents, Iceberg maintenance, and erasure requests.
contracts/ Avro event contracts and compatibility checks
db/ PostgreSQL migrations and deterministic seed data
generators/ Application-event and OLTP workload generation
ingestion/ Kafka topics and Debezium connector definitions
streaming/ Flink SQL jobs and stateful streaming primitives
lakehouse/ Iceberg DDL, Spark Silver jobs, SCD2, and maintenance
transform/ dbt staging, intermediate, and Gold data products
quality/ Great Expectations and Soda quality definitions
orchestration/ Shared pipeline graph plus Dagster and Airflow adapters
lineage/ OpenLineage transport and event emission
observability/ Prometheus metrics, alerts, and Grafana dashboards
infra/ Docker configuration and AWS Terraform modules
tests/ Unit, integration, and end-to-end validation
docs/ Architecture decisions, data contracts, and runbooks
make test
make lint
make contract-checkUnit tests run without infrastructure. Integration and end-to-end tests require the corresponding
Compose profiles and are enabled explicitly with RUN_E2E=1.
The CI pipeline validates Python formatting and tests, contract compatibility, Compose rendering, and Terraform configuration on every change.