Data engineers build the systems that turn raw records from applications into tables that people, dashboards and models can trust. The tools have converged over the last few years: open table formats on object storage, a log such as Kafka in the middle, SQL engines on top, and orchestration and quality checks around the edges. AI raised the stakes rather than replacing the job. A model is only as good as its training data, and a retrieval system answers only from what the pipeline indexed, with the permissions the pipeline carried along.

This roadmap is the companion to the backend engineer roadmap, which stops where data leaves the service, and a track of the 2026 developer roadmap. It runs through six stages, each with a skill, a check that proves it and reading on this site, then works through one end-to-end design: moving orders from Postgres into a lakehouse in near real time and feeding both a dashboard and an AI assistant.

Advertisement

The stage map

Learn the stages in order. Streaming without solid modelling produces fast wrong numbers; a lakehouse without quality gates is a data swamp with transactions. The diagram shows where each stage lives on a typical platform.

A modern data platform: batch and streaming land in one set of open tablesOLTP databasesPostgres, MySQLEvents / appsclicks, logs, IoTSaaS / filesAPIs, CSV, docsCDClogical decodingKafka topicsschema registryBatch ingestorchestrated loadsLakehouse tablesIceberg / Delta / HudiStream processorFlink, SparkBI / SQLTrino, Spark SQLML featurestraining setsAI dataembeddings, evalsCatalog, lineage, quality checks, access controlThe catalog and the quality gates span every arrow; they are what make the tables trustworthy.
Sources flow through CDC, Kafka or batch ingest into open lakehouse tables, then out to BI, ML and AI consumers. Catalog, lineage and quality span the whole path.
StageSkillYou have it when
1. SQL and modellingJoins, windows, keys, dimensional modelsYou can model a business process and defend its grain
2. Batch processingSpark, columnar files, partitioning, shufflesYou can read a query plan and fix a skewed join
3. Lakehouse tablesIceberg / Delta / Hudi, snapshots, compactionYou can explain what a commit writes and how readers stay consistent
4. CDC and streamingKafka, Flink or Spark streaming, watermarks, exactly-onceYou can reason about late, duplicate and out-of-order events
5. Orchestration and qualityDAGs, idempotent loads, tests, lineageA bad load is blocked before anyone reads it
6. Data for AIEmbedding pipelines, ACLs, evaluation dataYou can re-index incrementally and prove deletions propagate

Stage 1: SQL and data modelling

SQL is the language of the job, and fluency means more than SELECT. Learn window functions, because deduplication, sessionisation and slowly changing dimensions are all window problems. Learn how joins multiply rows when keys are not unique; most wrong dashboard numbers are fan-out bugs.

Modelling is deciding the grain: what one row means. A fact table of order lines at one row per line item, with dimensions for customer, product and date, answers questions that a table of whole orders cannot. Slowly changing dimensions record history so that last year's revenue is attributed to last year's region. Know the source side too: transactions, MVCC and why reading a busy OLTP database directly for analytics is a bad idea.

Read next on this site: window functions, MVCC, transactional databases, materialized views, hash joins.

Advertisement

Stage 2: batch processing and columnar storage

Analytics data lives in columnar files, usually Parquet, where each column is stored and compressed separately and carries min/max statistics per row group. Engines read only the columns a query needs and skip row groups whose statistics rule them out. Partitioning by a coarse column such as date skips whole directories. Both only work if files are a sensible size; thousands of tiny files make planning slow and scans inefficient.

Spark is the workhorse for large transformations. Learn its execution model: a job becomes stages separated by shuffles, and the shuffle, where data is redistributed across the network by key, is where most cost and most failures live. Skew, one key with far more rows than the rest, turns one task into the whole job's runtime. Read the physical plan before and after every change, and know when a broadcast join removes a shuffle entirely.

Read next on this site: Parquet, the small-files problem, Spark stages and tasks, the Spark shuffle, adaptive query execution, broadcast joins, reading explain plans, partition pruning.

Stage 3: lakehouse table formats

A directory of Parquet files is not a table: two writers can corrupt it, a reader can see half a write, and there is no way to update a row. Table formats such as Apache Iceberg, Delta Lake and Apache Hudi add a metadata layer that records which files make up each snapshot. A commit writes new data files and then atomically swaps the table's current metadata pointer, so readers see either the old snapshot or the new one. That gives ACID writes, time travel, schema evolution and row-level updates on object storage.

