The Query Transformation Pipeline
Hacker News

The Query Transformation Pipeline

How Readyset Rewrites Your SQL: Inside the Query Transformation Pipeline Why Query Rewriting Matters for Dataflow Most databases work the same way: a query arrives, the engine builds an execution plan, scans tables, joins rows, filters, aggregates, and returns the result. Every time the query runs, the work repeats from scratch. This pull-based model has served relational databases for decades, but it carries an inherent cost: read latency is proportional to the complexity of the query and the size of the data it touches. Readyset takes a fundamentally different ap Vassili Zarouba 2026-04-28 · 17 min read Why Query Rewriting Matters for Dataflow Most databases work the same way: a query arrives, the engine builds an execution plan, scans tables, joins rows, filters, aggregates, and returns the result. Every time the query runs, the work repeats from scratch. This pull-based model has served relational databases for decades, but it carries an inherent cost: read latency is proportional to the complexity of the query and the size of the data it touches. Readyset takes a fundamentally different approach. Instead of re-executing queries on demand, Readyset compiles each query into a dataflow graph - a network of operators (joins, filters, aggregations, projections) that continuously maintains the query's result as the underlying data changes. When a row is inserted, updated, or deleted in the upstream database, the change propagates through the graph, and the cached result is incrementally updated. Reads become lookups into a pre-computed materialized view, not full query executions. This is the key difference: traditional engines optimize how to execute a query each time it runs. Readyset optimizes once, at cache-creation time, and then maintains the result incrementally forever. The tradeoff is that the query must be expressed in a form the dataflow engine can compile - and that form is more restrictive than what SQL allows. What the Dataflow Engine Requires SQL is a declarative language. The same logical query can be written in many equivalent ways: correlated subqueries, derived tables, CTEs , LATERAL joins, nested aggregations. A traditional optimizer treats these as interchangeable representations and picks the best execution plan regardless of syntax. Readyset's dataflow compiler is not a traditional optimizer. It translates SQL into a directed acyclic graph of streaming operators, where each operator receives changes from its inputs and emits changes to its outputs. This architecture imposes structural constraints that SQL syntax doesn't: Binary joins with equality predicates. Each join in the dataflow graph connects exactly two inputs via column-equality predicates (a.id = b.id ). The engine uses these equalities to maintain hash-based join state. Range predicates, expression-based join keys, or multi-table ON conditions are not directly supported. No correlated execution. In a traditional engine, a correlated subquery runs once per outer row - a nested-loop pattern. The dataflow graph has no concept of "per outer row." Every operator sees the full stream of changes from its inputs. Correlated subqueries must be rewritten into equivalent joins that the dataflow can maintain incrementally. Flat join structure preferred. Derived tables (subqueries in FROM ) are supported by the dataflow engine, but with a cost: a derived table compiles into a fully materialized intermediate node that cannot be parameterized. Every distinct combination of input data produces a stored result, regardless of whether the outer query needs it. Inlining the derived table - absorbing its FROM items, WHERE filters, and projections into the outer query - eliminates this intermediate materialization and lets the engine build parameterized lookups directly against the base tables. The rewrite pipeline aggressively inlines derived tables where semantically safe, reserving the materialized form for cases where inlining would change the query's meaning. Supported join types. The engine supports INNER JOIN , LEFT OUTER JOIN , and CROSS JOIN . RIGHT JOIN and FULL OUTER JOIN are not supported because their incremental maintenance in a streaming context requires tracking absence of matches on both sides - a significantly harder problem. Aggregation boundaries. GROUP BY and aggregate functions (COUNT , SUM , etc.) compile into stateful operators that maintain running totals. The engine requires that each aggregated query projects at least one aggregate-derived expression, and that GROUP BY keys are explicit column references (not positional numbers or aliases). These constraints mean that a syntactically valid SQL query - one that PostgreSQL or MySQL would execute without complaint - may not be directly compilable by Readyset, or may compile into a less efficient dataflow graph than necessary. The query rewrite pipeline exists to bridge this gap. The Rewrite Pipeline: Bridging SQL to Dataflow When a user issues CREATE CACHE for a query, Readyset runs the query through a multi-pass rewrite pipeline that transforms arbitrary SQL into the canonical form the dataflow engine expects. The pipeline is organized into three blocks: | Block | Purpose | Example | |---|---|---| | A - Normalization | Desugar syntax, resolve schemas, qualify columns | SELECT * becomes SELECT t.id, t.name, ... | | B - Deep Rewrites | Decorrelate subqueries, flatten derived tables, optimize | WHERE id IN (SELECT ...) becomes a join | | C - Cleanup | Remove redundant clauses, parameterize literals | ORDER BY id LIMIT 10 removed when result is provably single-row | Each pass is semantics-preserving: the rewritten query returns the same result as the original for all possible data. The passes build on each other - Block A normalizes the SQL into a canonical form that Block B's transformations can reliably operate on, and Block C cleans up artifacts left by Block B. Block A: Making SQL Canonical Before any deep transformation can happen, the query must be in a predictable shape. Block A handles this: - Schema resolution binds table and column names to the actual schema metadata Readyset has replicated from the upstream database. This is essential for later passes that need to know primary keys, unique constraints, and column types. - Star expansion replaces SELECT * with the explicit column list. Every downstream pass expects to see named columns, not wildcards. - Column qualification ensures every column reference is prefixed with its table name ( id becomest.id ). This prevents ambiguity when multiple tables have columns with the same name. - USING desugaring converts JOIN ... USING(id) intoJOIN ... ON (a.id = b.id) . The dataflow engine works withON predicates, notUSING clauses. After Block A, the query is fully resolved, qualified, and desugared - a clean foundation for the transformations that follow. Block B: The Heavy Lifting Block B is where the real work happens. These passes transform SQL constructs that the dataflow engine cannot handle into equivalent constructs that it can. The ordering matters - each pass prepares the ground for the next. Array Constructor Rewrite The first Block B pass, and PostgreSQL-specific. PostgreSQL's ARRAY(SELECT ...) constructor produces an array value from a subquery's rows - a per-outer-row scalar aggregation shape the dataflow engine has no direct operator for. The pipeline rewrites each occurrence into a LATERAL LEFT JOIN whose body wraps the original subquery in array_agg(...) , with a COALESCE(..., ARRAY[]) on the outer side so empty subqueries yield an empty array rather than NULL . ORDER BY and DISTINCT inside the constructor are copied into the array_agg call (required for correctness when the subquery also has LIMIT /Top-K ). By running first, this pass turns a SQL construct the later passes wouldn't recognize into a LATERAL join they already know how to handle. -- Before: ARRAY(SELECT ...) constructor - no direct dataflow operator SELECT u.name, ARRAY(SELECT p.title FROM posts p WHERE p.user_id = u.id) AS post_titles FROM users u -- After: LATERAL + array_agg + COALESCE - shapes that the rest of the pipeline understands SELECT u.name, COALESCE(array_subq.agg_result, ARRAY[]) AS post_titles FROM users u LEFT JOIN LATERAL ( SELECT array_agg(inner_subq.title) AS agg_result FROM (SELECT p.title FROM posts p WHERE p.user_id = u.id) inner_subq ) array_subq ON TRUE Redundant Join Elimination Queries generated by ORMs often contain redundant self-joins - the same table joined to itself on its primary key, with all projected columns coming from one side. The pipeline detects and eliminates these early, reducing the join graph complexity before the heavier transformations that follow. Left-Spine Hoisting Before decorrelating subqueries, the pipeline attempts to inline the leftmost derived table in FROM - but only at the top level of the query. This is the one position where a Top-K pattern (ORDER BY ... LIMIT ) is most valuable: at the top level, the LIMIT can be parameterized (e.g., LIMIT ? ) and the dataflow compiler can deploy a native Top-K node - a streaming operator purpose-built for maintaining the top N rows incrementally. If the subquery were left nested and processed later by the general decorrelation pass, the Top-K would be replaced with a ROW_NUMBER() based filter - functionally correct but less efficient, and with the LIMIT no longer parameterizable. -- Before: Top-K nested in a derived table - dataflow would emit a ROW_NUMBER filter SELECT sq.id, sq.name, sq.score FROM (SELECT id, name, score FROM products ORDER BY score DESC LIMIT ?) AS sq -- After left-spine hoist: Top-K at the top level - native, parameterized Top-K node SELECT id, name, score FROM products ORDER BY score DESC LIMIT ? With the ORDER BY and LIMIT at the top level, the dataflow compiler recognizes the Top-K shape and deploys a native streaming operator. The LIMIT remains a parameter (? ), so a single cached dataflow graph serves every LIMIT 10 , LIMIT 50 ,

Read on Hacker News ↗ ← Back to News

Comments

No comments yet. Start the discussion.