Data modeling courseLesson 11 of 11
Data modeling course · Lesson 11 of 11
Modern Data Modeling: Streaming Models, Semantic Layers and dbt
Model streaming data, define metrics once in a semantic layer, structure a dbt project and build incremental models that are safe to rerun, with verified SQL.
On this page
- Sample data: an order event stream
- Data modelling for streaming
- What it is and why it matters
- How it works
- A worked example
- Joining streams to dimensions
- Pitfalls
- In interviews
- Semantic layer and metrics layer
- What it is and why it matters
- How it works
- A worked example: MetricFlow definitions
- Pitfalls
- In interviews
- Modelling in dbt
- What it is and why it matters
- How it works
- A worked example
- Pitfalls
- In interviews
- Idempotent incremental models
- What it is and why it matters
- The naive append, and why it breaks
- The idempotent version: merge with a lookback
- The same model in dbt
- Choosing a strategy
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
The modelling ideas in this course were developed for nightly batch loads into one database. Today the same models are fed by event streams, built by dbt, re-run many times a day and queried through semantic layers by BI tools and AI assistants. This lesson covers the four practices that make classic models work in that world, again using Kestrel Market (the fictional online shop from the rest of the course).
Sample data: an order event stream
Kestrel’s order service publishes an event every time an order is created, changed or cancelled. Events carry a unique id, the time the change happened in the shop (event time) and the time the warehouse received it (ingestion time). The stream delivers at least once, so duplicates happen, and the mobile app sends some events late.
CREATE TABLE raw_order_events (
event_id TEXT NOT NULL,
order_id TEXT NOT NULL,
op TEXT NOT NULL CHECK (op IN ('create', 'update', 'cancel')),
event_ts TIMESTAMP NOT NULL, -- when it happened
ingested_at TIMESTAMP NOT NULL, -- when the warehouse received it
customer_id TEXT,
net_amount NUMERIC(12,2)
);
INSERT INTO raw_order_events VALUES
('e1', 'O-1001', 'create', '2026-03-02 10:00', '2026-03-02 10:00:05', 'C1', 5797.00),
('e1', 'O-1001', 'create', '2026-03-02 10:00', '2026-03-02 10:00:09', 'C1', 5797.00), -- redelivered duplicate
('e2', 'O-1002', 'create', '2026-03-02 11:00', '2026-03-02 11:00:03', 'C2', 8999.00),
('e3', 'O-1001', 'update', '2026-03-02 18:00', '2026-03-02 18:00:02', 'C1', 5398.00),
('e4', 'O-1003', 'create', '2026-03-02 23:50', '2026-03-03 06:10:00', 'C1', 2898.00), -- phone was offline: late
('e5', 'O-1002', 'cancel', '2026-03-03 09:00', '2026-03-03 09:00:01', 'C2', NULL);
Data modelling for streaming
What it is and why it matters
A stream never “finishes”, so a streaming model has to say how the tables behave while data keeps arriving: what is a row, which time decides where it belongs, what happens to duplicates and to events that arrive after their window was reported. Most streaming bugs are modelling bugs: counting duplicates, grouping by arrival time instead of event time, or treating a change event as a new order.
How it works
Four design decisions shape a streaming model:
| Decision | Options | Kestrel choice |
|---|---|---|
| What a row means | Append-only event log (every event a row) or changelog / upsert table (one current row per key, rebuilt from events) | Keep the log; derive current state |
| Which time | Event time (when it happened) or processing / ingestion time (when it arrived) | Event time for business metrics; ingestion time for load bookkeeping |
| Duplicates | At-least-once delivery means duplicates; deduplicate on a unique event id | Unique key on event_id |
| Lateness | A watermark says “events older than this are probably all in”; later events either update results or are routed aside | Recompute affected windows; accept up to a stated delay |
The layers look like this:
raw_order_events (append-only, duplicates possible)
| deduplicate on event_id
v
order_events (append-only, unique events) -> event-level facts, activity streams
| latest event per order, cancels removed
v
orders_current (changelog materialised as state) -> current-state dimensions and facts
| aggregate by event-time window
v
revenue_by_hour (windowed aggregate, recomputed as late events arrive)
A worked example
Deduplicate on the event id. A unique constraint plus ON CONFLICT DO NOTHING makes the load safe to replay:
CREATE TABLE order_events (
event_id TEXT PRIMARY KEY, order_id TEXT NOT NULL, op TEXT NOT NULL,
event_ts TIMESTAMP NOT NULL, ingested_at TIMESTAMP NOT NULL,
customer_id TEXT, net_amount NUMERIC(12,2)
);
INSERT INTO order_events
SELECT DISTINCT ON (event_id) event_id, order_id, op, event_ts, ingested_at, customer_id, net_amount
FROM raw_order_events
ORDER BY event_id, ingested_at
ON CONFLICT (event_id) DO NOTHING;
SELECT (SELECT COUNT(*) FROM raw_order_events) AS raw_rows, (SELECT COUNT(*) FROM order_events) AS unique_events;
| raw_rows | unique_events |
|---|---|
| 6 | 5 |
Materialise the changelog into current state: the latest event per order wins, and cancelled orders drop out. Ordering by event time (with the event id as a tie-breaker) rather than arrival time means a late-arriving older event cannot overwrite a newer one:
CREATE VIEW orders_current AS
SELECT order_id, customer_id, net_amount, event_ts AS last_changed_at
FROM (
SELECT DISTINCT ON (order_id) *
FROM order_events
ORDER BY order_id, event_ts DESC, event_id DESC
) latest
WHERE op <> 'cancel';
SELECT * FROM orders_current ORDER BY order_id;
| order_id | customer_id | net_amount | last_changed_at |
|---|---|---|---|
| O-1001 | C1 | 5398.00 | 2026-03-02 18:00:00 |
| O-1003 | C1 | 2898.00 | 2026-03-02 23:50:00 |
Aggregate by event time. Daily revenue from current orders, grouped by the time the order happened:
SELECT o.created_at::DATE AS order_day, SUM(c.net_amount) AS revenue, COUNT(*) AS live_orders
FROM orders_current c
JOIN (SELECT order_id, MIN(event_ts) AS created_at FROM order_events GROUP BY order_id) o USING (order_id)
GROUP BY 1 ORDER BY 1;
| order_day | revenue | live_orders |
|---|---|---|
| 2026-03-02 | 8296.00 | 2 |
The late order O-1003 (placed at 23:50 on 2 March, received at 06:10 on 3 March) counts towards 2 March. Grouping by ingested_at would have put it on 3 March, and the 2 March number would change depending on when you ran the query. With event time, the 2 March figure changes only when late data arrives, which is honest, and the pipeline must recompute that day when it does.
Joining streams to dimensions
A streaming fact must pick up the dimension version that was valid at event time, exactly like the late-arriving fact lookup in the fact tables lesson. Stream processors call this a temporal (or “as of”) join against a versioned table. If the dimension row has not arrived yet, use an inferred member rather than dropping or delaying the event (see late-arriving dimensions). A pragmatic alternative many teams use is to denormalise at the source: include the few attributes analysts need (customer city, channel) in the event itself, so the event records the context as it was.
Pitfalls
- Aggregating by processing time and calling it daily revenue.
- No unique event id, so duplicates cannot be removed reliably. Ask producers for one; derive a deterministic hash only as a last resort.
- Treating updates as new facts, double counting an amended order. Model the stream as a changelog keyed by order.
- Unbounded state. Deduplicating “forever” in a stream processor needs unbounded memory; bound it with a time window and back it with a unique key in the warehouse.
- Out-of-order events applied in arrival order.
In interviews
Streaming modelling questions usually hide one of these traps: “count orders per minute from Kafka” (duplicates, event time, late events) or “keep a table of current order status” (changelog to state). Name event time versus processing time, deduplication by event id, watermarks and late-data handling, and how the same events feed both an append-only fact and a current-state table.
Semantic layer and metrics layer
What it is and why it matters
A semantic layer sits between the warehouse models and the tools that query them. It stores, in code, the definitions of entities (customer, order), dimensions (order date, channel) and metrics (revenue, average order value), and generates SQL for whatever slice a user asks for. Without it, “revenue” is redefined in every dashboard, notebook and spreadsheet, and the numbers disagree. With it, a BI tool, an API call or an AI assistant asking for “average order value by city last month” gets the same SQL.
How it works
A metrics definition has to capture what the earlier lessons taught about aggregation:
| Concern | What the definition captures | Kestrel example |
|---|---|---|
| Aggregation | How a measure aggregates | Revenue is SUM(net_amount) |
| Additivity | Which dimensions a measure can be summed across | Wallet balance is not summed over time; take the last value per period |
| Ratios | Numerator and denominator aggregated separately, then divided | AOV = revenue / distinct orders |
| Joins | Entities (keys) that let the layer join models safely | customer entity links orders to customers |
| Time | The default time dimension and grains | order_date by day, week, month |
The ratio rule is the one most often broken in hand-written dashboards. Average order value computed correctly, and computed as the average of daily averages:
CREATE TABLE fct_orders (order_id TEXT PRIMARY KEY, order_date DATE, city TEXT, net_amount NUMERIC(12,2));
INSERT INTO fct_orders VALUES
('O-1001', '2026-03-02', 'Pune', 5398.00),
('O-1002', '2026-03-02', 'Delhi', 8999.00),
('O-1003', '2026-03-02', 'Pune', 2898.00),
('O-1004', '2026-03-03', 'Mumbai', 499.00);
WITH daily AS (
SELECT order_date, SUM(net_amount) / COUNT(DISTINCT order_id) AS daily_aov
FROM fct_orders GROUP BY order_date
)
SELECT (SELECT ROUND(SUM(net_amount) / COUNT(DISTINCT order_id), 2) FROM fct_orders) AS correct_aov,
(SELECT ROUND(AVG(daily_aov), 2) FROM daily) AS average_of_daily_aov;
| correct_aov | average_of_daily_aov |
|---|---|
| 4448.50 | 3132.00 |
A semantic layer defines AOV once as a ratio of two metrics, and always generates the first query, at whatever grain is requested.
A worked example: MetricFlow definitions
The dbt Semantic Layer is powered by MetricFlow. Semantic models describe a dbt model’s entities, dimensions and measures; metrics are built from them. Other products (for example Cube, LookML in Looker, or warehouse-native semantic views) express the same ideas in their own syntax. The YAML below follows the MetricFlow spec documented by dbt; newer dbt releases have been revising this YAML, so check the docs for the version you run.
semantic_models:
- name: orders
model: ref('fct_orders')
defaults:
agg_time_dimension: order_date
entities:
- name: order
type: primary
expr: order_id
- name: customer
type: foreign
expr: customer_id
dimensions:
- name: order_date
type: time
type_params:
time_granularity: day
- name: city
type: categorical
measures:
- name: revenue
agg: sum
expr: net_amount
- name: orders
agg: count_distinct
expr: order_id
- name: wallet_balances
model: ref('fct_wallet_balance_daily')
defaults:
agg_time_dimension: snapshot_date
entities:
- name: customer
type: foreign
expr: customer_id
dimensions:
- name: snapshot_date
type: time
type_params:
time_granularity: day
measures:
- name: wallet_balance
agg: sum
expr: closing_balance
non_additive_dimension: # semi-additive: sum across customers, never across days
name: snapshot_date
window_choice: max # take the latest day in each requested period
metrics:
- name: revenue
label: Revenue
type: simple
type_params:
measure: revenue
- name: order_count
label: Orders
type: simple
type_params:
measure: orders
- name: average_order_value
label: Average order value
type: ratio
type_params:
numerator: revenue
denominator: order_count
- name: wallet_balance_owed
label: Wallet credit owed
type: simple
type_params:
measure: wallet_balance
A user then asks for metrics by dimensions, for example average_order_value by order__city and metric_time at month grain, and MetricFlow writes the joins and aggregations. The semi-additive wallet measure answers “credit owed per month” with each month’s last day, which is the logic from the fact tables lesson encoded once instead of in every dashboard.
Pitfalls
- A semantic layer on top of messy models. It encodes definitions; it cannot fix duplicated rows or mixed grains underneath. Model first.
- Defining ratios as averages, or pre-computing ratios in the table.
- Too many near-duplicate metrics (
revenue,revenue_v2,net_revenue_final). Treat metric definitions like an API: reviewed, documented, owned. - Bypassing it. If popular tools query tables directly, definitions drift again. Decide which consumers must go through the layer.
In interviews
“How do you make sure the CFO’s dashboard and the product team’s notebook show the same revenue?” The answer is conformed models plus a semantic layer: metrics defined once in version control, with aggregation, additivity and joins, generating SQL for every tool. Being able to explain ratio and semi-additive metrics in that layer shows real experience.
Modelling in dbt
What it is and why it matters
dbt turns SQL SELECT statements into tables and views in the warehouse, in dependency order, with tests and documentation. It has become the default way to build the transformation layer of a modern warehouse, so “how do you structure a dbt project?” is a standard modelling interview question.
How it works
Each model is a .sql file containing one SELECT. {{ ref('model') }} and {{ source('system', 'table') }} declare dependencies, which dbt uses to build the graph and to resolve the right schema per environment. The materialisation decides what dbt creates:
| Materialisation | Creates | Use for |
|---|---|---|
view |
A view | Light staging models |
table |
A table rebuilt on every run | Small or medium marts |
incremental |
A table that only processes new or changed rows | Large facts (next section) |
ephemeral |
Nothing; inlined as a CTE | Small reusable logic |
| snapshot | A Type 2 history table | Tracking changing source rows (see SCDs) |
A widely used layout, recommended in dbt’s own best-practice guides, maps directly onto the modelling layers of this course:
models/
staging/shop/ stg_shop__orders.sql, stg_shop__customers.sql one per source table: rename, cast, dedupe
intermediate/ int_order_lines_enriched.sql reusable joins and business logic
marts/core/ dim_customers.sql, fct_order_lines.sql Kimball stars (or OBTs) for consumers
snapshots/ customers_snapshot.yml Type 2 history of source rows
A worked example
A staging model renames, casts and deduplicates one source table, and nothing else:
-- models/staging/shop/stg_shop__orders.sql
with source as (
select * from {{ source('shop', 'order_events') }}
),
deduplicated as (
select *,
row_number() over (partition by event_id order by ingested_at) as rn
from source
)
select
event_id,
order_id,
lower(op) as op,
cast(event_ts as timestamp) as event_ts,
cast(ingested_at as timestamp) as ingested_at,
customer_id,
cast(net_amount as numeric(12,2)) as net_amount
from deduplicated
where rn = 1
A mart dimension built from the snapshot, with a deterministic surrogate key (the dbt_utils package’s generate_surrogate_key macro hashes the listed columns):
-- models/marts/core/dim_customers.sql
select
{{ dbt_utils.generate_surrogate_key(['customer_id', 'dbt_valid_from']) }} as customer_key,
customer_id,
customer_name,
city,
dbt_valid_from as valid_from,
dbt_valid_to as valid_to,
dbt_valid_to = to_date('9999-12-31') as is_current
from {{ ref('customers_snapshot') }}
Tests sit next to the models in YAML. The grain test and the relationship test are the two that catch most modelling bugs:
# models/marts/core/_core__models.yml
models:
- name: fct_order_lines
columns:
- name: order_line_id
data_tests:
- unique # the declared grain
- not_null
- name: customer_key
data_tests:
- not_null
- relationships: # every fact row finds a dimension row
to: ref('dim_customers')
field: customer_key
Older projects write tests: instead of data_tests:; recent dbt versions accept both. Recent versions also support unit tests for model logic with mocked inputs, which suit tricky SQL such as SCD or allocation logic.
Pitfalls
- Business logic in staging. Keep staging one-to-one with sources so every downstream model starts from the same clean base.
- Marts selecting from other teams’ marts in long chains, so one change breaks ten dashboards. Share logic through intermediate models and conformed dimensions.
- No grain tests. A
uniquetest on the fact’s key is the cheapest protection against fan-out joins. - Everything as
table. Large facts rebuilt in full every hour waste money; that is what incremental models are for.
In interviews
Expect “walk me through how you would structure a dbt project for this company”. Answer with sources, staging (one per source table), intermediate, marts (facts and dimensions with declared grains), snapshots for history, tests on keys and relationships, and the semantic layer on top. Mention materialisation choices and why.
Idempotent incremental models
What it is and why it matters
An incremental model processes only new or changed source rows and merges them into an existing table, instead of rebuilding it. It is how large facts stay affordable. It must also be idempotent: running it twice, re-running yesterday’s job, or replaying the last week after a bug fix must leave the table exactly as if each row had been processed once. The idempotency lesson covers the general principle; this section applies it to models.
Three ingredients make an incremental model idempotent:
- A unique key at the declared grain, so reprocessed rows replace rather than duplicate.
- A merge (upsert) or a delete-then-insert of a bounded slice, never a blind append.
- A lookback window, so late-arriving and updated rows inside the window are picked up on the next run.
The naive append, and why it breaks
CREATE TABLE fct_orders_append (order_id TEXT, net_amount NUMERIC(12,2), changed_at TIMESTAMP);
-- Run 1, then an automatic retry of the same run after a timeout
INSERT INTO fct_orders_append SELECT order_id, net_amount, last_changed_at FROM orders_current;
INSERT INTO fct_orders_append SELECT order_id, net_amount, last_changed_at FROM orders_current;
SELECT COUNT(*) AS rows, SUM(net_amount) AS revenue FROM fct_orders_append;
| rows | revenue |
|---|---|
| 4 | 16592.00 |
A retry doubled revenue. Appends are only safe when every row has a unique key that something downstream enforces.
The idempotent version: merge with a lookback
The model keeps a high-water mark (the latest ingested_at it has processed), re-reads a lookback window before it, and merges on order_id. Cancelled orders are deleted, updated orders replace their old values:
CREATE TABLE fct_orders_inc (
order_id TEXT PRIMARY KEY, -- the grain, and the merge key
customer_id TEXT,
net_amount NUMERIC(12,2) NOT NULL,
order_ts TIMESTAMP NOT NULL,
_ingested_at TIMESTAMP NOT NULL
);
CREATE VIEW incremental_batch AS
WITH watermark AS (
SELECT COALESCE(MAX(_ingested_at), TIMESTAMP '1900-01-01') - INTERVAL '1 day' AS since -- 1-day lookback
FROM fct_orders_inc
), changed_orders AS (
SELECT DISTINCT order_id FROM order_events, watermark WHERE ingested_at > watermark.since
)
SELECT DISTINCT ON (e.order_id)
e.order_id, e.op, e.customer_id, e.net_amount,
MIN(e.event_ts) OVER (PARTITION BY e.order_id) AS order_ts,
MAX(e.ingested_at) OVER (PARTITION BY e.order_id) AS _ingested_at
FROM order_events e
JOIN changed_orders USING (order_id)
ORDER BY e.order_id, e.event_ts DESC, e.event_id DESC;
MERGE INTO fct_orders_inc t
USING incremental_batch s ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'cancel' THEN DELETE
WHEN MATCHED THEN UPDATE SET customer_id = s.customer_id, net_amount = s.net_amount,
order_ts = s.order_ts, _ingested_at = s._ingested_at
WHEN NOT MATCHED AND s.op <> 'cancel' THEN
INSERT VALUES (s.order_id, s.customer_id, s.net_amount, s.order_ts, s._ingested_at);
SELECT order_id, net_amount, order_ts FROM fct_orders_inc ORDER BY order_id;
| order_id | net_amount | order_ts |
|---|---|---|
| O-1001 | 5398.00 | 2026-03-02 10:00:00 |
| O-1003 | 2898.00 | 2026-03-02 23:50:00 |
The batch reprocesses whole orders (every event for any order that changed inside the window) and keeps the latest state, so an update arriving without the original create still produces a correct row. Now a retry, then a new batch that includes a duplicate and a late update:
-- Retry of the same run: no change
MERGE INTO fct_orders_inc t
USING incremental_batch s ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'cancel' THEN DELETE
WHEN MATCHED THEN UPDATE SET customer_id = s.customer_id, net_amount = s.net_amount,
order_ts = s.order_ts, _ingested_at = s._ingested_at
WHEN NOT MATCHED AND s.op <> 'cancel' THEN
INSERT VALUES (s.order_id, s.customer_id, s.net_amount, s.order_ts, s._ingested_at);
-- New events: a new order, its redelivered duplicate, and a late price correction to O-1003
INSERT INTO order_events VALUES
('e6', 'O-1004', 'create', '2026-03-03 12:00', '2026-03-03 12:00:02', 'C3', 499.00),
('e7', 'O-1003', 'update', '2026-03-03 07:00', '2026-03-03 12:30:00', 'C1', 2799.00)
ON CONFLICT (event_id) DO NOTHING;
INSERT INTO order_events VALUES
('e6', 'O-1004', 'create', '2026-03-03 12:00', '2026-03-03 12:00:09', 'C3', 499.00)
ON CONFLICT (event_id) DO NOTHING;
MERGE INTO fct_orders_inc t
USING incremental_batch s ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'cancel' THEN DELETE
WHEN MATCHED THEN UPDATE SET customer_id = s.customer_id, net_amount = s.net_amount,
order_ts = s.order_ts, _ingested_at = s._ingested_at
WHEN NOT MATCHED AND s.op <> 'cancel' THEN
INSERT VALUES (s.order_id, s.customer_id, s.net_amount, s.order_ts, s._ingested_at);
SELECT order_id, net_amount, order_ts FROM fct_orders_inc ORDER BY order_id;
| order_id | net_amount | order_ts |
|---|---|---|
| O-1001 | 5398.00 | 2026-03-02 10:00:00 |
| O-1003 | 2799.00 | 2026-03-02 23:50:00 |
| O-1004 | 499.00 | 2026-03-03 12:00:00 |
A full rebuild from order_events would produce exactly this table, which is the test every incremental model should pass:
SELECT COUNT(*) AS differences
FROM (
(SELECT order_id, net_amount FROM fct_orders_inc
EXCEPT
SELECT order_id, net_amount FROM orders_current)
UNION ALL
(SELECT order_id, net_amount FROM orders_current
EXCEPT
SELECT order_id, net_amount FROM fct_orders_inc)
) d;
| differences |
|---|
| 0 |
The same model in dbt
In dbt the MERGE is generated for you. is_incremental() is true when the table already exists and the run is not a --full-refresh; {{ this }} refers to the existing table. The unique_key makes the merge an upsert:
-- models/marts/core/fct_orders.sql
{{
config(
materialized = 'incremental',
unique_key = 'order_id',
incremental_strategy = 'merge',
on_schema_change = 'append_new_columns'
)
}}
with events as (
select * from {{ ref('stg_shop__orders') }}
{% if is_incremental() %}
-- reprocess every order touched in the lookback window, not just the new events
where order_id in (
select order_id from {{ ref('stg_shop__orders') }}
where ingested_at > (select max(_ingested_at) - interval '1 day' from {{ this }})
)
{% endif %}
),
latest as (
select *,
min(event_ts) over (partition by order_id) as order_ts,
max(ingested_at) over (partition by order_id) as last_ingested_at,
row_number() over (partition by order_id order by event_ts desc, event_id desc) as rn
from events
)
select
order_id,
customer_id,
net_amount,
order_ts,
last_ingested_at as _ingested_at,
op = 'cancel' as is_cancelled -- soft delete: the merge strategy updates and inserts, it does not delete
from latest
where rn = 1
Two differences from the SQL version are deliberate. The default merge strategy updates and inserts but does not delete, so cancellations become a soft-delete flag that marts filter on. And the interval syntax in the watermark is warehouse-specific; adapters and cross-database macros differ.
For large event tables partitioned by time, dbt 1.9 added the microbatch strategy: you declare the event-time column and batch size, and dbt processes (and can retry or backfill) one time slice at a time, replacing each slice in full. No unique_key is needed because each batch overwrites its own period:
-- models/marts/core/fct_order_events_daily.sql
{{
config(
materialized = 'incremental',
incremental_strategy = 'microbatch',
event_time = 'event_ts',
batch_size = 'day',
lookback = 3,
begin = '2026-01-01'
)
}}
select event_id, order_id, op, event_ts, net_amount
from {{ ref('stg_shop__orders') }} -- dbt filters this ref to the batch's time range
Choosing a strategy
| Strategy | How it stays idempotent | Fits | Watch out for |
|---|---|---|---|
merge |
Upsert on unique_key |
Rows that update (orders, CDC tables) | Cost on very large targets; duplicates in the batch make the merge fail or behave unpredictably, so deduplicate first |
delete+insert |
Deletes target rows for the batch’s keys, then inserts | Engines where merge is slow or missing | Not atomic on every adapter; check transaction behaviour |
insert_overwrite |
Replaces whole partitions | Partitioned event tables (BigQuery, Spark) | Must always rebuild entire partitions, including late rows |
microbatch |
Replaces each event-time batch | Large time-series facts, backfills | Needs a reliable event-time column on inputs |
append |
Not idempotent by itself | Immutable events with a unique id deduplicated downstream | Retries duplicate rows |
Strategy support differs by adapter, so check the dbt documentation for your warehouse.
Pitfalls
- Watermark on event time instead of ingestion time, so late events with old timestamps are never picked up.
- No lookback, so a row updated just after a run’s watermark snapshot is missed.
- Filtering events instead of entities: processing only the new event for an order and losing the context of its earlier events. Reprocess the whole entity.
- Schema changes. Decide
on_schema_changedeliberately; the default (ignore) silently drops new columns. - Never testing a full refresh. Periodically compare the incremental table with a full rebuild, as shown above.
In interviews
“How do you make an incremental model safe to rerun?” Strong answers mention: a unique key at the grain, merge or partition overwrite rather than append, a watermark on ingestion time with a lookback window for late data, reprocessing whole entities, deletes or soft deletes for cancellations, and a reconciliation test against a full rebuild. Being able to say when microbatch or insert_overwrite is the better fit shows you have run these at scale.
Practice questions
You count orders per hour from a Kafka topic and the numbers are higher than the source database. List likely causes.
Duplicate deliveries (at-least-once) not removed by event id; update and cancel events counted as new orders instead of being applied as a changelog; grouping by processing time so late events land in the wrong hour; and replays after consumer restarts. Fix with deduplication on a unique event id, a changelog-to-state model keyed by order, and event-time windows that are recomputed when late events arrive.
What is the difference between event time and processing time, and which should a revenue report use?
Event time is when the business event happened; processing (or ingestion) time is when the system received or processed it. Revenue reports should use event time, so a late event is attributed to the period it belongs to. Ingestion time is still useful for incremental load watermarks and for auditing what was known when.
Why should average order value be defined as a ratio metric rather than computed in a table?
A ratio is non-additive. If you store a daily AOV and average it over a month, days with few orders get the same weight as busy days. Defining AOV as revenue divided by order count, both aggregated at the requested grain, gives the correct answer for any slice. A semantic layer enforces this by generating the SQL.
How would you structure a dbt project for an e-commerce company?
Sources declared in YAML; one staging model per source table (rename, cast, deduplicate, no business logic); intermediate models for reusable joins and logic; marts with Kimball facts and conformed dimensions at declared grains; snapshots for Type 2 history; tests for uniqueness, not-null and relationships on every key; incremental materialisation for large facts; and a semantic layer defining metrics on top of the marts.
Your incremental model uses WHERE event_ts > (SELECT MAX(event_ts) FROM this). What goes wrong, and how do you fix it?
Late events with an event time earlier than the current maximum are skipped forever, and a rerun after a partial failure may skip rows too. Use an ingestion-time watermark with a lookback window, reprocess every entity touched in that window, and merge on the unique key so reprocessing is harmless. Validate periodically against a full refresh.
When would you choose microbatch or insert_overwrite instead of merge?
When the table is a large, time-partitioned event fact where rows do not change individually and late data is bounded. Replacing whole time slices is cheaper than row-level merges on a huge target, makes backfills and retries simple, and needs no unique key. Use merge when individual rows change (orders, CDC-fed tables).
Key takeaways
- Streaming models separate an append-only, deduplicated event log from current-state tables and event-time aggregates that are recomputed when late data arrives.
- Use event time for business metrics, ingestion time for load watermarks, and a unique event id for deduplication.
- A semantic layer defines metrics once, including ratio and semi-additive behaviour, and generates consistent SQL for every tool.
- dbt projects map onto modelling layers: staging per source, intermediate logic, Kimball marts, snapshots for history, tests on every key.
- Incremental models are idempotent when they merge or overwrite on the grain, use an ingestion-time watermark with a lookback, reprocess whole entities and match a full rebuild.
Progress is saved in this browser only. No account needed.