Query Optimizer
The query-planning layer that turns a semantic-layer request or ad hoc SQL query into an efficient execution plan against the lakehouse.
High-Level Design
The Query Optimizer is why a well-formed question doesn't require the asker to know partition strategy.
💼 Business Context
- The difference between a query that returns instantly and one that times out is usually query planning, not raw compute — this is where that gap is closed
- Keeps query compute costs proportional to the question asked, not the total size of the lakehouse
- Owned by Analytics Engineering / Platform Engineering
🔌 Technical Overview
Queries compiled from the Semantic Layer or submitted ad hoc are planned by the Apache Arrow DataFusion engine running inside the Analytics/AI API (a Docker container on AKS), which pushes filter predicates down to Iceberg's partition metadata (from Partitioning & Clustering) to prune irrelevant files before any data is read, reorders joins based on table statistics, and can route especially large ad hoc analytical workloads to Snowflake or Spark via the same Iceberg tables when a query exceeds DataFusion's efficient working range.
Optimizations
💾 Query Plan (abridged)
EXPLAIN SELECT net_revenue FROM marts.order_fact WHERE event_date >= '2026-07-26'; -- Plan: Partition prune (7 of 400+ date partitions read) -- -> Projection push-down (order_total, refund_amount only) -- -> Aggregate (sum)
🔗 Integration Points
- Partitioning & Clustering (Unified Data Foundation) — the physical layout the optimizer prunes against
- Semantic Layer — the compiled SQL the optimizer plans and executes
- Snowflake / Spark — offload targets for workloads beyond DataFusion's efficient range, via the same Iceberg tables
- Caching — the layer the optimizer checks before planning a fresh execution
🧰 Services Consumed
- Owning microservice —
Cxos.Intelligence.Api(see the Full Application Service Map) - Database — Azure Database for PostgreSQL (semantic layer) + Azure Cache for Redis (query cache)
⚠️ Non-Functional Considerations
- Scale: partition pruning keeps query cost roughly constant as total lakehouse size grows, since only relevant partitions are ever touched
- Latency: p95 interactive dashboard queries return in under 2 seconds against multi-billion-row fact tables
- Reliability: a query plan that would scan an unreasonable number of partitions is flagged and rejected before execution rather than silently running for hours
- Security/Privacy: the optimizer plans within the row/column boundaries Access Control already applied — it never exposes data the caller lacks entitlement to see
🎯 Enterprise Example
A "last 7 days of revenue by region" dashboard query against a 3-billion-row order_fact table returns in under a second because the optimizer prunes to 7 of 400+ date partitions before ever touching storage, rather than scanning the full table history.