Row-level changes come in two styles. Copy-on-write rewrites affected files at write time, so reads are fast and writes are expensive. Merge-on-read writes small delete or change files and merges them at read time, so writes are cheap until compaction catches up. Either way, table maintenance is part of the job: compact small files, expire old snapshots and remove orphaned files, or storage and planning time grow without bound.

Read next on this site: Iceberg with Spark, Delta Lake, Hudi, Iceberg with Hive, compaction, the metastore as catalog.

Stage 4: change data capture and streaming

Change data capture reads the database's own replication log, such as Postgres logical decoding, and emits every insert, update and delete as an event, usually through a connector such as Debezium into Kafka. This avoids polling queries, captures deletes, and preserves order per key. Kafka then acts as a durable, replayable buffer that many consumers can read independently; a schema registry keeps producers and consumers agreeing on the shape of events.

Streaming adds time as a first-class problem. Event time, when something happened, differs from processing time, when you saw it. Watermarks declare how late data may arrive, which lets an engine close windows and discard state. Exactly-once results come from combining replayable sources, checkpointed state and idempotent or transactional sinks; no single component provides it alone.

from pyspark.sql import functions as F

events = (spark.readStream.format("kafka")
          .option("kafka.bootstrap.servers", "broker:9092")
          .option("subscribe", "orders.v1").load()
          .select(F.from_json(F.col("value").cast("string"), ORDER_SCHEMA).alias("e"))
          .select("e.*"))

revenue = (events
           .withWatermark("event_time", "10 minutes")      # bound state, drop very late data
           .groupBy(F.window("event_time", "5 minutes"), "region")
           .agg(F.sum("amount").alias("revenue")))

(revenue.writeStream.format("iceberg").outputMode("append")
        .option("checkpointLocation", "s3://lake/_chk/revenue_5m")   # offsets + state
        .toTable("lake.metrics.revenue_5m"))

The checkpoint location stores Kafka offsets and window state together, so a restart resumes exactly where the last committed micro-batch ended. Move or delete it and the job reprocesses or skips data.

Read next on this site: CDC via logical decoding, schema registries, exactly-once processing, windowing, watermarks, Apache Flink, backpressure, Spark streaming watermarks.

Stage 5: orchestration and data quality

Orchestrators such as Airflow or Dagster run pipelines as dependency graphs, retry failures and backfill history. The property that makes all of that safe is idempotency: running a load for a given partition twice must give the same result as running it once. Write loads as overwrite-this-partition or MERGE-by-key, never as blind appends.

Quality checks belong inside the pipeline as blocking gates, not in a dashboard someone might look at. Test keys, nulls, ranges, referential integrity, row-count deltas and freshness. When a check fails, the new data is not published and consumers keep reading the last good snapshot, which table formats make easy.

# A blocking quality gate that runs after each load, before consumers see the data.
CHECKS = {
    "no_null_keys":   "SELECT count(*) FROM {t} WHERE order_id IS NULL",
    "unique_keys":    "SELECT count(*) - count(DISTINCT order_id) FROM {t}",
    "fresh":          "SELECT CASE WHEN max(updated_at) < now() - INTERVAL '2' HOUR THEN 1 ELSE 0 END FROM {t}",
    "no_neg_amounts": "SELECT count(*) FROM {t} WHERE amount < 0",
}

def gate(run_sql, table):
    failures = {n: v for n, q in CHECKS.items() if (v := run_sql(q.format(t=table))) != 0}
    if failures:
        raise RuntimeError(f"{table} failed quality gate: {failures}")  # do not publish

Read next on this site: idempotency, dead-letter queues, SLO burn-rate alerting.

Stage 6: data for AI systems

AI applications add three data products. Training and fine-tuning sets need provenance, deduplication, licence tracking and a frozen version per model. Retrieval indexes need chunking, embedding and, above all, freshness and access control: a document deleted or restricted in the source must disappear from the index. Evaluation sets need curation and versioning, because they are how every model or prompt change is judged.

Treat an embedding index as a derived table with a pipeline, not a one-off script. Key every chunk by a hash of its text, its permissions and the embedding model version, so a refresh re-embeds only what changed, and carry the document's permissions into the index so retrieval can filter by the caller's identity.

import hashlib

