Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Building Blocks of Data Systems

Every data system design is assembled from roughly ten building blocks. For each one, know what it’s for, how it works inside, the knobs, and when NOT to use it.

Contents: Message log (Kafka) · Partitioning · Stream processors · Batch engines · Storage & table formats · Serving stores · Caches · Orchestrators · Schema registry · Probabilistic structures


1. The message log (Kafka)

Mental model

A Kafka topic is an append-only, partitioned, replicated log. Producers append; consumers read at their own offset. Data isn’t deleted when it’s read. It’s deleted by retention (time/size) or compaction.

flowchart LR
    P1[Producer A] -->|key=user_1| T0
    P2[Producer B] -->|key=user_2| T1
    P2 -->|key=user_3| T2
    subgraph Topic["Topic: clicks (3 partitions, RF=3)"]
        T0["P0: 0 1 2 3 4 5 →"]
        T1["P1: 0 1 2 3 →"]
        T2["P2: 0 1 2 3 4 →"]
    end
    subgraph CG1["Consumer group: flink-agg"]
        C1[consumer 1]
        C2[consumer 2]
    end
    subgraph CG2["Consumer group: lake-sink"]
        C3[consumer 1]
    end
    T0 --> C1
    T1 --> C2
    T2 --> C2
    T0 --> C3
    T1 --> C3
    T2 --> C3

Key facts interviewers probe

ConceptWhat to say
OrderingGuaranteed only within a partition. Same key → same partition (hash(key) % N) → per-key order.
ParallelismMax active consumers in a group = number of partitions. Partitions are the unit of scale.
Consumer groupsEach group gets every message once (load-balanced across its members); different groups are independent → fan-out for free.
Replicationreplication.factor=3, min.insync.replicas=2, producer acks=all → no data loss if one broker dies.
RetentionTime/size based (e.g. 7 days) → replay window for reprocessing/backfills. Tiered storage moves old segments to S3 for long retention.
Log compactionKeeps only the latest value per key; tombstone (key, null) deletes. Good for CDC/changelog topics and state restore.
Idempotent producerenable.idempotence=true → broker dedupes retries via (producer id, sequence number).
TransactionsAtomic writes across partitions + offset commits → exactly-once for read-process-write within Kafka.
Adding partitionsChanges hash(key) % N → breaks key ordering for existing keys. Over-provision up front.
Consumer laglog-end-offset − committed-offset: the health metric for streaming.

Kafka vs alternatives

KafkaKinesisPulsarPub/SubRabbitMQ/SQS
ModelLogLog (shards)Log + queueManaged log-ishQueue (delete on ack)
Replay✅ retention✅ up to 365d✅ tiered✅ seek (7d)❌
Orderingper partitionper shardper keyper ordering keylimited
Opsheavy (or MSK/Confluent)zeroheavyzerolight
Use whendefault for data pipelinesAWS-native, simplermulti-tenant, geoGCP-nativetask queues, not analytics

When NOT Kafka: low volume (< a few hundred events/s) where a DB table + CDC or SQS is simpler; pure batch file drops (land in object storage + AutoLoader).


2. Partitioning strategies

Partitioning shows up in Kafka topics, Spark shuffles, table layouts and serving stores. Same trade-off everywhere: locality vs balance.

StrategyExampleProsCons
Hash by keyhash(user_id) % 64Even spread; co-locates a key’s data (ordering, joins, state)Hot keys → hot partition; resizing reshuffles
Rangedates, A–F, G–MRange scans, pruningHotspots on the “latest” range (time-ordered writes)
Round-robinno keyPerfect balanceNo ordering; aggregations need a shuffle later
Time-based (tables)PARTITION BY event_datePruning, easy retention/backfill per dayToo fine-grained → small files
Composite / salted(campaign_id, salt 0–9)Spreads hot keysTwo-phase aggregation needed
Geo / tenantregion, tenant_idIsolation, data residencyUneven tenant sizes

Hot-key mitigation (classic deep dive)

