
π Data Engineering & Analytics Β· Data Flow
How every insert, update and delete in an operational database is streamed to the data lake and search index in near real time using change data capture.
Drawing diagramβ¦
CDC data flow: Debezium reads the PostgreSQL write-ahead log and publishes row changes to Kafka topics, with schemas in a schema registry. A stream job writes changes to an Iceberg table in the data lake, another updates an Elasticsearch index, and a third invalidates Redis cache keys.
flowchart LR DB[(Orders PostgreSQL)] -->|Write-ahead log| DZ[Debezium Connector] DZ -->|Row changes| K[Kafka Topics] SR[(Schema Registry)] -->|Schemas| DZ K -->|Changes| LJ[Lake Writer - Flink] K -->|Changes| SJ[Search Indexer] K -->|Changes| CJ[Cache Invalidator] LJ -->|Upserts| ICE[(Iceberg Tables in S3)] SJ -->|Documents| ES[(Elasticsearch)] CJ -->|Delete keys| RC[(Redis Cache)] SR -->|Schemas| LJ
A typical modern data stack: data is loaded from apps and SaaS tools into a cloud warehouse, modelled with dbt, orchestrated with Airflow and served to BI dashboards.
A classic star schema for sales analytics: one fact table of order lines surrounded by date, customer, product, store and promotion dimensions.
The steps of a nightly batch pipeline with data quality gates: extract, validate, transform, load and publish, stopping safely when checks fail.
Where a real-time analytics stack runs: Kafka for events, Flink for stream processing, a real-time OLAP database and live dashboards.
The states of a single pipeline run in an orchestrator like Airflow, including retries, upstream failures and manual reruns.
The medallion pattern for a data lake: raw data lands in bronze, is cleaned in silver, and turned into business-ready tables in gold.