def refresh_embeddings(docs, index, embed, model_version):
    """Re-embed only chunks whose text, ACL or embedding model changed."""
    for doc in docs:                                  # docs from the lakehouse, with ACLs
        chunks = chunk_text(doc.body)
        acl = sorted(doc.allowed_groups)              # permission change => re-upsert
        for i, chunk in enumerate(chunks):
            key = f"{doc.id}:{i}"
            h = hashlib.sha256(f"{model_version}|{acl}|{chunk}".encode()).hexdigest()
            if index.get_hash(key) == h:
                continue                              # unchanged, skip the API cost
            index.upsert(key, embed(chunk), hash=h,
                         acl=acl, source_version=doc.updated_at)
        index.delete_chunks_from(doc.id, first_stale=len(chunks))  # doc shrank: drop tail
    index.delete_missing(valid_doc_ids={d.id for d in docs})   # deleted docs propagate

Read next on this site: incremental embedding pipelines, ACL-aware retrieval, hybrid search, golden datasets, data labelling, drift detection.

Worked example: orders from Postgres to dashboard and assistant

An e-commerce team wants revenue dashboards within minutes and a support assistant that can answer questions about an order. Debezium reads Postgres logical decoding and writes order changes to a Kafka topic keyed by order_id. A Spark job reads the topic in micro-batches into a staging table and applies the MERGE below to an Iceberg table. The log sequence number guards against applying an older event after a newer one when a batch is replayed.

-- Apply a micro-batch of CDC events (Debezium-style op codes) to an Iceberg table.
-- Deduplicate first: keep only the latest event per key in this batch.
MERGE INTO lake.sales.orders AS t
USING (
  SELECT * FROM (
    SELECT *, row_number() OVER (PARTITION BY order_id ORDER BY source_lsn DESC) AS rn
    FROM staging.order_changes
  ) WHERE rn = 1
) AS s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'd' AND s.source_lsn > t.source_lsn THEN DELETE
WHEN MATCHED AND s.source_lsn > t.source_lsn THEN UPDATE SET *
WHEN NOT MATCHED AND s.op <> 'd' THEN INSERT *;

A streaming aggregation feeds the five-minute revenue table; a nightly job compacts files and expires snapshots older than seven days. Quality gates block publication on duplicate keys or stale data. The assistant's index is refreshed from the same Iceberg table, carrying the customer id as an access attribute, so support staff retrieve only orders for the customer they are serving. When a customer is deleted for privacy reasons, the delete flows through CDC to the table and from the table to the index, and a weekly check confirms that no deleted id survives anywhere.

Failure modes and trade-offs

  • Fan-out joins. A non-unique key doubles revenue. Test uniqueness of every join key.
  • Small files. Streaming commits every few seconds create millions of files. Tune the trigger interval and compact.
  • Schema drift. A renamed source column silently becomes null. Enforce compatibility in the registry and fail loudly.
  • Replay duplicates. Blind appends double-count after a retry. Use MERGE by key with a version guard.
  • Stale or leaky AI indexes. Deleted or restricted documents remain retrievable. Propagate deletes and ACLs, and test for it.

The standing trade-offs: batch against streaming (simpler and cheaper against fresher), copy-on-write against merge-on-read, managed platforms against open components you operate, and normalised models against wide denormalised tables that are faster to query and harder to keep consistent. Choose freshness per consumer; most dashboards do not need seconds.

What to do next

  • Model one business process as a fact table with an explicit grain and two dimensions.
  • Run a Spark job, read its physical plan, and remove one shuffle with a broadcast join.
  • Create an Iceberg or Delta table, run a MERGE, then query an older snapshot with time travel.
  • Set up CDC from a local Postgres into Kafka and apply changes with a version-guarded MERGE.
  • Write a windowed streaming aggregation with a watermark and restart it from its checkpoint.
  • Add blocking quality gates to one pipeline and make a load fail on purpose.
  • Build an incremental embedding refresh that carries ACLs and prove a deleted document disappears.
Key takeaway: Data engineering is making data correct, fresh and governed at scale. Learn SQL and modelling first, then columnar batch processing, then the table formats that make object storage transactional, then CDC and streaming with honest handling of time. Wrap every load in idempotency and blocking quality gates. AI adds embedding indexes and evaluation sets as first-class derived data, and they inherit the same duties: freshness, lineage, deletion and access control.