dremio reflections are the mechanism that lets an interactive query hit terabytes of Parquet and Iceberg data sitting on cheap object storage and still come back in under a second — and they are the single feature that decides whether an open lakehouse feels like a warehouse or feels like a slow file scan. The promise of the lakehouse is seductive on paper: keep one open copy of your data in S3 or ADLS, describe it with an open table format, and let any engine query it in place instead of paying to load it into a proprietary warehouse first. The catch that every team hits in week two is that raw object-storage scans are not free — a dashboard that aggregates a billion rows on every refresh will grind, and the naive fix is to give up and copy the data into a columnar warehouse after all. Query acceleration is the escape hatch: pre-compute the expensive parts once, keep them fresh, and let the engine silently reuse them.
This guide is the walkthrough you wished existed the first time an interviewer asked "how does an open lakehouse serve BI dashboards without copying data into a warehouse," or "explain the difference between a raw and an aggregation materialization," or "walk me through how the optimizer decides to use a pre-computed result the query never mentions." It opens the engine in layers: the decoupled architecture that queries Iceberg tables in place, the raw and aggregation materializations that pre-compute the heavy lifting, the cost-based substitution that rewrites a plan to read a materialization transparently, the Apache Arrow columnar format and vectorized execution that make every operator fast, and the semantic layer of virtual datasets where governance and acceleration meet. Each section pairs a teaching block with a Solution-Tail interview answer — code, a step-by-step trace, an output table, then a concept-by-concept breakdown of why it works.
When you want hands-on reps immediately after reading, drill the query optimization practice library →, sharpen your data-modeling instincts on the database practice library →, and rehearse the fundamentals on the SQL practice library →.
On this page
- Dremio & the open lakehouse architecture
- Reflections — raw & aggregation
- How the optimizer transparently substitutes reflections
- Apache Arrow & vectorized performance
- Semantic layer & governance patterns
- Cheat sheet — Dremio reflection & lakehouse recipes
- Frequently asked questions
- Practice on PipeCode
1. Dremio & the open lakehouse architecture
The open lakehouse queries data in place — and reflections are where the speed comes from
The one-sentence invariant: an open lakehouse keeps a single open copy of data as Parquet files described by an open table format like Apache Iceberg on object storage, and a decoupled query engine reads those files directly at query time instead of loading them into a proprietary store — which removes the copy, the vendor lock-in, and the nightly load window, but shifts the entire performance burden onto the engine, and that is exactly the gap that query acceleration through materializations is built to close. Every architectural decision downstream — whether you copy data into a warehouse, whether your dashboards are fast, whether your storage bill is one copy or five — flows from whether the engine can make an in-place scan feel interactive. Understanding where the acceleration comes from is understanding the whole platform.
The three planes of the open lakehouse.
- Storage. Parquet files in an object store (S3, ADLS, GCS, MinIO). Columnar on disk, compressed, immutable once written. This is the durable copy — the only copy — and it is open: any engine that speaks Parquet can read it.
-
Table format. Apache Iceberg (or Delta Lake, or plain Hive-style directories) layers table semantics over the raw files: a schema, atomic snapshots, hidden partitioning, and a manifest of which files belong to which snapshot. Iceberg is what turns "a folder of Parquet" into "a table you can
INSERT,UPDATE,DELETE, and time-travel." - Compute. A query engine — Dremio, Trino, Spark — that plans and executes SQL against the table format. Compute is stateless and elastic: it scales independently of storage, spins up for a workload, and holds no durable data of its own. The engine is where reflections, Arrow, and the semantic layer live.
Dremio's decoupled internals.
- Coordinator. Accepts SQL, holds the catalog and the semantic layer (spaces, virtual datasets), runs the planner/optimizer, and hands physical plan fragments to executors. It is the brain; it does not scan data.
- Executors. Stateless workers that scan Parquet/Iceberg files, run the operators (filters, joins, aggregations) over Apache Arrow record batches in memory, and stream results back. Executors are grouped into elastic engines you can size per workload.
- Reflection store. A configured location on the data lake where Dremio writes its materializations as Iceberg tables. Reflections are not a hidden cache in RAM — they are durable, open, columnar copies that survive restarts.
- Sources. Connectors to Iceberg catalogs, object stores, and relational databases. The lakehouse can federate: a query can join an Iceberg fact table on S3 to a dimension table in Postgres without copying either.
The 2026 reality — why query-in-place won.
- One copy, many engines. The lakehouse eliminated the "load it into the warehouse first" step. The same Iceberg table feeds Dremio for BI, Spark for ML feature pipelines, and Flink for streaming — no per-engine copy, no drift between copies.
- Open formats ended lock-in. Iceberg and Parquet are specifications, not products. Your data outlives any single engine choice; migrating engines is a config change, not a re-platform.
- The performance tax is real. Reading a billion rows off object storage on every dashboard refresh is slow and expensive. Without acceleration, teams quietly rebuild a warehouse. With it, the lakehouse serves sub-second BI on the open copy. Acceleration is the load-bearing feature, not a nice-to-have.
What interviewers listen for.
- Do you describe the lakehouse as "query the open copy in place" rather than "another warehouse"? — required framing.
- Do you name Iceberg's snapshot + manifest model as what enables in-place ACID, not just "files in S3"? — senior signal.
- Do you locate acceleration in the engine (materializations + Arrow), not in the storage? — required answer.
- Do you separate storage / table-format / compute as three independently-scaling planes? — senior signal.
- Do you say reflections are durable open Iceberg copies, not a RAM cache? — senior signal.
Worked example — mapping the lakehouse stack for a sales dataset
Detailed explanation. The most useful artifact for a lakehouse architecture interview is a clean three-plane map of a real dataset. Every senior lakehouse discussion converges on "where does the data live, what describes it, and what runs the query" within the first ten minutes. Walk through building the map for a sales.orders dataset that has to serve a BI dashboard and a data-science feature job at the same time.
-
The dataset.
sales.orders— 4 billion rows, ~1.2 TB of Parquet, partitioned by order date. - The consumers. A finance dashboard that aggregates revenue by region and day, and a churn-model feature job that reads raw rows.
- The constraint. One physical copy on S3; no nightly load into a separate warehouse.
Question. Map sales.orders onto the three planes and name what each consumer touches.
Input.
| Concern | Choice |
|---|---|
| Storage | Parquet on S3 (s3://lake/sales/orders/) |
| Table format | Apache Iceberg (catalog-managed snapshots) |
| Compute | Dremio (coordinator + elastic executors) |
| BI consumer | aggregation query, needs sub-second |
| ML consumer | raw row scan, tolerant of seconds |
Code.
-- The physical Iceberg table (one open copy)
CREATE TABLE lakehouse.sales.orders (
order_id BIGINT,
customer_id BIGINT,
order_ts TIMESTAMP,
order_date DATE,
region VARCHAR,
status VARCHAR,
amount DECIMAL(12,2)
)
PARTITION BY (order_date)
STORE AS (type => 'iceberg')
LOCATION 's3://lake/sales/orders/';
-- BI consumer — an aggregation the dashboard runs every refresh
SELECT region, order_date, SUM(amount) AS revenue, COUNT(*) AS orders
FROM lakehouse.sales.orders
WHERE order_date >= DATE '2026-01-01'
GROUP BY region, order_date;
-- ML consumer — a raw scan for features
SELECT order_id, customer_id, order_ts, amount, status
FROM lakehouse.sales.orders
WHERE order_date >= DATE '2026-06-01';
Step-by-step explanation.
- The
CREATE TABLE ... STORE AS icebergstatement writes Parquet files under the S3 location and records an Iceberg snapshot pointing at them. There is exactly one copy of the data; the "table" is the Iceberg metadata layered over those files. -
PARTITION BY (order_date)is Iceberg hidden partitioning — the engine can prune whole date partitions from a query's file list using table metadata, before reading a single row. This is the cheapest acceleration and it is free with the table format. - The BI aggregation touches every row in the date range and collapses them into a small grouped result. On raw Parquet this is I/O-heavy and repeated on every refresh — the exact shape that a materialization will later accelerate.
- The ML scan reads raw columns for a broad date range. It is tolerant of a few seconds and does not need pre-aggregation — it may benefit from a sorted/partitioned raw materialization but not from a pre-aggregated one.
- Both consumers point at the same logical table. Nothing is copied into a second store; the engine, not the storage, is responsible for making the BI path fast. That responsibility is the reflection.
Output.
| Plane | Component | BI consumer | ML consumer |
|---|---|---|---|
| Storage | Parquet on S3 | scans date range | scans date range |
| Table format | Iceberg | partition-prunes by date | partition-prunes by date |
| Compute | Dremio engine | needs acceleration | raw scan is fine |
| Acceleration | (reflection) | aggregation reflection | optional raw reflection |
Rule of thumb. Draw the three planes before proposing any acceleration. Storage and table-format give you free pruning; the engine gives you materializations. If a consumer is slow, the fix lives in the compute plane, not in copying data into a new store.
Worked example — the no-copy cost model vs ETL-into-warehouse
Detailed explanation. The senior argument for the lakehouse is economic, not just architectural. The classic pattern copies data from the lake into a proprietary warehouse nightly; the lakehouse queries the lake copy directly and accelerates it. Quantify the difference for the sales.orders dataset so the trade-off is concrete rather than ideological.
- Warehouse pattern. Nightly load of 1.2 TB into a columnar warehouse; storage billed twice (lake + warehouse); a load window; drift between copies until the next load.
- Lakehouse pattern. Query the lake copy; add a ~40 GB aggregation reflection for the dashboard; storage billed once plus the small reflection; no load window; freshness bounded by reflection refresh.
Question. Compare storage footprint, freshness, and the failure surface of the two patterns.
Input.
| Dimension | ETL-into-warehouse | Lakehouse + reflection |
|---|---|---|
| Copies of data | 2 (lake + warehouse) | 1 + small reflection |
| Extra storage | +1.2 TB | +~40 GB (reflection) |
| Freshness | last nightly load | reflection refresh interval |
| Load window | nightly batch | none (incremental refresh) |
| Drift | until next load | bounded by refresh |
Code.
-- Lakehouse: no load job. One reflection accelerates the dashboard.
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION orders_by_region_day
USING
DIMENSIONS (region, order_date)
MEASURES (amount (SUM, COUNT), order_id (COUNT));
-- Compare: the warehouse pattern would instead run, every night,
-- an EXTRACT from the lake + a LOAD into a separate columnar store,
-- doubling storage and adding a batch window. No such job here.
Step-by-step explanation.
- In the warehouse pattern, the 1.2 TB lake copy is extracted and loaded into a second columnar store nightly. You now pay for two copies and you inherit a load window during which the warehouse is stale relative to the lake.
- In the lakehouse pattern, the dashboard's aggregation is materialized once as
orders_by_region_day. Because it stores only the groupedregion × order_daterows with a handful of measures, it is roughly 40 GB, not 1.2 TB — a fraction of a full copy. - Freshness in the warehouse pattern is "as of last night's load." Freshness in the lakehouse pattern is "as of the last reflection refresh," and refresh can be incremental (only new partitions), so it runs in minutes and can be scheduled hourly.
- The failure surfaces differ. The warehouse pattern fails as a broken load job (dashboard silently stale until someone notices). The lakehouse pattern fails as a stale reflection, and if a reflection is stale or missing the engine simply falls back to scanning the raw table — slower, but never wrong.
- The economic punchline: the lakehouse serves the same dashboard at one-copy storage plus a small materialization, with no batch window and a graceful-degradation failure mode. That is why query-in-place plus acceleration displaced copy-into-warehouse for a large class of BI workloads.
Output.
| Metric | ETL-into-warehouse | Lakehouse + reflection |
|---|---|---|
| Storage billed | 1.2 TB × 2 | 1.2 TB + ~40 GB |
| Dashboard latency | fast (post-load) | fast (via reflection) |
| Freshness | ~24 h | ~1 h (incremental) |
| Failure mode | stale until load fixed | falls back to raw scan |
| Vendor lock-in | warehouse format | open Iceberg |
Rule of thumb. Frame the lakehouse choice as "one open copy plus targeted materializations" versus "a second proprietary copy plus a load window." The reflection is the small, purpose-built acceleration that replaces the whole second copy — never justify the lakehouse on openness alone; justify it on one-copy economics with acceleration filling the speed gap.
Worked example — the Iceberg table anatomy reflections build on
Detailed explanation. Reflections are stored as Iceberg tables and they accelerate queries against Iceberg tables, so the whole acceleration story rests on the Iceberg metadata model. A senior engineer can sketch it: a table points at a current snapshot, a snapshot points at a manifest list, manifests point at data files, and each manifest carries per-file column statistics used for pruning. Walk the model.
- Snapshot. An atomic version of the table. Every commit (insert, update, delete, refresh) creates a new snapshot; readers see a consistent snapshot.
- Manifest list + manifests. The snapshot references a manifest list; each manifest lists data files plus per-file stats (row counts, null counts, min/max per column, partition values).
- Metadata pruning. The engine uses those per-file min/max and partition values to skip files that cannot match a predicate — before any Parquet is read.
Question. Show how Iceberg metadata prunes files for a date-filtered query and why that matters for reflections.
Input.
| Iceberg object | Role |
|---|---|
| table metadata | points at current snapshot + schema |
| snapshot | atomic version; points at manifest list |
| manifest | lists data files + per-file min/max stats |
| data file | a Parquet file with actual rows |
Code.
-- Time-travel and metadata inspection make the model visible.
-- 1. Read the table as of a specific snapshot (consistent version)
SELECT COUNT(*)
FROM lakehouse.sales.orders
AT SNAPSHOT '8172049571043210';
-- 2. Read as of a wall-clock time (Iceberg picks the snapshot)
SELECT COUNT(*)
FROM lakehouse.sales.orders
AT TIMESTAMP '2026-09-01 00:00:00';
-- 3. A date-filtered query the engine prunes via manifest stats
SELECT SUM(amount)
FROM lakehouse.sales.orders
WHERE order_date = DATE '2026-08-15'; -- only files whose
-- min/max order_date span
-- Aug-15 are read
Step-by-step explanation.
-
AT SNAPSHOTandAT TIMESTAMPexpose Iceberg's snapshot model directly: the table is a pointer to an immutable version, and every read is against a consistent snapshot. Reflections capture data as of a specific source snapshot, which is how the engine knows whether a reflection is current. - The date-filtered query never reads the whole table. The planner consults each manifest's per-file min/max for
order_date, discards files whose range cannot contain2026-08-15, and reads only the survivors. This is metadata pruning and it happens before I/O. - Because partitioning is hidden, you do not write
WHERE partition_path = ...; you write the natural predicate onorder_dateand Iceberg maps it to partition pruning automatically. Reflections inherit this — a partitioned raw reflection prunes the same way. - Every commit creates a new snapshot, so an incremental reflection refresh can ask "which files/partitions changed since the snapshot I last materialized?" and refresh only those. The snapshot id is the anchor for incremental refresh.
- This anatomy is why reflections are cheap to keep fresh and cheap to substitute: they are Iceberg tables with the same pruning, the same snapshots, and the same statistics as the sources they accelerate.
Output.
| Query | Files scanned | Mechanism |
|---|---|---|
full table COUNT(*)
|
all data files | manifest metadata count (no row read) |
AT SNAPSHOT id |
files in that snapshot | snapshot pointer |
order_date = '2026-08-15' |
one partition's files | min/max + partition pruning |
| incremental refresh | changed files since snapshot | snapshot diff |
Rule of thumb. Learn the snapshot → manifest → data-file chain cold; it is the substrate for everything else. Reflections are Iceberg tables that inherit its pruning and snapshotting, which is exactly why they are fast to query and cheap to refresh incrementally.
Data engineering interview question on lakehouse architecture
A senior interviewer often opens with: "Your BI team wants sub-second dashboards on a 4-billion-row Iceberg orders table sitting in S3, and leadership refuses to fund a separate proprietary warehouse copy. Walk me through the open-lakehouse architecture you'd stand up, where the query acceleration comes from, and how you'd keep the accelerated results fresh without a nightly load window."
Solution Using a decoupled Dremio engine over Iceberg with a targeted aggregation reflection
-- 1. The single open copy: an Iceberg table over Parquet on S3
CREATE TABLE lakehouse.sales.orders (
order_id BIGINT, customer_id BIGINT, order_ts TIMESTAMP,
order_date DATE, region VARCHAR, status VARCHAR, amount DECIMAL(12,2)
)
PARTITION BY (order_date)
STORE AS (type => 'iceberg')
LOCATION 's3://lake/sales/orders/';
-- 2. The acceleration: one aggregation reflection for the dashboard shape
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION dash_region_day
USING
DIMENSIONS (region, order_date)
MEASURES (amount (SUM, COUNT));
-- 3. Freshness without a load window: incremental, hourly refresh
ALTER TABLE lakehouse.sales.orders
SET ACCELERATION REFRESH POLICY
REFRESH METHOD INCREMENTAL BY (order_date)
REFRESH EVERY 1 HOURS
EXPIRE AFTER 24 HOURS;
Engine topology (decoupled)
===========================
Coordinator -> catalog + semantic layer + planner (chooses reflections)
Executors -> stateless, scan Parquet, process Arrow batches, elastic
Reflection -> dash_region_day stored as an Iceberg table in the reflection
store store on S3 (durable, open, columnar — not a RAM cache)
Fallback -> if the reflection is stale/absent, plan reads raw orders
Step-by-step trace.
| Step | Before (warehouse copy) | After (lakehouse + reflection) |
|---|---|---|
| Copies of data | 2 (lake + warehouse) | 1 + small reflection |
| Dashboard aggregation | pre-loaded in warehouse | served by dash_region_day
|
| Freshness | nightly load | hourly incremental refresh |
| Load window | required | none |
| Failure mode | stale until load fixed | falls back to raw Iceberg scan |
| Storage delta | +1.2 TB | +~40 GB reflection |
After this design, the dashboard aggregation is answered from the ~40 GB dash_region_day reflection in well under a second, the raw 1.2 TB copy stays the single source of truth on S3, hourly incremental refresh keeps the reflection within an hour of source, and if the reflection is ever stale the planner transparently reads the raw table — slower but always correct.
Output:
| Metric | Value |
|---|---|
| Dashboard latency | < 1 s (via reflection) |
| Source copies | 1 (open Iceberg) |
| Reflection footprint | ~40 GB |
| Freshness p99 | ~1 h (incremental) |
| Degradation on stale reflection | raw scan (correct, slower) |
Why this works — concept by concept:
- Decoupled compute and storage — executors are stateless and elastic; the durable data lives once on S3 as Iceberg. Scaling compute for a workload never duplicates data, which is what makes the one-copy economics real.
- Iceberg snapshots + hidden partitioning — give the engine atomic versions and free partition/metadata pruning, so even the raw fallback scan reads only the relevant date partitions rather than the whole table.
-
Aggregation reflection — pre-computes the dashboard's
region × order_dateroll-up once and stores it as a small Iceberg table, collapsing a billion-row scan into a few-thousand-row read on every dashboard refresh. - Incremental refresh by order_date — uses the snapshot diff to re-materialize only new date partitions, so freshness is an hourly minutes-long job instead of a nightly full rebuild — no load window.
- Cost — one full copy plus a ~40 GB reflection and a small hourly refresh, versus two full copies plus a nightly load. Query cost drops from O(rows in range) to O(groups) for the dashboard; the raw fallback preserves correctness. Net: warehouse-class latency at lakehouse storage economics.
Database
Topic — database
Lakehouse and data-modeling problems
2. Reflections — raw & aggregation
A reflection is a materialization stored as Iceberg — raw ones accelerate scans and joins, aggregation ones accelerate roll-ups
The mental model in one line: a dremio reflection is a pre-computed, physically-optimized copy of a dataset — either a raw reflection that stores selected columns sorted and partitioned for fast scans/filters/joins, or an aggregation reflection that stores pre-computed dimensions and measures for fast group-bys — persisted as an Iceberg table in the reflection store and kept in sync with its source by a refresh policy, so the engine can answer expensive queries from the cheap materialization instead of the raw table. Reflections are the acceleration engine of the lakehouse: everything about substitution, freshness, and cost budgeting in the rest of this guide is about choosing, shaping, and maintaining them.
Raw reflection anatomy — a physically-optimized copy.
-
Display columns. The subset of columns the reflection stores (
USING DISPLAY (...)). A query that needs only those columns can be served entirely from the reflection. -
Sort.
LOCALSORT BY (...)clusters rows so range filters and merge joins on the sort key read fewer files. Sorting is the single biggest lever for filter-heavy queries. -
Partition.
PARTITION BY (...)splits the reflection so predicate values prune whole partitions, exactly like a source Iceberg table. -
Distribution.
DISTRIBUTE BY (...)co-locates rows by a join key so a downstream join can avoid a shuffle. Choose the common join key.
Aggregation reflection anatomy — a pre-computed cube.
-
Dimensions. The grouping columns (
DIMENSIONS (...)). Any query grouping by a subset of these dimensions can be served — you group up from a finer grain to a coarser one for free. -
Measures. The aggregated columns and the functions kept for each (
MEASURES (amount (SUM, COUNT, MIN, MAX))). StoreSUMandCOUNTtogether soAVGcan be derived (SUM/COUNT) without a raw scan. -
Additivity.
SUM,COUNT,MIN,MAXare additive/re-aggregatable across finer groups;COUNT(DISTINCT)is not exactly re-aggregatable, so Dremio keeps an approximate sketch (HyperLogLog) for distinct counts on aggregation reflections. - Grain. The reflection's grain is its full dimension set; queries at that grain or coarser match, queries at a finer grain do not.
The reflection store — where materializations live.
- Durable and open. Reflections are written as Iceberg tables (Parquet under the hood) to a configured reflection-store location on the data lake. They survive engine restarts and are readable as ordinary columnar files.
- Not a RAM cache. A common misconception is that reflections are an in-memory cache. They are persisted materializations; memory caching (the columnar cloud cache) is a separate layer that speeds up reading them.
- Footprint budgeting. Each reflection costs storage (the materialized rows) and refresh compute. Aggregation reflections are usually tiny relative to the source; raw reflections can approach source size if they keep most columns.
Refresh policies — keeping reflections in sync.
- Full refresh. Rebuild the entire reflection from the current source snapshot. Simple, correct, expensive on large sources.
- Incremental refresh. Materialize only what changed since the last refresh, anchored on a monotonically increasing column or on new Iceberg snapshots/partitions. Minutes instead of hours.
-
Schedule + expiry.
REFRESH EVERY Nsets the cadence;EXPIRE AFTER Nmarks a reflection too stale to use, forcing fallback to the raw source rather than serving stale data. -
Manual force.
ALTER TABLE ... REFRESH REFLECTIONStriggers an immediate refresh, e.g. after a large backfill.
Common interview probes on reflections.
- "Raw vs aggregation reflection — when each?" — raw for scans/filters/joins on selected columns; aggregation for group-by/roll-up shapes.
- "Where are reflections stored?" — durable Iceberg tables in the reflection store, not a RAM cache.
- "How do you keep
AVGcorrect in an aggregation reflection?" — storeSUMandCOUNT, deriveAVG = SUM/COUNT. - "Why is
COUNT(DISTINCT)special?" — not exactly re-aggregatable; served via an approximate HLL sketch.
Worked example — creating a raw reflection for filter-and-join queries
Detailed explanation. The canonical raw reflection accelerates a family of queries that filter on a date/region and join orders to customers. You store just the columns those queries touch, sort by the common filter column, and distribute by the join key so the join avoids a shuffle. Build it and reason about what it now serves.
-
Target queries. Filter by
order_date/region, project a handful of columns, join tocustomersoncustomer_id. - Display columns. Only the columns those queries need.
-
Sort/partition/distribute. Sort by
order_ts, partition byregion, distribute bycustomer_id.
Question. Create a raw reflection on orders tuned for the filter-and-join workload and state which queries it can serve.
Input.
| Knob | Choice | Why |
|---|---|---|
| DISPLAY | order_id, customer_id, order_ts, region, status, amount | columns the workload projects |
| PARTITION BY | region | filters on region prune partitions |
| LOCALSORT BY | order_ts | range filters on time read fewer files |
| DISTRIBUTE BY | customer_id | join to customers avoids a shuffle |
Code.
-- Raw reflection: a physically-optimized copy for filter + join queries
ALTER TABLE lakehouse.sales.orders
CREATE RAW REFLECTION orders_raw_fastscan
USING DISPLAY (order_id, customer_id, order_ts, region, status, amount)
PARTITION BY (region)
LOCALSORT BY (order_ts)
DISTRIBUTE BY (customer_id);
-- A query this raw reflection can fully serve
SELECT o.order_id, o.amount, c.segment
FROM lakehouse.sales.orders o
JOIN lakehouse.crm.customers c ON c.customer_id = o.customer_id
WHERE o.region = 'EU'
AND o.order_ts >= TIMESTAMP '2026-08-01 00:00:00';
Step-by-step explanation.
-
USING DISPLAY (...)restricts the reflection to the six columns the workload projects. Because a raw reflection is a column subset, it is smaller than the source and the scan reads fewer bytes — the first win. -
PARTITION BY (region)means theregion = 'EU'predicate prunes to the EU partition of the reflection, skipping every other region's files via manifest metadata. -
LOCALSORT BY (order_ts)clusters rows in time order within each partition, so theorder_ts >= ...range filter reads a contiguous tail of files rather than scanning the whole partition. -
DISTRIBUTE BY (customer_id)co-locates rows by the join key. When the query joinsorderstocustomersoncustomer_id, the engine can perform the join without reshuffling the orders side across executors — cutting the most expensive part of a distributed join. - The query never names
orders_raw_fastscan. It writes ordinary SQL againstorders; the planner recognizes the reflection covers the columns, partitioning, and join key, and substitutes it. The raw reflection is invisible acceleration for an entire query family, not a single query.
Output.
| Query shape | Served by raw reflection? | Why |
|---|---|---|
| filter region + time, project subset | yes | columns + partition + sort all match |
| join on customer_id | yes (no shuffle) | distributed by customer_id |
| project a column not in DISPLAY | no | column absent from reflection |
| filter on a non-sorted, non-partition column | partially | reflection used, but full scan of it |
Rule of thumb. Shape a raw reflection to a query family, not a single query: put the projected columns in DISPLAY, the filter column in LOCALSORT, the high-cardinality equality filter in PARTITION, and the join key in DISTRIBUTE. A raw reflection that mirrors the workload's access pattern is what turns an object-store scan into an interactive one.
Worked example — creating an aggregation reflection for a dashboard
Detailed explanation. The dashboard groups revenue by region and day and sometimes by status. An aggregation reflection stores those dimensions and the additive measures once, so every dashboard refresh reads a few thousand pre-aggregated rows instead of re-scanning billions. Build it and show what it can and cannot roll up.
-
Dimensions.
region,status,order_date— the union of everything the dashboard groups by. -
Measures.
amountwithSUM,COUNT,MIN,MAX;customer_idwithCOUNTfor order counts. -
Derivations.
AVG(amount) = SUM(amount)/COUNT(amount)derived from stored measures.
Question. Create an aggregation reflection that serves "revenue by region by day," "revenue by region," and "average order value by status," and identify a query it cannot serve.
Input.
| Component | Value |
|---|---|
| DIMENSIONS | region, status, order_date |
| MEASURES | amount (SUM, COUNT, MIN, MAX) |
| Grain | region × status × order_date |
| Derivable | AVG (SUM/COUNT), revenue roll-ups |
Code.
-- Aggregation reflection: a small pre-computed cube for the dashboard
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION orders_cube
USING
DIMENSIONS (region, status, order_date)
MEASURES (amount (SUM, COUNT, MIN, MAX));
-- Q1 — served: group by a subset of the cube's dimensions
SELECT region, order_date, SUM(amount) AS revenue
FROM lakehouse.sales.orders
GROUP BY region, order_date;
-- Q2 — served: coarser roll-up (region only) re-aggregates the cube
SELECT region, SUM(amount) AS revenue
FROM lakehouse.sales.orders
GROUP BY region;
-- Q3 — served: AVG derived from stored SUM and COUNT
SELECT status, SUM(amount) / NULLIF(COUNT(*),0) AS avg_order_value
FROM lakehouse.sales.orders
GROUP BY status;
-- Q4 — NOT served: groups by a dimension the cube doesn't carry
SELECT customer_id, SUM(amount)
FROM lakehouse.sales.orders
GROUP BY customer_id; -- customer_id is not a cube dimension
Step-by-step explanation.
- The reflection's grain is
region × status × order_date. Q1 groups by a subset of those dimensions (region,order_date), so the engine re-aggregates the cube up to the coarser grain — reading only cube rows, never the source. - Q2 groups by
regionalone. This is an even coarser roll-up;SUMis additive, so summing the cube's per-day-per-status rows up to region is exact. Coarser-than-cube queries always match; finer-than-cube queries never do. - Q3 asks for an average by status. Because the cube stores both
SUM(amount)andCOUNT(amount), the engine derivesAVG = SUM/COUNTwithout touching raw rows. This is why you storeSUMandCOUNTtogether rather than a rawAVGmeasure. - Q4 groups by
customer_id, which is not a cube dimension. The cube has already collapsed away per-customer detail, so it cannot answer this — the planner falls back to the raw table (or a raw reflection if one exists). Aggregation reflections trade detail for size; you cannot recover a dimension you did not keep. - The cube is small: its row count is bounded by
|region| × |status| × |dates|, typically thousands of rows even over billions of source rows. That size is why aggregation reflections give the biggest acceleration for the least storage on BI workloads.
Output.
| Query | Grain vs cube | Served? |
|---|---|---|
| region × day revenue | subset (coarser) | yes |
| region revenue | coarser roll-up | yes |
| avg order value by status | derived SUM/COUNT | yes |
| revenue by customer | finer (missing dim) | no → raw fallback |
Rule of thumb. Size the aggregation reflection's DIMENSIONS to the union of the dashboard's group-by columns and always keep SUM + COUNT so AVG derives for free. Remember the grain rule: the cube serves queries at its grain or coarser, never finer — a dimension you drop is detail you cannot get back without the raw table.
Worked example — incremental refresh anchored on Iceberg snapshots
Detailed explanation. A full refresh of a reflection over a 4-billion-row source is expensive and defeats the "no load window" promise. Incremental refresh materializes only what changed since the last refresh. On Iceberg sources, the anchor is the snapshot: Dremio asks "which partitions/files are new since the snapshot I last materialized?" and re-aggregates only those. Configure it and reason about correctness.
-
Anchor. Incremental by
order_date(append-mostly) or by Iceberg snapshot diff. - Cadence. Refresh every hour; expire after 24 hours to prevent serving badly-stale data.
- Correctness. Additive measures re-aggregate cleanly; late-arriving updates to old partitions need a full refresh of those partitions.
Question. Configure incremental refresh for orders_cube and describe the refresh math for a new day's data.
Input.
| Parameter | Value |
|---|---|
| Refresh method | incremental |
| Anchor | order_date (new partitions) |
| Cadence | every 1 hour |
| Expiry | 24 hours |
| Force refresh | after backfills |
Code.
-- Configure incremental, scheduled refresh with an expiry guard
ALTER TABLE lakehouse.sales.orders
SET ACCELERATION REFRESH POLICY
REFRESH METHOD INCREMENTAL BY (order_date)
REFRESH EVERY 1 HOURS
EXPIRE AFTER 24 HOURS;
-- After a large historical backfill that rewrote old partitions,
-- force a full re-materialization so old grains are correct
ALTER TABLE lakehouse.sales.orders REFRESH REFLECTIONS;
Incremental refresh math (per hour), aggregation reflection
===========================================================
last materialized snapshot : S_prev (covers order_date <= 2026-09-04)
current source snapshot : S_now (adds order_date = 2026-09-05 rows)
diff(S_prev, S_now) -> new files in the 2026-09-05 partition only
re-aggregate ONLY those rows into cube grain (region,status,2026-09-05)
append the new cube rows; existing cube rows for older dates untouched
=> refresh cost = O(rows in new partition), not O(source)
Step-by-step explanation.
-
REFRESH METHOD INCREMENTAL BY (order_date)tells Dremio the source grows by neworder_datepartitions. Each refresh diffs the current Iceberg snapshot against the last materialized one and processes only the new partition's files. - For the aggregation reflection, the engine aggregates only the new day's rows to
region × status × order_dategrain and appends those cube rows. Older cube rows are untouched because their source partitions did not change — the refresh is O(new rows), not O(source). -
REFRESH EVERY 1 HOURSsets the cadence; the dashboard is therefore at most an hour behind source.EXPIRE AFTER 24 HOURSis the safety guard: if refresh fails for a day, the reflection is marked expired and the planner falls back to the raw table rather than serving day-old numbers as if current. - Incremental refresh is exact only for append-mostly or additive changes. If a backfill updates old partitions (restating last month's revenue), the incremental diff may miss the semantic change, so you force
REFRESH REFLECTIONSto fully re-materialize the affected grains. Knowing when incremental is unsafe is the senior distinction. - The net effect is warehouse-like freshness without a load window: hourly, minutes-long refreshes keep a tiny cube current over a billion-row source, and the expiry guard converts a refresh outage into a correctness-preserving slowdown rather than a silent data bug.
Output.
| Refresh | Snapshot processed | Cube rows touched | Cost |
|---|---|---|---|
| hourly (new day) | diff since last | one date's grains | O(new partition) |
| backfill (old dates rewritten) | full re-materialize | affected grains | O(affected) |
| refresh failed > 24 h | none | reflection expired | raw fallback |
Rule of thumb. Prefer incremental refresh anchored on the partition/snapshot the source actually grows by, always pair it with an EXPIRE AFTER guard, and force a full refresh after any backfill that rewrites historical partitions. Incremental keeps the reflection cheap; the expiry guard keeps a stale reflection from becoming a correctness incident.
Data engineering interview question on reflections
A senior interviewer might ask: "You have a 4-billion-row Iceberg orders table. Finance runs a revenue-by-region-by-day dashboard that must be sub-second, and a fraud team runs ad-hoc filter-and-join queries on recent orders. Design the reflections you'd create, justify raw vs aggregation for each consumer, and specify the refresh strategy so neither reflection serves stale numbers."
Solution Using a raw reflection for ad-hoc scans and an aggregation reflection for the dashboard
-- 1. Aggregation reflection for the finance dashboard (roll-up shape)
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION fin_cube
USING
DIMENSIONS (region, status, order_date)
MEASURES (amount (SUM, COUNT, MIN, MAX));
-- 2. Raw reflection for the fraud team's filter + join queries
ALTER TABLE lakehouse.sales.orders
CREATE RAW REFLECTION fraud_raw
USING DISPLAY (order_id, customer_id, order_ts, region, status, amount)
PARTITION BY (region)
LOCALSORT BY (order_ts)
DISTRIBUTE BY (customer_id);
-- 3. One refresh policy: incremental hourly, expire after a day
ALTER TABLE lakehouse.sales.orders
SET ACCELERATION REFRESH POLICY
REFRESH METHOD INCREMENTAL BY (order_date)
REFRESH EVERY 1 HOURS
EXPIRE AFTER 24 HOURS;
Consumer -> reflection routing (chosen by the planner, not the user)
===================================================================
finance dashboard (GROUP BY region, order_date) -> fin_cube (tiny)
finance roll-up (GROUP BY region) -> fin_cube (coarser)
fraud filter+join (region + recent time + join) -> fraud_raw (sorted)
fraud full-detail (project columns not in raw) -> raw orders (fallback)
Step-by-step trace.
| Consumer | Query shape | Reflection chosen | Why |
|---|---|---|---|
| Finance | group by region, day | fin_cube | grain matches / coarser |
| Finance | group by region only | fin_cube | additive roll-up |
| Fraud | filter region + time, join | fraud_raw | partition + sort + distribute |
| Fraud | project non-DISPLAY column | raw orders | column absent from reflection |
| Both | after hourly refresh | fresh materializations | incremental snapshot diff |
After deployment, the finance dashboard reads a few thousand rows from fin_cube and returns sub-second; the fraud team's recent-orders filter-and-join reads the sorted, distributed fraud_raw and skips both irrelevant partitions and the join shuffle; incremental hourly refresh keeps both within an hour of source; and any query the reflections cannot serve (or a reflection that has expired) falls back to the raw Iceberg table — correct, just slower.
Output:
| Metric | fin_cube | fraud_raw |
|---|---|---|
| Type | aggregation | raw |
| Footprint | ~thousands of rows | column subset of source |
| Serves | roll-ups at/above grain | filter + join on selected cols |
| Refresh | incremental hourly | incremental hourly |
| Fallback | raw table | raw table |
Why this works — concept by concept:
-
Aggregation reflection (fin_cube) — collapses billions of rows into a
region × status × order_datecube; the dashboard's group-bys match its grain or roll up coarser, so a scan becomes a tiny read. StoringSUM+COUNTletsAVGderive for free. -
Raw reflection (fraud_raw) — a physically-optimized column subset: PARTITION prunes by region, LOCALSORT clusters by time for range filters, DISTRIBUTE co-locates by
customer_idso the join skips a shuffle — the three levers that make ad-hoc filter-and-join interactive. -
One shared refresh policy — both reflections refresh incrementally by
order_date, so the maintenance cost is O(new partition) per hour, not O(source), and a single policy governs both consumers. -
Expiry guard + fallback —
EXPIRE AFTER 24 HOURSturns a refresh outage into a correctness-preserving fallback to the raw table instead of silently serving stale revenue. - Cost — two reflections sized to their consumers (a tiny cube plus a column-subset raw), refreshed incrementally, versus copying the whole dataset into a warehouse. Dashboard cost drops to O(groups), fraud queries to O(pruned partition), and correctness never degrades — worst case is a raw scan. That is the raw-plus-aggregation pairing senior interviews look for.
Optimization
Topic — optimization
Materialization and pre-aggregation problems
3. How the optimizer transparently substitutes reflections
The planner rewrites your query to read a matching reflection — you never name it, cost decides it
The mental model in one line: reflection substitution is the planner step where Dremio takes an incoming query against a logical dataset, searches the set of available reflections for one whose definition subsumes the query (same or superset of columns, same or finer-than-needed grain, compatible filters and joins), rewrites the logical plan to read the reflection instead of the raw source, costs both the reflected and unreflected plans, and picks the cheapest — all without the query ever referencing the reflection by name. This transparency is the whole point: application code and BI tools issue plain SQL against curated datasets, and acceleration is applied by the optimizer, so you can add, drop, or re-shape reflections without touching a single query.
The substitution pipeline — match, rewrite, cost, choose.
- Match. For each candidate reflection, the planner checks algebraic compatibility: does the reflection contain the columns the query projects and filters on? Is its grain equal to or finer than the query's group-by? Do its joins cover the query's joins?
- Rewrite. If compatible, the planner substitutes the reflection's scan for the raw dataset's scan and adds any residual operations — an extra roll-up if the reflection is finer-grained, an extra filter if the reflection is a superset.
- Cost. Both the reflected plan and the raw plan are costed using statistics. Substitution is not automatic acceptance — a reflection is only used if the reflected plan is cheaper.
- Choose. The cheapest plan wins. The query profile records reflections considered, matched, and chosen, which is your debugging surface.
Subsumption — when a reflection can answer a query.
- Column subsumption. A raw reflection can serve a query only if it contains every column the query reads. A superset is fine (residual projection drops extras); a missing column disqualifies it.
- Grain subsumption. An aggregation reflection can serve a query grouped at its grain or coarser (roll up), never finer. Coarser is a residual re-aggregation; finer is impossible.
-
Filter subsumption. A reflection filtered to a subset (e.g. only
status='paid') can serve only queries whose predicates fall inside that subset. An unfiltered reflection subsumes filtered queries via a residual filter. - Join subsumption. A reflection that materializes a join can serve queries over that join; a reflection on a single table can serve the un-joined portion and let the planner join the rest.
Why substitution beats named materialized views.
- Decoupled from queries. With classic materialized views, a query must reference the view to benefit. With substitution, the query references the logical dataset and the optimizer picks the acceleration — so you tune reflections without rewriting queries or BI dashboards.
- Multiple candidates. Several reflections may match one query; the cost model picks the best. You can layer a broad raw reflection and a narrow aggregation reflection and let the planner route each query to the cheaper one.
- Graceful fallback. If no reflection subsumes the query or all are expired, the plan reads the raw source. Acceleration is an optimization, never a correctness dependency.
Common interview probes on substitution.
- "Do users reference reflections?" — required answer: no; substitution is transparent and cost-based.
- "When is a reflection considered but not chosen?" — matched algebraically but the raw plan costed cheaper, or the reflection is expired.
- "Can an aggregation reflection serve a finer-grained query?" — no; only equal or coarser grain.
- "How do you debug a reflection that isn't used?" — read the query profile's considered/matched/chosen and check columns, grain, filters, and freshness.
Worked example — reading a plan that substitutes a reflection
Detailed explanation. The way to see substitution is to EXPLAIN a dashboard query and read where the scan comes from. Without a matching reflection, the plan scans raw orders; with orders_cube present, the plan reads the reflection and does a small residual aggregation. Walk the two plans side by side.
- The query. Revenue by region by day — the dashboard shape.
-
Without reflection. Scan raw
orders(billions of rows) → aggregate. -
With
orders_cube. Scan the cube (thousands of rows) → residual roll-up.
Question. Show the logical difference between the unaccelerated and accelerated plans for the dashboard query.
Input.
| Plan | Leaf scan | Rows scanned | Extra step |
|---|---|---|---|
| unaccelerated | raw orders | ~4e9 | full aggregation |
| accelerated | orders_cube | ~thousands | residual roll-up |
Code.
-- Inspect the plan for the dashboard query
EXPLAIN PLAN FOR
SELECT region, order_date, SUM(amount) AS revenue
FROM lakehouse.sales.orders
GROUP BY region, order_date;
-- Unaccelerated plan (no matching reflection)
Screen
Project(region, order_date, revenue)
HashAgg(group=[region, order_date], SUM(amount))
TableScan(table=[lakehouse.sales.orders], rows=4,000,000,000)
-- Accelerated plan (orders_cube subsumes the query)
Screen
Project(region, order_date, revenue)
HashAgg(group=[region, order_date], SUM(sum_amount)) <- residual roll-up
TableScan(table=[__accelerator..orders_cube], rows=6,800) <- reflection
[reflections: considered=2, matched=1, chosen=orders_cube]
Step-by-step explanation.
- In the unaccelerated plan, the leaf
TableScanreads the rawordersIceberg table — four billion rows — and aHashAggcollapses them to the grouped result. Every dashboard refresh pays the full scan. - In the accelerated plan, the leaf scan is
__accelerator...orders_cube— the reflection's Iceberg table — reading only a few thousand pre-aggregated rows. The planner recognized that the cube's grain (region × status × order_date) is finer than the query's group-by (region × order_date). - Because the cube is finer-grained, the plan keeps a residual
HashAggthat rolls the cube's per-status rows up toregion × order_date, summing the storedsum_amount. This is grain subsumption in action: coarser query, residual re-aggregation. - The annotation
considered=2, matched=1, chosen=orders_cubeis the substitution audit trail. Two reflections were considered, one matched algebraically, and the cost model chose it because scanning 6,800 rows is far cheaper than scanning four billion. - The query text is identical in both cases. The user did not change anything; adding the reflection changed the plan. That is the transparent-acceleration contract — you optimize by managing reflections, not by rewriting SQL.
Output.
| Metric | Unaccelerated | Accelerated |
|---|---|---|
| Leaf scan | raw orders | orders_cube |
| Rows scanned | ~4e9 | ~6,800 |
| Residual op | full aggregation | roll-up of sums |
| Query text changed? | — | no |
Rule of thumb. Always EXPLAIN a query you expect to be accelerated and confirm the leaf scan names a reflection and the profile shows chosen=. If the leaf still names the raw table, the reflection did not win — and the plan, not the docs, tells you why.
Worked example — subsumption rules with a filtered reflection
Detailed explanation. Subsumption is not "same query"; it is "the reflection's data is a superset of what the query needs." A reflection filtered to status='paid' can only serve queries confined to paid orders; an unfiltered reflection can serve both via a residual filter. Walk three queries against two reflections to make the rule concrete.
- Reflection A. Unfiltered aggregation cube over all statuses.
-
Reflection B. Aggregation cube filtered to
status='paid'(smaller). -
The rule. B serves only queries whose predicate is inside
status='paid'; A serves anything at its grain or coarser.
Question. For three dashboard queries, decide which reflection(s) subsume each and what residual operation the planner adds.
Input.
| Query | Predicate | Group by |
|---|---|---|
| Q1 | (none) | region, order_date |
| Q2 | status='paid' | region, order_date |
| Q3 | status='refunded' | region |
Code.
-- Reflection A: all statuses
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION cube_all
USING DIMENSIONS (region, status, order_date) MEASURES (amount (SUM, COUNT));
-- Reflection B: paid only (a filtered, smaller cube)
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION cube_paid
USING DIMENSIONS (region, status, order_date) MEASURES (amount (SUM, COUNT));
-- (filter applied via the reflection's defining view: WHERE status='paid')
Subsumption decision table
==========================
Q1 no filter, group region+date -> cube_all (residual: roll up over status)
-> cube_paid NOT valid (missing non-paid rows)
Q2 status=paid, group region+date -> cube_paid (exact-ish; smallest scan)
-> cube_all also valid (residual filter status=paid)
-> planner costs both, picks cube_paid (fewer rows)
Q3 status=refunded, group region -> cube_all (residual filter + roll up to region)
-> cube_paid NOT valid (refunded not in it)
Step-by-step explanation.
- Q1 has no status filter, so it needs every status.
cube_paidis missing refunded/pending rows and therefore cannot subsume Q1 — using it would return wrong totals. Onlycube_allqualifies, with a residual roll-up that sums away thestatusdimension. - Q2 is confined to paid orders. Both cubes can serve it:
cube_paiddirectly, andcube_allvia a residualstatus='paid'filter. Both are algebraically valid, so the cost model decides —cube_paidscans fewer rows, so it wins. - Q3 wants refunded orders rolled up to region.
cube_paidobviously cannot help (no refunded rows).cube_allsubsumes it with two residual steps: astatus='refunded'filter, then a roll-up toregion. - The disqualifier in every case is missing rows: a filtered reflection can never serve a query that needs rows outside its filter. The enabler is superset: an unfiltered reflection serves filtered queries by adding the filter back as a residual.
- This is why a broad, unfiltered reflection is a safe general accelerator and a narrow, filtered one is a targeted optimization for a hot predicate — the planner will use the narrow one when it applies and fall back to the broad one otherwise.
Output.
| Query | Valid reflections | Chosen | Residual added |
|---|---|---|---|
| Q1 | cube_all | cube_all | roll up over status |
| Q2 | cube_paid, cube_all | cube_paid | (none / minimal) |
| Q3 | cube_all | cube_all | filter refunded + roll up |
Rule of thumb. A filtered reflection subsumes only queries whose predicates live inside its filter; an unfiltered reflection subsumes filtered queries via a residual filter. Keep one broad reflection as the safety net and add narrow filtered reflections only for genuinely hot predicates — and let the cost model choose between them.
Worked example — debugging a reflection that isn't used
Detailed explanation. The most common Dremio support question is "I made a reflection but my query still scans the raw table." Substitution can fail at any of four gates: the reflection is expired, it lacks a needed column, its grain is too fine-or-coarse, or its filter excludes needed rows. Walk the diagnostic order using the query profile.
- Gate 1 — freshness. Expired reflections are skipped entirely; check the reflection's status and last refresh.
- Gate 2 — columns. A projected/filtered column absent from a raw reflection disqualifies it.
- Gate 3 — grain. An aggregation reflection coarser than the query cannot serve it.
- Gate 4 — cost. Matched but the raw plan costed cheaper (rare, but happens for tiny sources).
Question. Given a query that unexpectedly scans raw orders, produce the diagnostic checklist that finds the failing gate.
Input.
| Gate | Symptom in profile | Fix |
|---|---|---|
| freshness | reflection status = expired/failed | fix refresh; force REFRESH REFLECTIONS |
| columns | considered but not matched | add column to DISPLAY |
| grain | not matched (agg too coarse) | add a finer reflection |
| cost | matched but not chosen | usually fine; source is small |
Code.
-- 1. Confirm the reflection exists, is enabled, and is fresh
SELECT reflection_name, type, status, last_refresh, expiration
FROM SYS.reflections
WHERE dataset_name = 'lakehouse.sales.orders';
-- 2. Re-run EXPLAIN and read the reflections annotation
EXPLAIN PLAN FOR
SELECT customer_id, SUM(amount)
FROM lakehouse.sales.orders
GROUP BY customer_id;
-- [reflections: considered=1, matched=0] <- matched=0 means a gate failed
-- 3. If columns/grain are wrong, re-shape the reflection
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION cube_by_customer
USING DIMENSIONS (customer_id, order_date) MEASURES (amount (SUM, COUNT));
-- 4. If it was expired, force a refresh
ALTER TABLE lakehouse.sales.orders REFRESH REFLECTIONS;
Step-by-step explanation.
- Start at freshness: query
SYS.reflectionsforstatusandexpiration. An expired or failed reflection is invisible to substitution regardless of how well it matches — this is the most common and most overlooked cause. - If the reflection is healthy, read the
EXPLAINannotation.matched=0means no reflection was even algebraically compatible — a columns or grain problem.matched>=1, chosen=nonemeans it matched but lost on cost. - In the example, the query groups by
customer_id, but the existing cube's finest dimension set does not includecustomer_id. Grain subsumption fails (matched=0) because you cannot recover a dimension the cube dropped. The fix is a new reflection that carriescustomer_idas a dimension. - If the profile showed the reflection matched but not chosen, the raw plan was cheaper — usually because the source is small enough that scanning it beats the reflection's overhead. That is the optimizer working correctly; no action needed.
- Work the gates in order — freshness, columns, grain, cost — and the profile tells you exactly which one failed. Guessing (dropping and recreating reflections at random) is the anti-pattern; the profile is deterministic.
Output.
| Observation | Diagnosis | Action |
|---|---|---|
| status=expired | freshness gate | fix/force refresh |
| matched=0, missing column | column gate | add to DISPLAY |
| matched=0, wrong grain | grain gate | add dimension / finer reflection |
| matched≥1, chosen=none | cost gate | none (raw is cheaper) |
Rule of thumb. Debug substitution in fixed order — freshness, columns, grain, cost — using SYS.reflections and the EXPLAIN reflections annotation. The plan and the system table are authoritative; never diagnose "reflection not used" by intuition when the profile names the exact failing gate.
Data engineering interview question on reflection substitution
A senior interviewer might ask: "A BI dashboard against your Iceberg orders table is slow even though you created an aggregation reflection matching its group-by. The team insists the reflection 'isn't working.' Walk me through how substitution actually decides to use a reflection, how you'd prove from the query plan whether it's being used, and the ordered checklist you'd run to find why it isn't."
Solution Using EXPLAIN, the query profile, and the four subsumption gates
-- 1. Prove whether the reflection is chosen (plan is authoritative)
EXPLAIN PLAN FOR
SELECT region, order_date, SUM(amount) AS revenue
FROM lakehouse.sales.orders
GROUP BY region, order_date;
-- Read: leaf scan name + [reflections: considered/matched/chosen]
-- 2. Check freshness and shape in the system table
SELECT reflection_name, type, status, last_refresh, expiration
FROM SYS.reflections
WHERE dataset_name = 'lakehouse.sales.orders';
-- 3. Fix by gate:
-- columns/grain -> reshape the reflection to subsume the query
ALTER TABLE lakehouse.sales.orders
CREATE AGGREGATE REFLECTION dash_cube
USING DIMENSIONS (region, status, order_date) MEASURES (amount (SUM, COUNT));
-- freshness -> force refresh
ALTER TABLE lakehouse.sales.orders REFRESH REFLECTIONS;
The four gates, in the order the planner (and you) apply them
=============================================================
1 FRESHNESS expired/failed reflection -> skipped entirely
2 COLUMNS query projects/filters a column the reflection lacks -> no match
3 GRAIN agg reflection coarser than the query's group-by -> no match
4 COST matched, but raw plan costed cheaper -> not chosen
profile line: [reflections: considered=N, matched=M, chosen=X]
Step-by-step trace.
| Step | What you inspect | Signal | Conclusion |
|---|---|---|---|
| 1 | EXPLAIN leaf scan | names raw orders | reflection not chosen |
| 2 | profile annotation | matched=0 | algebraic mismatch (cols/grain) |
| 3 | SYS.reflections.status | expired | freshness gate failed |
| 4 | reshape + refresh | leaf now names dash_cube | substitution restored |
| 5 | re-EXPLAIN | chosen=dash_cube | proven accelerated |
After running the checklist, the plan's leaf scan flips from lakehouse.sales.orders to __accelerator...dash_cube, the profile reports chosen=dash_cube, and the dashboard drops from a multi-second raw scan to a sub-second cube read — with no change to the dashboard's SQL. The team's "it isn't working" resolves to a specific, named gate rather than a guess.
Output:
| Before fix | After fix |
|---|---|
| leaf scan = raw orders | leaf scan = dash_cube reflection |
| matched=0 (grain/cols) | matched=1, chosen=dash_cube |
| multi-second dashboard | sub-second dashboard |
| SQL unchanged | SQL unchanged |
Why this works — concept by concept:
-
EXPLAIN is authoritative — the leaf scan name and the
considered/matched/chosenannotation are ground truth for whether substitution fired; you never diagnose acceleration by latency alone. -
Freshness gate — expired or failed reflections are invisible to substitution, so
SYS.reflections.statusis the first thing to check; a perfectly-shaped reflection that failed to refresh will still be ignored. -
Column and grain subsumption — a reflection must contain every column the query touches and be at a grain equal-to-or-finer-than the group-by;
matched=0almost always means one of these failed, and the fix is reshaping, not re-running. -
Cost gate —
matchedbutchosen=nonemeans the optimizer correctly preferred the raw plan for a small source; that is healthy behavior, not a bug. - Cost — the diagnosis is O(1): two reads (EXPLAIN + one system-table query) localize the failure to one of four gates, versus the anti-pattern of blindly dropping and recreating reflections. Once the right reflection subsumes the query and is fresh, the dashboard is O(groups) instead of O(rows) — permanently, and transparently to every consumer.
Optimization
Topic — optimization
Optimizer plan-reading and rewrite problems
4. Apache Arrow & vectorized performance
Columnar in memory, vectorized on the CPU — Arrow is why an in-place lakehouse scan can feel interactive
The mental model in one line: Apache Arrow is a language-agnostic columnar in-memory format, and Dremio processes every operator over Arrow record batches so that filters, projections, and aggregations run vectorized — one tight CPU loop over a contiguous column of values using SIMD instructions — while Arrow's zero-serialization design lets data move between operators, between executor nodes, and out to clients (via Arrow Flight) without ever being copied into and out of a row format. Arrow is the reason a reflection scan or a raw Iceberg scan turns bytes into answers fast; reflections decide what to read, Arrow decides how fast you can process it once read.
Columnar in memory — why the layout matters.
-
Contiguous columns. Arrow stores each column as a contiguous buffer of values plus a validity bitmap for nulls. A
SUM(amount)reads one tight array, not scattered fields across rows. - Parquet ↔ Arrow affinity. Parquet on disk is columnar; Arrow in memory is columnar; decoding Parquet into Arrow is close to a direct transfer, with no row-to-column transpose. The lakehouse's on-disk format and the engine's in-memory format are the same shape.
- Cache and SIMD friendliness. A contiguous column streams through CPU cache and feeds SIMD lanes; a row layout interleaves unrelated fields and wastes both cache lines and vector width.
- Cheap projection. Reading three columns of a fifty-column table reads three buffers; the other columns are never touched. Column pruning is nearly free.
Vectorized execution — one loop over a batch, not a call per row.
- Record batches. Operators process Arrow record batches (thousands of rows at a time), amortizing per-row interpreter overhead across the whole batch.
- SIMD. A single CPU instruction applies to multiple column values at once (e.g. compare 8 amounts to a threshold in one op). Vector width turns into throughput.
-
Gandiva. Dremio compiles SQL expressions to native code with an LLVM-based engine (Gandiva), so
amount * 1.2 WHERE region='EU'becomes tight machine code over Arrow buffers rather than an interpreted per-row expression tree. - Branch-light filters. Vectorized filters produce selection vectors instead of branching per row, keeping the CPU pipeline full.
Zero-serialization movement — Arrow on the wire.
- Between operators. Batches pass from scan to filter to aggregate by reference; no per-stage serialize/deserialize.
- Between nodes. When a plan fragment shuffles data across executors, Arrow buffers move without a row-format conversion at each hop.
- Arrow Flight. A high-throughput RPC protocol that streams Arrow batches to clients and BI tools directly, so a large result set is transferred as columnar batches instead of row-by-row over a legacy driver.
- The cost avoided. Serialization is often the hidden majority of a query's wall-clock at high row counts; Arrow's shared format removes it end to end.
Common interview probes on Arrow.
- "Why is columnar faster for analytics?" — contiguous columns feed cache and SIMD; only referenced columns are read.
- "What does vectorized execution mean?" — process a batch per operator call with SIMD, not a function call per row.
- "What is Gandiva?" — LLVM compilation of SQL expressions to native code over Arrow buffers.
- "What is Arrow Flight for?" — zero-serialization, high-throughput transfer of Arrow batches to clients and between nodes.
Worked example — why columnar + vectorized beats row-at-a-time
Detailed explanation. The clearest way to feel Arrow's advantage is to contrast how a filtered aggregation executes over a row store versus over Arrow columnar batches. The query sums amount where region='EU'. In a row engine, each row is a function call touching every field; in a vectorized columnar engine, two columns stream through SIMD loops. Walk both.
-
Row engine. For each of N rows: materialize the row, read
region, branch, readamount, add. Per-row interpreter overhead × N. -
Vectorized engine. Load the
regioncolumn batch, SIMD-compare to'EU'producing a selection vector, apply it to theamountcolumn batch, SIMD-sum. Overhead amortized per batch.
Question. Compare the per-row work and CPU behavior of the two execution models for SELECT SUM(amount) WHERE region='EU'.
Input.
| Aspect | Row-at-a-time | Vectorized columnar (Arrow) |
|---|---|---|
| Unit of work | one row | one batch (~thousands of rows) |
| Columns touched | all fields per row | region + amount buffers only |
| CPU | branch per row, cache misses | SIMD lanes, cache-friendly |
| Expression eval | interpreted per row | Gandiva native code |
Code.
-- Row-at-a-time (conceptual)
sum = 0
for row in table: # N iterations, per-row overhead
r = deserialize(row) # touch every field
if r.region == 'EU': # branch per row
sum += r.amount
-- work ~ O(N) rows x (deserialize + branch + add)
-- Vectorized columnar over Arrow (conceptual)
sum = 0
for batch in arrow_batches: # N/4096 iterations
reg = batch.column('region') # contiguous buffer
amt = batch.column('amount') # contiguous buffer
mask = simd_eq(reg, 'EU') # 1 SIMD op per lane-group
sum += simd_masked_sum(amt, mask) # 1 SIMD reduction per lane-group
-- work ~ O(N / lanes) with amortized loop overhead, no per-row branch
Step-by-step explanation.
- The row engine pays per-row overhead N times: it deserializes each row (touching every column even though only two are used), branches on
region, and adds. The branch and the scattered field access cause CPU cache misses and pipeline stalls. - The vectorized engine reads only the
regionandamountcolumns as contiguous Arrow buffers. Because the other 48 columns are never loaded, I/O and memory traffic drop to what the query actually needs — columnar pruning is automatic. -
simd_eq(reg, 'EU')compares many values per instruction and emits a selection mask instead of branching. Vector width (say 8 lanes) turns into roughly an 8× reduction in comparison instructions, and the mask keeps the CPU pipeline branch-free. -
simd_masked_sumreduces the selectedamountvalues, again many per instruction. Gandiva compiles the whole expression (region='EU'and the masked sum) to native code, removing interpreter overhead entirely. - The combined effect — fewer columns read, batch amortization, SIMD, and native-compiled expressions — is why a well-shaped scan over Arrow is often an order of magnitude faster than a row engine on the same data, and it stacks on top of whatever a reflection already saved on rows scanned.
Output.
| Metric | Row-at-a-time | Vectorized columnar |
|---|---|---|
| Columns read | all | region, amount |
| Per-row branches | N | 0 (selection mask) |
| Instruction count | ~O(N) | ~O(N / lanes) |
| Expression eval | interpreted | Gandiva native |
Rule of thumb. Analytics is a columnar, batch, SIMD problem — not a row-loop problem. Project only the columns you need so Arrow reads only those buffers, and let vectorized operators plus Gandiva do the rest; the fewer columns and the larger the batch, the more the CPU runs at full width.
Worked example — Arrow Flight for moving large results
Detailed explanation. A dashboard or an ML client that pulls a million-row result over a legacy row-based driver spends most of its time serializing rows. Arrow Flight streams the same result as Arrow record batches — the exact format the engine already holds — so there is no server-side row conversion and no client-side re-parse. Walk the transfer.
- Legacy driver. Server converts columnar → rows → wire bytes; client parses wire bytes → rows → back to columns for analysis.
- Arrow Flight. Server streams Arrow batches as-is; client receives Arrow batches ready to use.
- The win. Zero serialization on both ends; throughput bounded by the network, not by CPU conversion.
Question. Contrast a legacy row-driver fetch with an Arrow Flight fetch for a large analytical result.
Input.
| Stage | Legacy JDBC/ODBC | Arrow Flight |
|---|---|---|
| Server encode | columnar → rows → bytes | Arrow batches as-is |
| Wire format | row-encoded | Arrow columnar |
| Client decode | bytes → rows → columns | Arrow batches ready |
| Bottleneck | CPU serialization | network bandwidth |
Code.
# Arrow Flight client — receive columnar batches, no row re-parse
import pyarrow.flight as flight
client = flight.connect("grpc+tls://dremio-coordinator:32010")
# (auth elided) obtain a FlightInfo for a SQL query
info = client.get_flight_info(
flight.FlightDescriptor.for_command(
b"SELECT region, order_date, revenue FROM sales.dash_vds"))
# Stream Arrow record batches straight into a columnar table
reader = client.do_get(info.endpoints[0].ticket)
table = reader.read_all() # already Arrow columnar; no row loop
# Hand directly to a columnar analytics library — zero conversion
df = table.to_pandas(zero_copy_only=False)
Step-by-step explanation.
- The client asks the coordinator for a
FlightInfo, which describes where and how to fetch the result. The engine already holds the result as Arrow batches, so there is nothing to convert server-side. -
do_getstreams Arrow record batches over gRPC. The bytes on the wire are the Arrow columnar format — no row encoding, no schema re-declaration per row. -
reader.read_all()assembles the batches into an ArrowTableon the client with no per-row parsing; the client receives columns, not rows it must transpose back into columns. - Handing the Arrow table to a columnar library (pandas/Polars/DuckDB) is cheap because the memory is already in the target layout — often zero-copy for numeric columns.
- Because both ends skip serialization, the transfer is limited by network bandwidth rather than CPU. For large result sets — feature extracts, big BI pulls — Flight is frequently several times faster than a legacy driver, and it composes with reflections: the reflection makes the result small and cheap to compute, Flight makes it fast to move.
Output.
| Metric | Legacy driver | Arrow Flight |
|---|---|---|
| Server-side conversion | columnar → rows | none |
| Client-side parse | rows → columns | none |
| Limiting resource | CPU | network |
| Large-result throughput | baseline | multiples faster |
Rule of thumb. For large analytical result sets, move data with Arrow Flight so neither side pays serialization; reflections shrink what you compute, Arrow Flight speeds moving it. Reach for a legacy row driver only for tiny results or tools that cannot speak Flight.
Worked example — projection and batch size as performance levers
Detailed explanation. Two knobs decide how well vectorization pays off: how many columns you project (fewer columns, fewer Arrow buffers read) and how large the record batches are (bigger batches amortize per-call overhead, up to a cache-size limit). Walk how a wide SELECT * sabotages vectorization and how batch size trades overhead for memory.
-
Projection.
SELECT *reads every column buffer;SELECT region, amountreads two. Wide projection defeats columnar pruning. - Batch size. Tiny batches pay loop/dispatch overhead per batch; huge batches blow the CPU cache. There is a sweet spot (Arrow batches are typically a few thousand rows).
Question. Show why narrowing projection and right-sizing batches improves a vectorized scan, and what breaks at the extremes.
Input.
| Lever | Bad extreme | Good setting |
|---|---|---|
| Projection | SELECT * (50 cols) | only needed columns |
| Batch size | 1 row (row-like) | ~thousands of rows |
| Batch size | 10M rows | cache thrash |
Code.
-- Anti-pattern: wide projection reads all 50 column buffers
SELECT * FROM lakehouse.sales.orders WHERE region = 'EU';
-- Good: read only the two buffers the workload needs
SELECT region, amount FROM lakehouse.sales.orders WHERE region = 'EU';
Vectorized cost model (intuition)
=================================
cost ~= (bytes_read: columns_projected) # projection lever
+ (batch_overhead: rows / batch_size) # too-small batches hurt
+ (cache_pressure: batch_size x row_width) # too-large batches hurt
sweet spot: project few columns; batch a few thousand rows so a batch's
working set fits in CPU cache while amortizing per-call overhead
Step-by-step explanation.
-
SELECT *forces Arrow to materialize all 50 column buffers even though the filter and result use two. Bytes read, memory bandwidth, and cache pressure all scale with columns projected — so wide projection is the most common self-inflicted slowdown. - Narrowing to
SELECT region, amountreads exactly two buffers. This is the columnar dividend: projection is a first-class performance lever, not a cosmetic one, and it compounds with reflection column subsets. - Batch size too small (near one row) reverts the engine toward row-at-a-time: the fixed per-batch dispatch and loop-setup overhead is paid too many times, and SIMD lanes are underfilled.
- Batch size too large blows the CPU cache — a batch's working set (batch_size × row width) must fit in cache to keep the SIMD loops hot. Beyond that, you thrash and lose the vectorization win.
- The sweet spot is few columns and batches of a few thousand rows, which is why Arrow's default batch sizing plus disciplined projection gives most of the available speed. These levers sit underneath reflections: a reflection reduces rows and columns, and good projection/batching extracts full throughput from whatever it hands the executor.
Output.
| Setting | Effect |
|---|---|
| SELECT * | reads all buffers; cache pressure high |
| narrow projection | reads only needed buffers |
| batch = 1 row | per-batch overhead dominates |
| batch = few thousand | overhead amortized, cache-hot |
| batch = 10M rows | cache thrash, throughput drops |
Rule of thumb. Never SELECT * on a wide analytical table, and let the engine batch at a few thousand rows rather than forcing tiny or gigantic batches. Projection is the cheapest vectorization lever you control; combined with a column-subset reflection, it makes the executor read the least data at full CPU width.
Data engineering interview question on Apache Arrow
A senior interviewer might ask: "Explain why an in-place lakehouse engine like Dremio can process a filtered aggregation over Iceberg/Parquet nearly as fast as a warehouse, even though the data was never loaded into a proprietary store. Cover the in-memory format, the execution model, expression compilation, and how large results reach a BI client — and where reflections fit relative to Arrow."
Solution Using Arrow columnar batches, Gandiva-compiled vectorized operators, and Arrow Flight
End-to-end fast path for SELECT region, SUM(amount) WHERE region='EU'
======================================================================
1 STORAGE Parquet on S3 is columnar; only the region + amount column
chunks are read (projection + Iceberg predicate pruning).
2 DECODE Parquet chunks decode directly into Arrow record batches
(columnar -> columnar; no row transpose).
3 EXECUTE Vectorized operators run over batches: SIMD compare on region
-> selection mask -> SIMD masked SUM on amount. Gandiva
compiles the expression to native code (no per-row interpret).
4 MOVE Result batches stream to the BI client via Arrow Flight with
zero serialization (network-bound, not CPU-bound).
5 ACCEL If an aggregation reflection subsumes the query, step 1 reads a
few thousand pre-aggregated Arrow rows instead of billions —
Arrow makes each row fast; the reflection makes rows few.
# Client side: pull the (already small, already fast) result via Flight
import pyarrow.flight as flight
client = flight.connect("grpc+tls://dremio-coordinator:32010")
info = client.get_flight_info(flight.FlightDescriptor.for_command(
b"SELECT region, SUM(amount) FROM sales.orders "
b"WHERE region='EU' GROUP BY region"))
table = client.do_get(info.endpoints[0].ticket).read_all() # Arrow columnar
Step-by-step trace.
| Stage | Mechanism | Why it's fast |
|---|---|---|
| read | Parquet → Arrow, projected columns only | columnar affinity, no transpose |
| filter | SIMD compare → selection mask | branch-free, many values per op |
| aggregate | SIMD masked sum, Gandiva native | no per-row interpretation |
| move | Arrow Flight batches | zero serialization, network-bound |
| accelerate | reflection subsumes query | rows scanned drop to O(groups) |
After this path, the filtered aggregation reads only two column buffers, evaluates them branch-free with SIMD and Gandiva-compiled code, and streams a columnar result to the client with no serialization — and when a reflection matches, the leaf scan is a few thousand pre-aggregated rows instead of the raw billions. Arrow makes every row cheap to process; reflections make the number of rows small; together they deliver warehouse-class latency on an in-place open copy.
Output:
| Metric | Value |
|---|---|
| Columns read | 2 (region, amount) |
| Filter/aggregate | SIMD + Gandiva native code |
| Result transfer | Arrow Flight (zero-serialization) |
| With reflection | leaf scan ~thousands of rows |
| Net latency | interactive / sub-second |
Why this works — concept by concept:
- Columnar in memory (Arrow) — Parquet on disk and Arrow in memory are both columnar, so decoding is a near-direct transfer and only projected columns are read; the layout feeds CPU cache and SIMD instead of fighting them.
- Vectorized execution — operators process record batches with one SIMD loop per batch and emit selection masks instead of per-row branches, turning vector width into throughput and amortizing interpreter overhead.
- Gandiva LLVM compilation — SQL expressions become native machine code over Arrow buffers, removing per-row expression-tree interpretation entirely.
- Arrow Flight — streams results as Arrow batches with no server encode or client parse, so large transfers are bounded by the network rather than CPU serialization.
- Cost — Arrow drives per-row cost toward the hardware limit (few columns, SIMD, native code, zero-serialization movement), and reflections drive the row count down to O(groups) or a pruned partition. The two are orthogonal multipliers: fast rows × few rows = interactive queries on data that was never copied out of the lake.
Optimization
Topic — optimization
Vectorized-execution and scan-cost problems
5. Semantic layer & governance patterns
Curate once as virtual datasets, govern access and masking in one place, accelerate the same layer with reflections
The mental model in one line: the semantic layer is a curated tree of virtual datasets — saved SQL views organized into spaces and folders — that presents business-friendly, consistent, governed datasets over the physical lakehouse sources without copying any data, and because reflections attach to those virtual datasets, the very same layer that enforces row/column access, column masking, and lineage is also the layer that gets accelerated, so speed and governance are defined once, together, rather than bolted on downstream. The semantic layer is where the lakehouse stops being "files and an engine" and becomes a product data consumers can safely self-serve.
Virtual datasets — curation without copies.
- A VDS is a saved view. A virtual dataset is a named SQL definition over physical tables or other VDS — no materialized data by default, just logic. Consumers query the VDS as if it were a table.
-
Spaces and folders. VDS live in spaces (e.g.
Finance,Marketing) and folders, giving a browsable, governed catalog instead of raw source paths. - Layered modeling. Build a staging VDS over raw sources, a business VDS over staging, and a reporting VDS over business — each layer curates and constrains, and none copies data.
- Consistency. Business logic (revenue definitions, currency conversion, deduplication) lives once in a VDS, so every consumer sees the same numbers.
Governance — enforced where curation happens.
-
Access control.
GRANT SELECT ONa VDS to roles; consumers see only what they are entitled to, at the curated layer rather than on raw files. - Column masking. A masking policy (a UDF bound to a column) redacts or hashes sensitive values for unentitled roles — the same VDS returns full emails to admins and masked emails to analysts.
- Row access. A row-access policy (a boolean UDF over row columns) filters rows per role — a regional analyst sees only their region without a separate dataset.
- Lineage and documentation. The engine computes a lineage graph from VDS definitions (which physical columns feed which VDS column) and lets you attach wikis and tags, so governance includes provenance and discoverability, not just access.
Reflections on the semantic layer — fast and governed at once.
- Attach to a VDS. A reflection can be defined on a virtual dataset, so the curated business view — not just the raw table — is what gets accelerated.
- Governance is preserved. Substitution respects access and masking: a reflection accelerates the plan, but row/column policies still apply on top, so acceleration never leaks data a policy would hide.
- One definition, two benefits. Because the VDS is the unit of both governance and acceleration, you define correct business logic once and get speed and security on the same object.
- Consumers stay decoupled. BI tools query the VDS; the reflection is chosen by substitution; policies are enforced by the engine — none of it is visible in the dashboard's SQL.
Common interview probes on the semantic layer.
- "What is a virtual dataset?" — a saved SQL view; curation and business logic without copying data.
- "How do you enforce column-level security?" — a masking policy (UDF) bound to the column, applied per role.
- "Where does governance live in a lakehouse?" — at the semantic layer (VDS + policies), not scattered across consumers.
- "Can you accelerate a governed view?" — yes; attach a reflection to the VDS; substitution respects the policies.
Worked example — a governed VDS with column masking
Detailed explanation. Finance needs a customer orders view where analysts see masked emails and only admins see raw emails, all from one virtual dataset. Build the VDS, define a masking policy, bind it to the email column, and grant access — governance defined once at the curated layer.
-
VDS.
Finance.customer_ordersjoining orders to customers with business columns. -
Masking policy. A UDF that returns the raw email for the
adminrole and a masked form otherwise. -
Grant.
analystandadminroles both getSELECT; the mask differentiates what they see.
Question. Create a governed VDS that masks email for non-admins while serving the same view to all analysts.
Input.
| Object | Purpose |
|---|---|
Finance.customer_orders |
curated VDS (orders ⋈ customers) |
governance.mask_email |
masking UDF (role-aware) |
| masking binding | email column → mask_email |
| grants | SELECT to analyst + admin roles |
Code.
-- 1. The curated virtual dataset (no data copied)
CREATE VDS Finance.customer_orders AS
SELECT o.order_id, o.order_date, o.region, o.amount,
c.customer_id, c.segment, c.email
FROM lakehouse.sales.orders o
JOIN lakehouse.crm.customers c ON c.customer_id = o.customer_id;
-- 2. A role-aware masking policy (UDF)
CREATE FUNCTION governance.mask_email (email VARCHAR)
RETURNS VARCHAR
RETURN CASE
WHEN is_member('admin') THEN email
ELSE REGEXP_REPLACE(email, '^[^@]+', '****')
END;
-- 3. Bind the policy to the email column of the VDS
ALTER VDS Finance.customer_orders
MODIFY COLUMN email SET MASKING POLICY governance.mask_email(email);
-- 4. Grant access to both roles; the mask differentiates the view
GRANT SELECT ON Finance.customer_orders TO ROLE analyst;
GRANT SELECT ON Finance.customer_orders TO ROLE admin;
Step-by-step explanation.
-
CREATE VDS Finance.customer_ordersdefines the business view as pure SQL over the physicalordersandcustomerstables. No data is copied; the VDS is logic plus a place in the governed catalog. -
governance.mask_emailis a masking UDF that inspects the caller's role withis_member('admin'). Admins get the raw email; everyone else gets the local part replaced with****. The policy is defined once, independent of any consumer. -
ALTER VDS ... MODIFY COLUMN email SET MASKING POLICYbinds the UDF to theemailcolumn of the VDS. From now on, any query that readsemailthrough this VDS has the policy applied by the engine — there is no way for a consumer to bypass it. - Both roles get
SELECTon the same VDS. An analyst's dashboard and an admin's dashboard issue identical SQL; the engine returns masked or raw emails based on role. One curated object, correct behavior for both. - Governance is centralized at the semantic layer: the masking lives on the VDS, not duplicated in every dashboard or enforced by convention. Change the policy once and every consumer updates — which is exactly the property you want when audit or privacy rules change.
Output.
| Role | Query on customer_orders | email returned |
|---|---|---|
| admin | SELECT email ... | raw (alice@corp.com) |
| analyst | SELECT email ... | masked (****@corp.com) |
| unauthorized | SELECT ... | denied (no grant) |
Rule of thumb. Put business logic and data protection on the same virtual dataset: define the VDS once, bind masking/row policies to its columns, and grant on the VDS. Never enforce masking in individual dashboards — centralize it at the semantic layer so one change propagates to every consumer.
Worked example — a row-access policy for regional scoping
Detailed explanation. A regional analyst should see only their region's rows from the same customer_orders VDS, without maintaining a separate per-region dataset. A row-access policy — a boolean UDF over the row's region — filters rows by the caller's entitlements. Build it and reason about how it composes with reflections.
-
Policy. A boolean UDF returning true when the caller is entitled to the row's
region. - Binding. Attach the row-access policy to the VDS.
- Composition. Substitution still applies a matching reflection; the row policy filters on top, so acceleration never returns rows the policy hides.
Question. Add a row-access policy so each analyst sees only their entitled regions from the shared VDS.
Input.
| Component | Value |
|---|---|
| Policy UDF | governance.region_visible(region) |
| Entitlement source | role membership (e.g. region_EU) |
| Binding | row-access policy on the VDS |
| Interaction | applied on top of any reflection |
Code.
-- 1. A boolean row-access policy: is the caller entitled to this region?
CREATE FUNCTION governance.region_visible (region VARCHAR)
RETURNS BOOLEAN
RETURN is_member('admin')
OR is_member('region_' || region); -- e.g. role region_EU sees EU rows
-- 2. Attach the row-access policy to the curated VDS
ALTER VDS Finance.customer_orders
ADD ROW ACCESS POLICY governance.region_visible(region);
-- 3. Same query, different rows per caller
SELECT region, SUM(amount) AS revenue
FROM Finance.customer_orders
GROUP BY region;
Step-by-step explanation.
-
governance.region_visiblereturns true for admins (see everything) or for a caller who holds theregion_<X>role matching the row's region. It encodes entitlement as role membership, so access is managed in the identity system, not in copies of data. -
ALTER VDS ... ADD ROW ACCESS POLICYbinds the boolean UDF to the VDS. The engine now injects the policy as a filter on every query against the VDS, transparently — consumers cannot see or disable it. - The identical
GROUP BY regionquery returns different rows to different callers: an EU analyst gets only EU rows (and therefore only EU revenue), while an admin gets all regions. One VDS, per-caller row scoping, zero duplicated datasets. - Crucially, the row policy composes with substitution. If a reflection on the VDS accelerates the aggregation, the engine still applies
region_visibleon top of the accelerated plan, so the EU analyst never receives non-EU rows even though the reflection materialized all regions. Acceleration and governance are independent and both enforced. - This is the payoff of unifying curation, governance, and acceleration on the VDS: you scale to many analysts and regions with one governed, fast object instead of a proliferation of per-team copies that drift and leak.
Output.
| Caller role | Rows visible | Revenue seen |
|---|---|---|
| admin | all regions | total |
| region_EU | EU only | EU revenue |
| region_US | US only | US revenue |
| none | none | empty |
Rule of thumb. Use a row-access policy on the shared VDS to scope rows by entitlement instead of forking per-region datasets, and trust that substitution applies the policy on top of any reflection. One governed, accelerated VDS beats N copies every time — for correctness, for maintenance, and for security.
Worked example — accelerating a governed VDS with a reflection
Detailed explanation. The final pattern ties the guide together: attach an aggregation reflection to the curated, governed customer_orders VDS so the finance dashboard is both fast and policy-enforced. The reflection materializes the business view's roll-up; substitution routes the dashboard to it; masking and row policies still apply. Build it.
-
Reflection on VDS. Aggregation reflection defined on
Finance.customer_orders, not the raw table. -
What it accelerates. The dashboard's
region × order_daterevenue roll-up over the curated join. - Governance preserved. Row/column policies apply on top of the accelerated plan.
Question. Attach an aggregation reflection to the governed VDS and confirm governance still holds under acceleration.
Input.
| Component | Value |
|---|---|
| Target |
Finance.customer_orders (VDS) |
| Reflection | aggregation (region, order_date; SUM/COUNT amount) |
| Consumer | finance dashboard |
| Policies | masking (email) + row access (region) |
Code.
-- 1. Aggregation reflection ON the governed virtual dataset
ALTER VDS Finance.customer_orders
CREATE AGGREGATE REFLECTION vds_rev_cube
USING
DIMENSIONS (region, order_date)
MEASURES (amount (SUM, COUNT));
-- 2. The dashboard query against the governed VDS (unchanged SQL)
SELECT region, order_date, SUM(amount) AS revenue
FROM Finance.customer_orders
GROUP BY region, order_date;
-- 3. Confirm acceleration + governance in one plan
EXPLAIN PLAN FOR
SELECT region, order_date, SUM(amount)
FROM Finance.customer_orders
GROUP BY region, order_date;
-- leaf scan -> vds_rev_cube ; row-access filter region_visible still applied
Step-by-step explanation.
-
ALTER VDS ... CREATE AGGREGATE REFLECTIONdefines the materialization on the curated VDS. The reflection captures the roll-up of the business view — the join, the business columns — not the raw physical table, so the accelerated result is the governed result's shape. - The dashboard queries
Finance.customer_orderswith plain SQL. Substitution recognizesvds_rev_cubesubsumes the group-by and rewrites the plan to read the reflection — the same transparent acceleration, now on a semantic-layer object. - The
EXPLAINshows the leaf scan asvds_rev_cubeand theregion_visiblerow-access filter still present in the plan. The engine applies governance on top of the accelerated scan, so an EU analyst's dashboard reads the fast cube but still sees only EU rows. - Column masking similarly composes: if a query through the VDS reads
email, the masking policy applies regardless of whether a reflection served other columns. Acceleration changes where rows come from, never what a role is allowed to see. - The result is the guide's thesis in one object: define business logic once as a VDS, govern it with row/column policies once, accelerate it with a reflection once — and every consumer gets consistent, fast, secure self-serve analytics over a single open copy of data.
Output.
| Property | Result |
|---|---|
| Dashboard latency | sub-second (via vds_rev_cube) |
| Business logic | defined once in the VDS |
| Column masking | applied on top of reflection |
| Row access | applied on top of reflection |
| Data copies | one (lake) + small reflection |
Rule of thumb. Accelerate the governed VDS, not the raw table, so speed and security live on the same object and substitution composes with your policies. When the curated layer is also the accelerated and governed layer, consumers get fast, consistent, safe self-serve analytics without a single extra copy of data.
Data engineering interview question on the semantic layer
A senior interviewer might ask: "Business users want self-serve dashboards on your lakehouse, but security requires that analysts see masked PII and only their region's rows, and finance insists the numbers be consistent and fast. Design the semantic layer: how you'd curate the data, where you'd enforce masking and row access, and how you'd make it fast without letting acceleration bypass governance."
Solution Using layered virtual datasets with masking, row-access policies, and a VDS-attached reflection
-- 1. Curate: a governed business VDS over physical sources (no copy)
CREATE VDS Finance.customer_orders AS
SELECT o.order_id, o.order_date, o.region, o.amount,
c.customer_id, c.segment, c.email
FROM lakehouse.sales.orders o
JOIN lakehouse.crm.customers c ON c.customer_id = o.customer_id;
-- 2. Govern: column masking + row access, bound to the VDS
CREATE FUNCTION governance.mask_email (e VARCHAR) RETURNS VARCHAR
RETURN CASE WHEN is_member('admin') THEN e
ELSE REGEXP_REPLACE(e, '^[^@]+', '****') END;
ALTER VDS Finance.customer_orders
MODIFY COLUMN email SET MASKING POLICY governance.mask_email(email);
CREATE FUNCTION governance.region_visible (r VARCHAR) RETURNS BOOLEAN
RETURN is_member('admin') OR is_member('region_' || r);
ALTER VDS Finance.customer_orders
ADD ROW ACCESS POLICY governance.region_visible(region);
GRANT SELECT ON Finance.customer_orders TO ROLE analyst;
-- 3. Accelerate: an aggregation reflection ON the governed VDS
ALTER VDS Finance.customer_orders
CREATE AGGREGATE REFLECTION vds_rev_cube
USING DIMENSIONS (region, order_date) MEASURES (amount (SUM, COUNT));
One object, three responsibilities (defined once, enforced by the engine)
=========================================================================
CURATE Finance.customer_orders = business join + columns (no copy)
GOVERN masking(email) + row-access(region) bound to the VDS
ACCELERATE vds_rev_cube reflection on the VDS -> substitution routes here
ENFORCE plan = reflection scan + row-access filter + column mask
=> fast AND governed; acceleration never bypasses policy
Step-by-step trace.
| Layer | Mechanism | Guarantee |
|---|---|---|
| Curate | VDS over orders ⋈ customers | consistent business logic, no copy |
| Mask | masking UDF on email | PII hidden from non-admins |
| Scope | row-access UDF on region | analysts see only their region |
| Accelerate | reflection on the VDS | sub-second roll-ups |
| Compose | policies applied atop reflection | acceleration never leaks data |
After deployment, an EU analyst opens the finance dashboard: identical SQL against Finance.customer_orders is transparently served by vds_rev_cube (sub-second), the row-access policy trims the plan to EU rows, and any email column is masked — all enforced by the engine, none of it visible in the dashboard. Finance gets consistent, fast numbers; security gets masking and row scoping; the platform stores one open copy plus a small reflection.
Output:
| Consumer | Sees | Speed | Copies |
|---|---|---|---|
| EU analyst | EU rows, masked email | sub-second | 1 + reflection |
| US analyst | US rows, masked email | sub-second | 1 + reflection |
| admin | all rows, raw email | sub-second | 1 + reflection |
Why this works — concept by concept:
- Virtual datasets — curate the business join and columns once as a saved view; consumers query consistent logic with no data copied, so "the numbers" are defined in exactly one place.
-
Column masking policy — a role-aware UDF bound to the
emailcolumn redacts PII for non-admins at the curated layer, so no dashboard can bypass it. -
Row-access policy — a boolean UDF on
regionscopes rows per entitlement, replacing N per-region datasets with one governed object. - Reflection on the VDS — accelerates the governed view; substitution routes the dashboard to the cube while the engine still applies masking and row filters on top, so speed never widens the security boundary.
- Cost — one curated VDS, two policy UDFs, and one small aggregation reflection replace a sprawl of per-team copies and dashboard-side security. Storage is one lake copy plus the reflection; queries are O(groups) and sub-second; governance is centralized and audit-friendly. Curate once, govern once, accelerate once — the whole guide, composed on a single object.
Database
Topic — database
Semantic-layer and view-modeling problems
SQL
Topic — sql
SQL views, joins, and access-policy problems
Cheat sheet — Dremio reflection & lakehouse recipes
-
Which reflection when. Use a raw reflection for queries that filter, project, and join on selected columns — shape it with
USING DISPLAY (cols),PARTITION BY (high-cardinality equality filter),LOCALSORT BY (range filter),DISTRIBUTE BY (join key). Use an aggregation reflection for group-by/roll-up dashboards —DIMENSIONS (union of group-by cols)+MEASURES (col (SUM, COUNT, MIN, MAX)). KeepSUMandCOUNTtogether soAVG = SUM/COUNTderives; remember an aggregation reflection serves its grain or coarser, never finer. -
Raw reflection DDL template.
ALTER TABLE t CREATE RAW REFLECTION r USING DISPLAY (c1, c2, ...) PARTITION BY (p) LOCALSORT BY (s) DISTRIBUTE BY (j);— DISPLAY drives which projections it serves; the physical clauses drive how much it prunes and whether joins shuffle. -
Aggregation reflection DDL template.
ALTER TABLE t CREATE AGGREGATE REFLECTION a USING DIMENSIONS (d1, d2) MEASURES (m (SUM, COUNT, MIN, MAX));— dimensions set the grain, measures set what rolls up;COUNT(DISTINCT)is served via an approximate HLL sketch, not exact re-aggregation. -
Incremental refresh setup.
ALTER TABLE t SET ACCELERATION REFRESH POLICY REFRESH METHOD INCREMENTAL BY (order_date) REFRESH EVERY 1 HOURS EXPIRE AFTER 24 HOURS;— anchor incremental on the partition/snapshot the source grows by; always setEXPIRE AFTERso a refresh outage falls back to the raw table instead of serving stale numbers; force full re-materialization withALTER TABLE t REFRESH REFLECTIONSafter any backfill that rewrites old partitions. -
Substitution is transparent + cost-based. Users never name a reflection. The planner matches (columns/grain/filter/join subsumption), rewrites, costs the reflected vs raw plan, and picks the cheaper. Confirm with
EXPLAIN PLAN FOR ...— the leaf scan should name the reflection and the annotation should readchosen=<reflection>. -
Debug "reflection not used" in gate order. 1) Freshness —
SELECT status, expiration FROM SYS.reflections; expired reflections are skipped. 2) Columns — a projected/filtered column absent from a raw reflection disqualifies it. 3) Grain — an aggregation reflection coarser than the query cannot serve it. 4) Cost —matchedbutchosen=nonemeans the raw plan was cheaper (fine for small sources). - Subsumption rules. A raw reflection needs every column the query touches (superset OK). An aggregation reflection serves equal-or-coarser grain (residual roll-up), never finer. A filtered reflection serves only queries inside its filter; an unfiltered one serves filtered queries via a residual filter. A join reflection serves that join.
-
Arrow / vectorization checklist. Project only needed columns (never
SELECT *on wide tables); rely on Parquet↔Arrow columnar affinity (no row transpose); let vectorized operators + Gandiva LLVM run SIMD over record batches of a few thousand rows; move large results with Arrow Flight for zero-serialization transfer. Arrow makes each row fast; reflections make rows few — they multiply. -
Iceberg fundamentals reflections rely on. Snapshot → manifest list → manifests (per-file min/max stats) → data files; hidden partitioning + metadata pruning skip files before I/O;
AT SNAPSHOT/AT TIMESTAMPtime travel; every commit is a new snapshot, which is the anchor for incremental refresh. -
Semantic-layer governance checklist. Curate business logic once as virtual datasets (spaces/folders, no copy); bind column masking (
MODIFY COLUMN c SET MASKING POLICY udf(c)) and row access (ADD ROW ACCESS POLICY udf(col)) to the VDS;GRANT SELECTon the VDS to roles; rely on auto-computed lineage. Enforce governance at the curated layer, never in individual dashboards. - Fast + governed on one object. Attach a reflection to the governed VDS, not the raw table. Substitution accelerates the plan while the engine applies masking and row-access on top — acceleration never widens the security boundary. Curate once, govern once, accelerate once.
- One-copy economics. The lakehouse serves warehouse-class BI on a single open Iceberg copy plus targeted reflections (a tiny cube, a column-subset raw), refreshed incrementally — versus a second proprietary copy plus a nightly load window. Worst-case failure is a raw scan (correct, slower), never a wrong number.
Frequently asked questions
What is a Dremio reflection in one sentence?
A dremio reflection is a pre-computed, physically-optimized copy of a dataset — either a raw reflection (selected columns, sorted and partitioned for fast scans, filters, and joins) or an aggregation reflection (pre-computed dimensions and measures for fast roll-ups) — that Dremio persists as an Iceberg table in a reflection store and keeps in sync with its source via a refresh policy. Its purpose is query acceleration: the optimizer can answer an expensive query from the cheap materialization instead of scanning the raw table, and it does so transparently, so the query never references the reflection by name. Reflections are the mechanism that lets an open lakehouse serve interactive BI on data that was never copied into a proprietary warehouse.
Raw vs aggregation reflection — when do I use each?
Use a raw reflection when queries filter, project, and join on specific columns — it stores a column subset that you sort (LOCALSORT), partition (PARTITION BY), and distribute (DISTRIBUTE BY) to prune files and avoid join shuffles, so ad-hoc filter-and-join workloads become interactive. Use an aggregation reflection when queries group and roll up — it stores a small cube of DIMENSIONS and additive MEASURES (SUM, COUNT, MIN, MAX), so a billion-row GROUP BY becomes a few-thousand-row read. The decisive rule is grain: an aggregation reflection serves any query at its grain or coarser (it rolls up), but never one finer (a dropped dimension is gone). Many deployments run both — an aggregation reflection for dashboards and a raw reflection for ad-hoc detail — and let cost-based substitution route each query to the cheaper one.
How does Dremio decide to use a reflection?
Through cost-based substitution, entirely transparently. When a query arrives, the planner checks each available reflection for algebraic subsumption — does it contain every column the query needs, is its grain equal-to-or-finer-than the group-by, do its filters and joins cover the query's — then rewrites the plan to read a matching reflection (adding residual roll-ups or filters as needed), costs the reflected plan against the raw plan using table statistics, and picks the cheaper. The query text never mentions the reflection; you tune acceleration by managing reflections, not by rewriting SQL or dashboards. You can verify what happened with EXPLAIN PLAN FOR ...: the leaf scan should name the reflection and the profile annotation reports reflections considered, matched, and chosen.
Do reflections copy my data, and where do they live?
Yes — a reflection is a real materialization, not a hint, but it is a small, purpose-built copy rather than a duplicate of the whole dataset. Reflections are written as Iceberg tables (Parquet under the hood) to a configured reflection-store location on your data lake, so they are durable, open, columnar files that survive engine restarts — not an in-memory cache (memory caching of reflection files is a separate speed layer). An aggregation reflection is usually tiny relative to the source because it stores only grouped rows; a raw reflection can approach source size if it keeps most columns. Because they are Iceberg tables, reflections inherit snapshotting and metadata pruning, which is exactly what makes them cheap to refresh incrementally and fast to substitute.
How does Apache Arrow make Dremio fast?
Apache Arrow is a columnar in-memory format, and processing data column-by-column instead of row-by-row unlocks three multipliers. First, contiguous column buffers feed the CPU cache and SIMD instructions, so a filter or sum runs vectorized — many values per instruction — and only the columns a query references are read at all. Second, Dremio compiles SQL expressions to native machine code with Gandiva (LLVM), removing per-row interpretation. Third, because Parquet on disk, Arrow in memory, and Arrow Flight on the wire are all the same columnar shape, data moves between operators, between executor nodes, and out to BI clients with near-zero serialization. Arrow makes each row cheap to process; reflections make the number of rows small — together they deliver warehouse-class latency on an in-place lakehouse copy.
What is the Dremio semantic layer and why does it matter for governance?
The semantic layer is a curated tree of virtual datasets — saved SQL views organized into spaces and folders — that present consistent, business-friendly datasets over physical lakehouse sources without copying data. It matters for governance because it is the single place to enforce security: bind column masking policies (role-aware UDFs that redact PII) and row-access policies (boolean UDFs that scope rows by entitlement) to a virtual dataset, GRANT SELECT on the VDS to roles, and rely on the auto-computed lineage graph for provenance — so one analyst sees masked emails and only their region while another sees everything, all from the same object. Crucially, you can attach a reflection to a governed VDS: substitution accelerates the plan while the engine still applies masking and row filters on top, so speed and security are defined once, together, and acceleration never bypasses governance.
Practice on PipeCode
- Drill the query optimization practice library → for the reflection-shaping, substitution, cost-model, and vectorized-scan problems senior lakehouse interviewers love.
- Sharpen your data modeling on the database practice library → for lakehouse layout, Iceberg table design, cube design, and semantic-layer modeling.
- Rehearse the fundamentals on the SQL practice library → for the aggregation, GROUP BY, join, and view/access-policy problems that underpin every reflection and virtual dataset.
- Stack the prerequisites against PipeCode's broader 450+ data-engineering catalogue to anchor the raw-vs-aggregation reflection decision and the substitution gates against real graded inputs.
Lock in Dremio reflection muscle memory
Docs explain reflections. PipeCode drills explain the decision — when a raw reflection beats an aggregation one, why the optimizer substitutes a materialization the query never named, how Apache Arrow turns an in-place scan interactive, and how a governed virtual dataset stays fast without leaking data. Pipecode.ai is Leetcode for Data Engineering — pattern-first practice tuned for the production trade-offs senior data engineers actually face.
Practice optimization problems →
Practice database problems →





Top comments (0)