flowchart LR
    E[Events<br/>key = campaign_id] --> S{Hot key?}
    S -->|no| H1[hash campaign_id]
    S -->|yes| SALT[append salt 0..N-1<br/>campaign_id#7]
    H1 --> A1[Partial aggregate per key]
    SALT --> A1
    A1 --> M[Second stage:<br/>strip salt, merge partials]
    M --> OUT[(Final per campaign)]
  1. Detect: per-partition lag/throughput skew; top-K key metrics.
  2. Salt the hot keys (or all keys) with a random suffix → partial aggregates → merge in a second stage. Works because SUM/COUNT/MIN/MAX/HLL are mergeable.
  3. Local pre-aggregation (combiner) before shuffle: Spark does this automatically for reduceByKey/groupBy().agg() partial aggregation.
  4. Split hot keys onto dedicated partitions/consumers.
  5. In joins: broadcast the small side, or AQE skew join (Spark splits skewed partitions automatically).

3. Stream processors

Spark Structured StreamingApache FlinkKafka StreamsManaged (DLT, Dataflow)
ModelMicro-batch (default), continuous experimentalTrue streaming, event-at-a-timeLibrary inside your appSpark / Beam under the hood
Latency~100 ms – seconds (Real-Time Mode lowers this)msmsseconds
StateRocksDB state store, applyInPandasWithState / transformWithStateBest-in-class: keyed state, timers, savepointsRocksDB, changelog topicsmanaged
Exactly-oncecheckpoint + idempotent/transactional sinkscheckpoint barriers (Chandy-Lamport) + 2PC sinksKafka transactions✅
Batch + stream same code✅ (same DataFrame API)✅ (Table API)❌✅
Choose whenLakehouse/Delta shop, seconds latency OK, team knows SparkSub-second, complex event processing, huge state, timersKafka-in Kafka-out microservicesWant zero ops

Talking point: “Our freshness SLA is under a minute and we’re a Databricks shop, so Structured Streaming with Delta sinks is the pragmatic choice. If we needed sub-100 ms fraud scoring with per-user timers, I’d pick Flink.”

Core streaming concepts (detail in streaming deep dive)

Streaming concepts


4. Batch engines

EngineSweet spot
SparkLarge-scale ETL, joins over TBs, ML feature pipelines, lakehouse
dbt (on a warehouse/Spark)SQL transformations with tests, docs, lineage, incremental models
Warehouse SQL (Snowflake, BigQuery, Databricks SQL, Redshift)ELT, analytics, BI
Trino/PrestoFederated interactive SQL over lakes
DuckDB / PolarsSingle-node, up to ~100s GB, insanely fast; great for tests and small jobs

Key batch design rules:

  1. Idempotent partitions: each run overwrites exactly one logical partition (INSERT OVERWRITE ... PARTITION (dt='2026-10-01') or MERGE / replaceWhere).
  2. Incremental by default: process only new data (watermark column, Delta CDF, AutoLoader), with a full-refresh escape hatch.
  3. Deterministic: no now() inside transformations. Pass the logical date in as a parameter so backfills reproduce history.

5. Storage layer

flowchart TB
    subgraph Catalog["Catalog (Unity / Glue / Polaris / Hive Metastore)"]
        CT[table name → metadata location, permissions]
    end
    subgraph Table["Open table format (Delta / Iceberg / Hudi)"]
        LOG[Transaction log / metadata tree<br/>versions, schema, file list, stats]
    end
    subgraph Files["Object storage (S3 / ADLS / GCS)"]
        F1[part-0001.parquet]
        F2[part-0002.parquet]
        F3[part-0003.parquet]
    end
    CT --> LOG
    LOG --> F1
    LOG --> F2
    LOG --> F3

File-size rule: target 128 MB–1 GB files. Thousands of tiny files are the most common lakehouse performance bug (metadata overhead, task overhead, S3 request costs) → compaction (OPTIMIZE), optimized writes, auto-compaction.


6. Serving stores

The access pattern picks the store:

NeedStoreWhy
Sub-second slice-and-dice on fresh events, high concurrency (user-facing analytics)Apache Pinot / Druid / ClickHouseColumnar + indexes + pre-aggregation (star-tree, rollups), real-time ingestion from Kafka
Internal BI, complex SQL, joins, moderate concurrencyWarehouse / lakehouse SQLFlexible, cheaper per TB, seconds latency
Point lookup by key in ms (features, profile, counters)Redis / DynamoDB / Cassandra / BigtableKV, predictable latency, horizontal scale
Search, text, logsElasticsearch / OpenSearchInverted index
Similarity searchVector DB (pgvector, Pinecone, Milvus, Mosaic AI Vector Search)ANN indexes (HNSW, IVF)
Graph traversalsNeo4j / NeptuneRelationships
Transactions (source of truth for an app)Postgres / MySQLOLTP, ACID
flowchart LR
    Q{Query pattern?}
    Q -->|"key → value, < 10 ms"| KV[Redis / DynamoDB]
    Q -->|"aggregations on fresh data,<br/>user-facing, 1000s QPS"| OLAP[Pinot / Druid / ClickHouse]
    Q -->|"ad-hoc SQL, joins,<br/>analysts"| WH[Warehouse / Lakehouse SQL]
    Q -->|"nearest neighbours"| VEC[Vector index]
    Q -->|"full text"| ES[Elasticsearch]

Senior nuance: you often need two serving stores: raw history in the lakehouse for flexibility, and pre-aggregated hot data in an OLAP/KV store for latency. Say explicitly which one is the source of truth (usually the lakehouse) and that serving stores can be rebuilt from it.


7. Caches


8. Orchestrators

AirflowDagsterDatabricks Workflows / Lakeflow JobsPrefect
ParadigmTask DAGs, schedule-centricAsset-centric (data-aware)Jobs + tasks, native to the lakehousePython flows
StrengthEcosystem, ubiquityLineage, partitions, testingZero infra, cluster mgmt, DLT/dbt tasksDynamic workflows
Watch outScheduler scale, DAG parsing, XCom misuseSmaller ecosystemVendor-specificSmaller ecosystem

Design rules: tasks idempotent, parameterised by logical date, retries with backoff, SLAs/alerts on freshness not just failure, data-aware triggers (run when upstream table updates) over cron chains. More in Orchestration & backfills.


9. Schema registry & contracts


10. Probabilistic data structures

StructureAnswersMemoryErrorMergeable?
HyperLogLogCount distinct~1.5–12 KB~0.8–2%✅ (union) → great for rollups
Count-Min SketchFrequency of item (heavy hitters)KBs–MBsOverestimates only✅
Bloom filter”Have I seen X?” (dedup, join pre-filter)~10 bits/item for 1% FPFalse positives, no false negatives✅ (OR)
t-digest / KLLPercentiles (p50/p95/p99)KBssmall rank error✅
Reservoir samplingUniform sample of a streamk itemsexact sample–

Use in design: “Real-time unique visitors per page per minute: I’ll store HLL sketches per minute; the hourly/daily numbers are unions of the minute sketches, so I never re-scan raw data. Finance’s exact numbers come from batch.”

SQL examples: approx_count_distinct(user_id), percentile_approx(latency, 0.95), Databricks hll_sketch_agg / hll_union_agg.


Cheat table: “If they say X, reach for Y”

They say…Reach for…
”Replay”, “multiple consumers”, “decouple”Kafka
”Changes from the database”Log-based CDC (Debezium) → Kafka → MERGE
”Sub-second dashboards for customers”Pinot/Druid/ClickHouse fed from Kafka
”Exactly once”Replayable source + checkpoint + idempotent sink; dedup keys
”Late events”Event-time windows + watermarks + restatement batch
”History of changes”SCD Type 2 / Delta CDF / snapshots
”Unique users at scale”HLL sketches
”Top-K trending”Count-min sketch + heap, or windowed aggregation + rank
”Point lookup features at inference”Online feature store (Redis/DynamoDB) synced from offline
”Search docs with LLM”RAG pipeline: chunk → embed → vector index + ACLs
”PII / GDPR”Classification tags, masking, crypto-shredding, deletion pipeline