diff --git a/docs/README.md b/docs/README.md index 9f802e094..ec15798f5 100644 --- a/docs/README.md +++ b/docs/README.md @@ -11,7 +11,7 @@ Choose the path that matches what you need to do. ## Run or inspect a query -Follow [Run and inspect a query](user_guide_docs/run-a-query.md). Tool-specific setup stays with the tool, including the [DAG viewer instructions](../tools/dag-viewer/RUNNING.md). +Follow [Run and inspect a query](user_guide_docs/run-a-query.md). To see how the planner plans a workload stage by stage, and why it selects its plan, use [the Stage Viewer](user_guide_docs/stage-viewer.md). ## Embed the library diff --git a/docs/design_docs/proposals/planner-layering.md b/docs/design_docs/proposals/planner-layering.md index 5273bfb3f..3102e8b3a 100644 --- a/docs/design_docs/proposals/planner-layering.md +++ b/docs/design_docs/proposals/planner-layering.md @@ -3,20 +3,24 @@ Status: proposal. Audience: designers and developers of ASAPPlanner and of deployments such as ASAPQuery-backend. -Read Goal for the motivation and assumptions, Stages for the overview, the -stage sections for the rules, the examples for why the rules are needed, and -Scenarios for how the design is extended. +Read Problem Definition for the motivation and requirements, ASAPPlanner Design +for the design intuition, assumptions and stage overview, the stage sections for +the rules, the examples for why the rules are needed, and Scenarios for how the +design is extended. ## Contents -- [Goal](#goal) +- [Problem Definition](#problem-definition) - [Motivation](#motivation) + - [Design Requirements for ASAPPlanner](#design-requirements-for-asapplanner) +- [ASAPPlanner Design](#asapplanner-design) + - [Design Intuition](#design-intuition) - [Assumptions](#assumptions) -- [Stages](#stages) -- [Stages and their decisions](#stages-and-their-decisions) + - [ASAPPlanner System Overview: Planning Stages](#asapplanner-system-overview-planning-stages) +- [ASAPPlanner Detailed Design: Planning Stages and their decisions/Strategies](#asapplanner-detailed-design-planning-stages-and-their-decisionsstrategies) - [0. Language-specific frontends](#0-language-specific-frontends) - [1. Logical ASAP-aware optimization](#1-logical-asap-aware-optimization) - - [Pass 1: Local candidate generation](#pass-1-local-candidate-generation) + - [Pass 1: Per-computation candidate generation](#pass-1-per-computation-candidate-generation) - [Pass 2: ASAP-aware common-subexpression elimination](#pass-2-asap-aware-common-subexpression-elimination) - [2. Physical ASAP-aware optimization](#2-physical-asap-aware-optimization) - [Materialization](#materialization) @@ -34,49 +38,101 @@ Scenarios for how the design is extended. - [Supporting a new query construct](#supporting-a-new-query-construct) - [Adding a new summary family](#adding-a-new-summary-family) - [Adding a better cost or accuracy estimation](#adding-a-better-cost-or-accuracy-estimation) +- [Appendix: Windows and window summaries](#appendix-windows-and-window-summaries) -## Goal +## Problem Definition ### Motivation -* Existing database query engines and optimizers do not consider ASAP - primitives for optimizing the queries: summaries such as sketches that trade bounded error for lower - cost. -* They also do not consider the query and data workloads of different use - cases, which have the potential to share the common optimization with ASAP primitives. Example use case workloads can be streaming or batch data input, and repeated, batch or ad hoc queries. -* So existing query planners miss the opportunity to share the benefits of ASAP primitives across - domains and use cases. +* Existing query planners miss the opportunity to share the benefits of ASAP primitives across domains and use cases. ASAP Primitives can be sketches, sampling, statistical models, wavelets, machine learning generative models, etc, such algorithms preserving the a certain semantic information from the raw data. -In order to achieve the goals, ASAPPlanner needs to abstract the modeling the following for use cases: +* ASAP Primitives can accelerate query execution, but the existing efforts are pretty ad-hoc and a point solution to a point use case scenario. For example, sketch for network monitoring, wavelets or sampling for AQP Query. There is no unified framework for leveraging the acceleration and cost-reduction benefits from mulitple primitive types. + * Lacking a unified way / systematical way to leverage existing primitives for benefits, where the primitives are already a lot of algorithms proposed in the literature. + * Duplicated efforts will be rediscovered when various use case scenarios trying to optimize with ASAP Primitives, they will re-discover and add the similar optimization rules again and again. But among use cases, they have the potential to share the common optimization with ASAP primitives, benefiting more use cases. + +* Therefore, ASAPPlanner aims for building a framework with the capbility to support ASAP Primitives within query execution plans, preserving the original query semantics, as well as fitting to differnet use case scenarios, for reusing the primitive acceleration ideas. + +### Design Requirements for ASAPPlanner + +* Query planner/optimizer should consider the query and data workloads of different use cases when doing query optimization. + * Example use case workloads can be streaming or batch data input, and repeated, batch or ad hoc queries. The query engine and optimzer should be more specialized for targeted workload modeling. (This is not a novelty argument but we should model the workloads and optimize based on specific workload.) + * For example, repeated dashboard queries over time should be optimized by considering the overlapping computation between two same query expression over time, which is partially considered in time series databases. Modeling the workloads is the first step for futher optimization plan generation. + +* In order to preserve the query semantics/application semantics and generate valid query plans, ASAPPlanner should model: + * *Logical Plan*: what data is needed and what computation is needed over the data. + * *Physical Plan*: How to retrieve/process data and what algorithm to use for a computation. + * We assume the use case deployment will generate the *execution plan* given a physical plan for now, given different use cases have different deployment constraints. + +* In order to leverage the benefits from ASAP Primitives for query acceleration and cost reduction, we need to + - (a) express the operations over ASAP Primitives, as well as + - (b) constructing the rules (replacement strategies) for replacing the original logical query plan with ASAP Primitive aware logical plans, where in this stage, summary/synopsis built from raw data is needed and computation over them is needed, in addition to over raw data. + +* ASAPPlanner should selecting the optimal physical plan. + * The optimization goals include reducing total query execution costs, obeying the application accuracy target. + * The use case deployment can take the physical plan and determine how the engines runs the plan at runtime (execution). So deployment constraints can be the input to ASAPPlanner for some invalid plans early pruning (e.g., if the deployment has no materialized view computation engine, MV is not an optimziation option). + * *Efficiency of selecting optimal physical plan*: ASAPPlanner currently doesn't consider how to early remove sub-optimal plans but just generates all possible candidates and then select an optimal one among all of them. + +As a summary, in order to achieve the requirements, ASAPPlanner needs to abstract the following modeling: | Modeling | Contents | |---|---| | Query workload | Queries with recurrence (repeated, batch, ad hoc), predictability, time selection, and accuracy and latency requirements | | Data workload | Arrival (streaming, at rest, or both), volume, rate, cardinality and distribution | | Deployment inputs | Empirical cost model, empirical accuracy model and execution capabilities | -| ASAP replacement strategies | Rules that replace a sub-DAG of the query expression with summary expressions, and the summary families each computation may use | +| ASAP Primitive operators | Supported operation abstraction over ASAP Primitives, such as SummaryCreation, SummaryUpdate, SummaryMerge, SummaryDeletion, SummarySubtraction| +| ASAP Primitive operation replacement strategies | Rules that replace a sub-DAG of the query expression with summary operators, forming a new ASAP-primitive-aware sub-DAG, preserving the same query semantic | -ASAPPlanner takes a [query workload](https://github.com/ProjectASAP/ASAPPlanner/blob/main/crates/types/src/workload.rs), a [data workload](https://github.com/ProjectASAP/ASAPPlanner/blob/main/crates/types/src/workload.rs#L531) and the deployment's -inputs (TODO: define this data structure, [#525](https://github.com/ProjectASAP/ASAPPlanner/issues/525)), and returns one optimal physical plan. It decides what is computed, how -it is computed, and which plan is best. The deployment only supplies inputs and -executes the plan: it provides its empirical cost model, empirical accuracy -model and capabilities, but never plans queries or selects plans. +## ASAPPlanner Design + +### Design Intuition + +The design of ASAPPlanner is largely inspired by the existing DB Query engine work. + +* *Semantic preserving*: We should borrow the DB Query Engine plan stages, i.e., the logical planning, physical planning separation. + * So that we can have a framework for analyzing the semantics of the queries, and easily design replacement strategies for replacing a query intent with a set of ASAP Primitive operations with the same semantic. + * We can add information of workload patterns over the planning framework and adds our own optimization rules targeting specific workloads. + +* *Materialized View (MV)*: Given that ASAP Primitives mostly work as summary over data, we leverage the idea of Materilzed View (MV) in Database world, and treat ASAP Primitives as MV, and automatically design optimization rules for rewriting queries, mapping queries, based on the operations ASAP primitives can support. + * MV serves as an good integration opportunity for ASAP Primitives to work for semantic-preserving query computation. ASAPPlanner designs MV rules with aware of ASAP Primitives, as well as workload, e.g., we can build MV/incremental MV for repeated queries over time, specifically considering the window summary in ASAP Primitive library. The assumption that we can leverage MV is also, we observed the query and data workloads, where querie can be repated or sharing computation among a batch of queries, the query patterns persist long enough, and queries can be answered by the stored information based on ASAP primitives. + * ASAPPlanner also provides the option of not computing materialized view at all but just execute the query over ASAP Primitive operators once and no reusing; this can potentially accelerate the queries as well due to the algorithmic complexity of ASAP Primitive computation and space is usually smaller than computing from raw data exactly. + +> https://dl.acm.org/doi/10.1145/376284.375706 +> https://cloudberry.apache.org/docs/performance/optimize-queries/use-auto-materialized-view-to-answer-queries/ +> https://www.alibabacloud.com/blog/detailed-explanation-of-query-rewriting-based-on-materialized-views_598129 +> https://docs.aws.amazon.com/redshift/latest/dg/materialized-view-auto-rewrite.html + +* *Common Subexpression Elimination (CSE)*: In DB, CSE refers to identifying the common sub-query-expression within a query or multiple queires, and compute the same subexpression once and reuse to eliminate redundant processing. + * One example can be the `(price * (1 - discount))` is computed once and being reused, rather than twice. + + ```SQL + SELECT (price * (1 - discount)) AS net_price + FROM orders + WHERE (price * (1 - discount)) > 100; + ``` + + * CSE also serves as a good integration opportunity for ASAP Primitives, because many ASAP Primitives support the pattern of one data structure supporting multiple query intent. For example, one UnivMon sketch supports L2 norm, entropy, and cardinality 3 query intents; one arbitrary sub-window query framework, e.g., Exponential Histogram suports queries within arbitrary sub-window within the outter-most window; wavelets coefficients set in a way can support all linear operations, such as sum, avg, count; and many other examples. ### Assumptions ASAPPlanner relies on the assumptions below. -1. **Query rewrite rules are given.** Rewrite and replacement rules (for - example, `avg` as `sum`/`count`, or a TopK as a Count-Min Sketch with a - heap) are written by developers or algorithm designers. ASAPPlanner applies them; it does not - discover or generate them. -2. **Summary-family capabilities are given.** For each summary family, a +1. **Summary-family capabilities are given.** For each summary family, an algorithm developer declares which computations it can answer, which estimates it can read out, how it is sized for an accuracy target, its error bound, and whether it can be merged, subtracted or deleted. ASAPPlanner does not automatically discover these capabilities. -3. **Frontend semantics are given.** Each frontend preserves its source + +2. **Query replacement strategies are given, ASAPPlanner implements some strategies as a default set of strategies** Rewrite and replacement strategies (for + example, `avg` as `sum`/`count`, or a TopK as a Count-Min Sketch with a + heap) are written by developers or algorithm designers. ASAPPlanner does not automatically discover or generate them. + * For ASAPPlanner, given one batch of queries (as a multi-root DAG), the strategies can be applied to a sub-DAG inside the DAG, one sub-DAG (part of the query execution) can be mapped to different candidates based on the strategies, this gives us an initial brute-force version of plan candidate generation. + * The semantic/structure of the strategies (rules) are based on the logical DAG representation, where a pattern of a sub-DAG in the DAG can be identified, and being replace by another sub-DAG with ASAP-aware operators (primitive operators). + * These strategies can be potentially shared in different use cases and workloads, we depend on the cost model to rank all possible plans. For example, both the repeated dashboard query over time series data, and the batch query execution over data at rest, can leverage the TopK as a Count-Min Sketch replacement strategy. + The strategies we have applied in ASAPPlanner, see Section [Stages and their decisions](#asapplanner-detailed-design-planning-stages-and-their-decisionsstrategies). + +3. **Query-language frontend (e.g., SQL or PromQL query language themselves) semantics are given.** Each frontend preserves its source language's behavior, and ASAPPlanner obeys their semantic behavior. + 4. **Workload descriptions are inputs.** Recurrence, predictability, requirements and the data workload are supplied with the workload. ASAPPlanner does not infer them from traffic. @@ -87,102 +143,109 @@ ASAPPlanner relies on the assumptions below. cheapest valid plan among the candidates produced by the given rules and capabilities, as estimated by the given models. It is not optimal over plans those rules cannot produce. -7. **The deployment executes the plan as given.** It does not change summary +7. **Efficiency of selecting optimal physical plan**: ASAPPlanner currently doesn't consider how to early prune sub-optimal plans but just generates all possible candidates and then select an optimal one among all of them. +8. **The deployment executes the plan as given.** It does not change summary choices or materialization. -## Stages +### ASAPPlanner System Overview: Planning Stages -In the diagram, × means the Cartesian product: each stage combines every option along one dimension with every option along the others. +ASAPPlanner takes a [query workload](https://github.com/ProjectASAP/ASAPPlanner/blob/main/crates/types/src/workload.rs), a [data workload](https://github.com/ProjectASAP/ASAPPlanner/blob/main/crates/types/src/workload.rs#L531) and the deployment's +inputs (TODO: define this data structure, issue [#525](https://github.com/ProjectASAP/ASAPPlanner/issues/525)), and returns one optimal physical plan. It decides what is computed, how it is computed, and which plan is best. The deployment only supplies data and query inputs and +executes the plan: it provides its empirical cost model, empirical accuracy +model and capabilities, but never optimize queries. -```text - Query workload - (PromQL / SQL / MetricsQL, - query recurrence, - accuracy requirements, - latency requirements) - - + Data workload - (streaming vs. data at rest, - data distribution, - cardinality) - - + Deployment inputs - (cost model, - accuracy model, - deployment capabilities) - │ - ▼ -┌────────────────────────────── ASAPPlanner ──────────────────────────────┐ -│ │ -│ 0. Language-specific frontends │ -│ Parse and convert source-language queries into a common logical │ -│ representation. Reject unsupported query expressions. │ -│ │ -│ Output: CandidateLogicalDAGs │ -│ │ -│ │ │ -│ ▼ │ -│ Logical planning — what to compute │ -│ │ -│ 1. Logical ASAP-aware optimization │ -│ Explore semantically equivalent and legal logical candidates: │ -│ │ -│ summary families │ -│ × query rewrites │ -│ × sharing one summary across multiple computations │ -│ │ -│ Output: CandidateLogicalASAPDAGs │ -│ │ -│ │ │ -│ ▼ │ -│ Physical planning — how to compute │ -│ │ -│ 2. Physical ASAP-aware optimization │ -│ Explore physical implementations of each logical candidate: │ -│ │ -│ materialization decisions │ -│ × physical operator implementations │ -│ × parallelism and partitioning │ -│ × resource management │ -│ │ -│ Output: CandidatePhysicalASAPDAGs │ -│ │ -│ │ │ -│ ▼ │ -│ 3. Plan selection │ -│ Evaluate complete physical candidates using the deployment's │ -│ empirical cost and accuracy models. Reject candidates that violate │ -│ accuracy, latency, or capability constraints. │ -│ │ -│ Choose the cheapest valid plan for the whole workload. │ -│ │ -└────────────────────────────────┬───────────────────────────────────────┘ - │ - ▼ - one selected PhysicalASAPDAG - │ - ▼ -┌────────────────────────────── Deployment ───────────────────────────────┐ -│ │ -│ 4. Execution │ -│ deployment executes the selected DAG (plan). │ -│ │ -└────────────────────────────────────────────────────────────────────────┘ -``` +In the diagram, × means the Cartesian product: each stage combines every option based on the replacement strategies along one dimension with every option along the others. -The planner receives three groups of inputs: +```mermaid +%%{init: {"flowchart": {"wrappingWidth": 900, "nodeSpacing": 30, "rankSpacing": 25}}}%% +flowchart TB + subgraph INPUTS["Inputs"] + direction LR + QW["Query workload
(PromQL / SQL / MetricsQL,
query recurrence,
accuracy requirements,
latency requirements)"]:::input + DW["+ Data workload
(streaming vs. data at rest,
data distribution,
cardinality)"]:::input + DI["+ Deployment inputs
(cost model,
accuracy model,
deployment capabilities)"]:::input + end -| Input | Contents | -|---|---| -| Query workload | Source-language expressions, recurrence, predictability, time selection, accuracy requirements, and latency requirements | -| Data workload | Data arrival pattern, sampling cadence, ingestion volume and rate, cardinality, and distribution | -| Deployment inputs | Empirical cost model, empirical accuracy model, and execution capabilities | + subgraph PLANNER["ASAPPlanner"] + direction TB + TOP[" "]:::anchor + TOP ~~~ H0 + subgraph ST0[" "] + H0["0. Query-language-specific frontends"]:::title + B0["Parse and convert source-language queries into a common logical
representation. Reject unsupported query expressions.

Output: CandidateLogicalDAGs (a set of LogicalDAG)"]:::body + H0 ~~~ B0 + end -## Stages and their decisions + subgraph LOGICAL[" "] + LT["Logical planning — what to compute"]:::group + LT ~~~ H1 + H1["1. Logical ASAP-aware optimization"]:::title + B1["Explore semantically equivalent and legal logical candidates:

summary families, summary operations
× query rewrites
× sharing one common subexpresion/summary across multiple computations

Output: CandidateLogicalASAPDAGs (a set of LogicalASAPDAG)"]:::body + H1 ~~~ B1 + end + + subgraph PHYSICAL[" "] + PT["Physical planning — how to compute"]:::group + PT ~~~ H2 + H2["2. Physical ASAP-aware optimization"]:::title + B2["Explore physical implementations of each logical candidate:

materialization decisions
× physical operator implementations
× parallelism and partitioning
× resource management

Output: CandidatePhysicalASAPDAGs (a set of PhysicalASAPDAG)"]:::body + H2 ~~~ B2 + end + + subgraph ST3[" "] + H3["3. Plan selection"]:::title + B3["Evaluate complete physical candidates using the deployment's
empirical cost and accuracy models. Reject candidates that violate
accuracy, latency, or capability constraints.

Choose the cheapest valid plan for the whole workload."]:::body + H3 ~~~ B3 + end + + B0 --> LT + B1 --> PT + B2 --> H3 + end + + subgraph DEPLOY[" "] + DT["Deployment"]:::group + DT ~~~ H4 + H4["4. Execution"]:::title + B4["deployment executes the selected DAG (plan)."]:::body + H4 ~~~ B4 + end + + QW ~~~ TOP + DW ~~~ TOP + DI ~~~ TOP + INPUTS --> TOP + B3 -- "one selected PhysicalASAPDAG" --> DT + + click H0 href "#0-language-specific-frontends" + click H1 href "#1-logical-asap-aware-optimization" + click H2 href "#2-physical-asap-aware-optimization" + click H3 href "#3-plan-selection" + click H4 href "#4-execution" + + classDef input fill:#f1f3f4,stroke:#5f6368,color:#000; + classDef title fill:none,stroke:none,color:#0969da,font-weight:bold; + classDef anchor fill:none,stroke:none,font-size:1px; + classDef group fill:none,stroke:none,color:#333; + classDef body fill:none,stroke:none,color:#000; + style INPUTS fill:#fff,stroke:#5f6368; + style ST0 fill:#fff,stroke:#5f6368; + style LOGICAL fill:#fff,stroke:#5f6368; + style PHYSICAL fill:#fff,stroke:#5f6368; + style ST3 fill:#fff,stroke:#5f6368; + style DEPLOY fill:#fff,stroke:#5f6368; +``` + +Stage details: +- [0. Frontends](#0-language-specific-frontends) · +- [1. Logical ASAP-aware optimization](#1-logical-asap-aware-optimization) · +- [2. Physical ASAP-aware optimization](#2-physical-asap-aware-optimization) · +- [3. Plan selection](#3-plan-selection) · +- [4. Execution](#4-execution) -Stages 0 to 2 each output a candidate set holding every semantically equivalent -and legal candidate DAG of that stage; stage 3 is the only step that chooses -one candidate DAG as output. Candidate sets are internal to ASAPPlanner and may +## ASAPPlanner Detailed Design: Planning Stages and their decisions/Strategies + +Stages 0 to 2 (logical and physical planning stages) each output a candidate set holding every semantically equivalent and legal candidate DAG of that stage; stage 3 is the only step that chooses one candidate DAG as output currently (TODO: early selection for valid and more efficient candidates of stage 0-2 can be designed later). Candidate sets are internal to ASAPPlanner and may be shared or enumerated lazily. Two kinds of removal are kept apart: * **Pruning** removes an **invalid** candidate. Any stage may prune, but only @@ -195,12 +258,12 @@ be shared or enumerated lazily. Two kinds of removal are kept apart: A candidate is a DAG for the **whole workload**, not for one query. Each stage combines its choices for every sub-DAG with the candidates it receives (the × in the diagram), so the candidate set grows from stage to stage until -selection picks one. Example 1 traces this growth step by step. +selection picks one. [Example 1](#example-1-aggregation-over-dimensions--the-candidate-set-through-every-stage) traces this growth step by step. | Stage | Input | Decides | Output | |---|---|---|---| | 0. Frontends | `query`, `language` | Parse and convert to a common logical form; reject what cannot be represented | `CandidateLogicalDAGs` | -| 1. Logical ASAP-aware optimization | Logical DAGs; accuracy requirements, `time_selection`, repetition interval | Summary replacement (Pass 1); ASAP-aware CSE (Pass 2) | `CandidateLogicalASAPDAGs` | +| 1. Logical ASAP-aware optimization | Logical DAGs; accuracy requirements, `time_selection`, repetition interval | Summary replacement (Pass 1); ASAP-aware Common Subexpression Elimination (CSE) (Pass 2) | `CandidateLogicalASAPDAGs` | | 2. Physical ASAP-aware optimization | Logical ASAP DAGs; `recurrence`, `predictability`, `DataWorkload` | Materialization; physical operators; parallelism and resources (TODO) | `CandidatePhysicalASAPDAGs` | | 3. Plan selection | Physical candidates; `requirements`; cost model, accuracy model, capabilities | Reject invalid candidates; pick the cheapest plan for the whole workload | One `PhysicalASAPDAG` | | 4. Execution (deployment) | The selected `PhysicalASAPDAG` | Run ingestion, storage and query-time computation | Query results | @@ -224,7 +287,11 @@ computation on its own; Pass 2 finds candidates that share computation across sub-DAGs and queries. Materialization and execution placement are decided in later stages. -#### Pass 1: Local candidate generation +#### Pass 1: Per-computation candidate generation + +**Per-computation** means each sub-DAG is considered on its own: its candidates +depend only on its own computation and accuracy requirement, not on any other +sub-DAG or query in the workload, so pass 1 doesn't consider the common subexpression sharing optimization. Sharing across sub-DAGs and queries is left to Pass 2. For each eligible sub-DAG, Pass 1 identifies its computation semantics, applies rewrite rules, and generates every candidate that is not provably unable to @@ -232,7 +299,7 @@ meet its accuracy requirement. Example for summary candidates: -| Original computation | Local candidates | +| Original computation | Per-computation candidates | |---|---| | `Sum(x) by (g)` | Exact grouped sum | | `TopK(k, x) by (g)` | Exact sort and limit per group, Count-Min Sketch with a top-*k* heap per group, Hydra over all groups | @@ -259,108 +326,357 @@ A summary-based candidate uses three kinds of summary nodes: One summary build node can feed several estimation nodes, which is what Pass 2 exploits. +**Algorithm 1: Logical Pass 1 — Candidate generation without considering CSE** + +```text +Input: D — CandidateLogicalDAGs from stage 0 (each a whole-workload LogicalDAG) + R — rewrite / replacement rules + F — summary-family capabilities + acc — accuracy requirement of each computation +Output: C1 — CandidateLogicalASAPDAGs with each sub-DAG replaced independently (no sharing) + +1: C1 ← ∅ +2: for each LogicalDAG d in D do +3: for each eligible sub-DAG s in d do +4: sem(s) ← IDENTIFY_SEMANTICS(s) // e.g., input expression, filter, grouping, window, computation +5: Replacements(s) ← { EXACT(s) } // the exact computation is always a candidate +6: for each rule r in R such that r.pattern matches s do +7: for each candidate c in r.APPLY(s, F) do // Algorithm 1.1 +8: // c is exact operators plus summary build / estimation nodes (no merge nodes yet) +9: if c provably cannot meet acc(s) then +10: PRUNE(c, reason) +11: else +12: ANNOTATE(c, input, filter, grouping, window, estimates, acc(s)) +13: Replacements(s) ← Replacements(s) ∪ { c } +14: end if +15: end for +16: end for +17: end for +18: // one whole-workload candidate per combination of per-sub-DAG replacements +19: for each choice (c_1, …, c_n) in Replacements(s_1) × … × Replacements(s_n) do +20: C1 ← C1 ∪ { SUBSTITUTE(d, s_1 ↦ c_1, …, s_n ↦ c_n) } +21: end for +22: end for +23: return C1 +``` + +*Example for Algorithm 1* (the workload of [Example 1](#example-1-aggregation-over-dimensions--the-candidate-set-through-every-stage)). Q1 needs an +exact answer, so every summary replacement for it is pruned and only the exact +computation remains. Q2 tolerates error and gets three replacements. The +Cartesian product gives 1 × 3 = 3 whole-workload candidates. + +```mermaid +flowchart LR + subgraph IN["d: one workload LogicalDAG"] + direction TB + s1["s₁ = Q1 · sum by (job) (rate(…[1m]))
accuracy: exact"]:::exact + s2["s₂ = Q2 · topk by (job) (10, sum_over_time(…[1m]))
accuracy: ε = 0.01"]:::exact + end + subgraph R1["Replacements(s₁)"] + direction TB + a1["Exact rate + sum"]:::exact + a2["summary candidates
✗ pruned: cannot be exact"]:::pruned + end + subgraph R2["Replacements(s₂)"] + direction TB + b1["Exact sort + limit per job"]:::exact + b2["CMS + top-10 heap per job"]:::summary + b3["Hydra over all jobs"]:::summary + end + subgraph OUT["C1 = Replacements(s₁) × Replacements(s₂) = 3 candidates"] + direction TB + c1["① exact Q1 + exact Q2"]:::ok + c2["② exact Q1 + CMS Q2"]:::ok + c3["③ exact Q1 + Hydra Q2"]:::ok + end + s1 --> R1 + s2 --> R2 + R1 --> OUT + R2 --> OUT + classDef data fill:#f1f3f4,stroke:#5f6368,color:#000; + classDef exact fill:#fff,stroke:#5f6368,color:#000; + classDef summary fill:#e8f0fe,stroke:#1a73e8,stroke-width:2px,color:#000; + classDef estimate fill:#e6f4ea,stroke:#188038,color:#000; + classDef pruned fill:#fff,stroke:#d93025,stroke-dasharray:4 3,color:#d93025; + classDef ok fill:#fff,stroke:#188038,stroke-width:2px,color:#000; +``` + +**Algorithm 1.1: r.APPLY(s, F) — Applying one rewrite rule to a sub-DAG** + +A rule `r` has a **pattern**, the shape of logical sub-DAG it matches (for +example `TopK(k, x) by (g)`), and one or more **replacement templates**, each a +sub-DAG of ASAP-aware operators with a summary slot to fill (for example "a +summary that estimates per-key frequency, one per group, plus a top-*k* heap"). +A template with no summary slot is a pure query rewrite, such as `avg` as +`sum` / `count`. + +```text +Input: r — a rewrite rule (pattern, replacement templates) + s — a sub-DAG that matches r.pattern + F — summary-family capabilities +Output: Out — candidate replacement sub-DAGs for s + +1: b ← MATCH(r.pattern, s) // binds input expression, filter, key / value, grouping, +2: // window, and parameters such as k or the quantile q +3: Out ← ∅ +4: for each template t in r.templates do // e.g. TopK: CMS + heap per group, Hydra over all groups +5: if t has no summary slot then +6: Out ← Out ∪ { t.INSTANTIATE(b) } // pure query rewrite +7: continue +8: end if +9: for each family f in F such that f supports t.required_estimates(b) +10: and f supports t.grouping_mode(b) do // per group or over all groups +11: build ← SUMMARY_BUILD(f, b.input, b.filter, b.key_or_value, b.grouping, b.window) +12: est ← SUMMARY_ESTIMATE(f, build, t.estimate(b)) // e.g. p99, top-k, entropy +13: c ← t.INSTANTIATE(b, build, est) // wires in the remaining exact operators, e.g. the heap +14: Out ← Out ∪ { c } +15: end for +16: end for +17: return Out +``` + +Summary nodes are created unsized. Algorithm 1 checks whether a family can +meet `acc(s)` at all, and the summary is sized for the accuracy requirement +later, since Pass 2 may tighten it to the strictest requirement among shared +consumers. Pass 1 creates no merge nodes: each candidate summarizes exactly its +own window, and merging across windows comes from the window-composition rule +in Pass 2. + +*Example for Algorithm 1.1* (Q2 of [Example 1](#example-1-aggregation-over-dimensions--the-candidate-set-through-every-stage)). MATCH binds the +parameters of Q2. The TopK rule has two templates: one summary per group, +or one summary over all groups. For each template, only families that support +the needed estimate are used: the Count-Min Sketch fills the per-group +template, Hydra fills the all-groups template, and KLL is skipped because it +cannot estimate per-key frequencies. + +```mermaid +flowchart LR + S["s · topk by (job) (10,
sum_over_time(http_requests_total[1m]))"]:::exact + B["b = MATCH(r.pattern, s)
input: http_requests_total
key: series · grouping: by job
window: 1m · k = 10"]:::exact + T1["template t₁:
frequency summary per group
+ top-k heap"]:::exact + T2["template t₂:
one summary over all groups"]:::exact + X["f = KLL
✗ skipped: no per-key frequency"]:::pruned + S -->|"line 1"| B + B --> T1 + B --> T2 + T1 -.->|"line 9"| X + subgraph C1["c₁ (f = Count-Min Sketch)"] + direction LR + m1["build: CMS per job"]:::summary --> m2["estimate: per-series sums"]:::estimate --> m3["top-10 heap per job"]:::exact + end + subgraph C2["c₂ (f = Hydra)"] + direction LR + h1["build: Hydra over job"]:::summary --> h2["estimate: top-10 per job"]:::estimate + end + T1 -->|"lines 11–13"| C1 + T2 -->|"lines 11–13"| C2 + classDef data fill:#f1f3f4,stroke:#5f6368,color:#000; + classDef exact fill:#fff,stroke:#5f6368,color:#000; + classDef summary fill:#e8f0fe,stroke:#1a73e8,stroke-width:2px,color:#000; + classDef estimate fill:#e6f4ea,stroke:#188038,color:#000; + classDef pruned fill:#fff,stroke:#d93025,stroke-dasharray:4 3,color:#d93025; + classDef ok fill:#fff,stroke:#188038,stroke-width:2px,color:#000; +``` + #### Pass 2: ASAP-aware common-subexpression elimination **Sub-DAG sharing** means several consumers reference one operator and its -upstream dependencies. Traditional CSE provides common sub-DAG sharing for +upstream dependencies. Traditional Common-Subexpression Elimination (CSE) provides common sub-DAG sharing for eligible, structurally identical computations. ASAP-aware CSE extends it with summary-specific sharing rules. Computations can share work when they use identical expressions, when one summary build node supports several estimates, or when one window summary can answer their overlapping windows. -The rules compare computations by their **summary input data**: what a summary for -that computation would ingest, namely the data source, the filters, and the key -or value being summarized together with its grouping. The summary input data does -not include the window; the window-composition rule compares windows -separately. - -The window-composition rule distinguishes the window a query reads from the -window summary that answers it: - -* A **window** is the time range one query evaluation reads, for example the - last 5 min. Consecutive evaluations of a repeating query read overlapping - windows. Most summaries cannot remove old data, so one summary cannot simply - slide forward with the window. -* A **window summary** keeps summaries so that many windows can be answered. - Three window summaries are considered for now: - * **Sliding window:** summaries over windows of a fixed length L that start - every s (the slide), so several windows are active at once. Each arriving - sample is inserted into every active window that contains it, at the cost - of more ingestion work and memory. A query window of length W is answered - from completed windows: - * **L = W:** each evaluation reads one completed window, with no merge. - For a 5-min window evaluated every 1 min, L = 5 min and s = 1 min, so 5 - windows are active and each sample updates all 5. This works even for - summaries that cannot be merged. - * **L shorter than W:** the query window is covered by W / L - non-overlapping completed windows, which are merged[^sliding-merge]. For - example, a 10-min window evaluated every 1 min merges two 5-min windows - with a 1-min slide. This needs a mergeable summary. - - L must divide W, and s must divide both L and the evaluation interval, so - the windows a query needs have always just completed. - * **Tumbling window:** back-to-back, non-overlapping windows of one fixed - length, each with one summary; a sliding window whose slide equals its - length. A longer query window is answered by - merging the tumbling windows it covers. The tumbling length must divide - both the query window length and the evaluation interval, so that every - query window starts and ends on a tumbling boundary: a 5-min window - evaluated every 1 min uses 1-min tumbling windows and merges exactly 5 of - them. It needs a mergeable summary. - * **Exponential Histogram (EH):** a sequence of EH buckets that covers a - long history. A query window is answered by merging the EH buckets it - covers. Few EH buckets cover a long history, at the cost that an old - query-window boundary may fall inside an EH bucket and is then - approximate. - * An **EH bucket** is one non-overlapping time range of the history with - one summary of the data in it. Unlike tumbling windows, EH buckets are - not all the same length: they grow with age, so recent data sits in - short EH buckets and older data in longer ones. Adjacent EH buckets are - merged into a longer one as they age. - - TODO: evaluate other sliding-window frameworks for sketches as further - window summaries, such as Smooth Histograms[^smooth-histograms], - MicroscopeSketch[^microscope-sketch] and Sliding Sketches[^sliding-sketches]. - | ASAP-aware CSE rule | Sharing condition | Shared computation | |---|---|---| | Identical-expression rule | The input and computation semantics are identical. | One common computation node serving multiple consumers. | | Summary-capability rule | The computations have the same summary input data and the same window, and one summary supports all requested computations and their accuracy requirements. | One summary build node feeding several estimation nodes, e.g. UnivMon → distinct count, entropy, L2 norm. | | Window-composition rule | The computations have the same summary input data, and one window summary can answer the requested windows within their accuracy requirements. | One window summary feeding per-query merge (where needed) and estimation nodes, e.g. a sliding-window or tumbling-window KLL, or an Exponential Histogram with a KLL per EH bucket. | +The rules compare computations by their **summary input data**: what a summary for +that computation would ingest, namely the data source, the filters, and the key +or value being summarized together with its grouping. The summary input data does +not include the window; the window-composition rule compares windows +separately. + The examples behind these rules: -* **Summary-capability rule (Example 2).** One UnivMon over `src_ip` from +* **Summary-capability rule ([Example 2](#example-2-one-summary-for-several-computations--the-summary-capability-rule-in-pass-2)).** One UnivMon over `src_ip` from `flows` in the last minute serves three queries refreshed every 10 s: `COUNT(DISTINCT src_ip)`, the entropy of the `src_ip` distribution, and the L2 norm of per-`src_ip` counts. Each flow record updates the UnivMon once; a distinct-count, an entropy and an L2 estimation node each compute their statistic from it. The UnivMon is sized for the strictest of the three accuracy requirements. -* **Window-composition rule, sliding or tumbling window (Example 3, - Pattern B).** For `quantile_over_time(0.99, latency_ms[5m])` repeated every +* **Window-composition rule, sliding or tumbling window ([Example 3, + Pattern B](#example-3-pattern-b)).** For `quantile_over_time(0.99, latency_ms[5m])` repeated every minute, one window summary serves every evaluation. With a sliding-window KLL, each sample updates the 5 active windows, and each evaluation reads the one that has just completed. With 1-min tumbling-window KLLs, each sample updates one window, and each evaluation merges the latest 5 with a merge node; consecutive evaluations share 4 of them. -* **Window-composition rule, Exponential Histogram (Example 3, Pattern A).** +* **Window-composition rule, Exponential Histogram ([Example 3, Pattern A](#example-3-pattern-a)).** One Exponential Histogram over the last 5 years, with a KLL per EH bucket, serves the p99 queries over `[5y]`, `[1y]`, `[1y] offset 1y`, `[1y] offset 2y` and `[3y] offset 2y`. Each query's merge node merges the EH buckets covering its interval, and its estimation node computes p99 from the merged KLL. -* **Other quantiles share for free.** One KLL answers every quantile, so adding +* **Other quantiles share for free ([Example 3, Pattern B](#example-3-pattern-b)).** One KLL answers every quantile, so adding `quantile_over_time(0.5, latency_ms[5m])` to the sliding-window dashboard adds only a p50 estimation node next to the p99 one, reading the same KLL, with no new summary. +Windows, sliding windows, tumbling windows and Exponential Histograms are defined in +[Appendix: Windows and window summaries](#appendix-windows-and-window-summaries). + Rules are defined by each summary family's capabilities and semantic requirements. A shared summary must meet the strictest accuracy requirement among its consumers. Applying a rule adds a shared candidate and keeps the independent candidates, so selection can compare both. +**Algorithm 2: Logical Pass 2 — ASAP-aware common-subexpression elimination** + +```text +Input: C1 — CandidateLogicalASAPDAGs from Pass 1 + F — summary-family capabilities (estimates, sizing, error bound, mergeable) +Output: C2 — CandidateLogicalASAPDAGs with sharing + +1: C2 ← C1 // independent candidates are kept +2: for each candidate d in C1 do +3: Opts ← ∅ // sharing options found in d +4: +5: // Identical-expression rule +6: for each group G of nodes in d with identical input and computation semantics, |G| ≥ 2 do +7: Opts ← Opts ∪ { one common node serving all consumers in G } +8: end for +9: +10: // Summary-capability rule +11: for each group G of computations in d with the same summary input data and window, |G| ≥ 2 do +12: for each family f in F that supports every estimate requested in G do +13: a ← strictest accuracy requirement in G +14: if f can be sized to meet a for every computation in G then +15: Opts ← Opts ∪ { one build node of f sized for a → one estimation node per computation } +16: end if +17: end for +18: end for +19: +20: // Window-composition rule (G may be a single repeating query) +21: for each group G of computations in d with the same summary input data do +22: for each family f in F that supports every estimate requested in G do +23: for each window summary w in WINDOW_SUMMARIES(G) do +24: if CAN_SHARE(w, f, G) then +25: Opts ← Opts ∪ { SHARED_WINDOW_SUMMARY(w, f, G) } +26: end if +27: end for +28: end for +29: end for +30: +31: for each non-empty, non-conflicting subset O ⊆ Opts do +32: C2 ← C2 ∪ { APPLYSHARING(d, O) } +33: end for +34: end for +35: return C2 + + +WINDOW_SUMMARIES(G): // candidate window summaries for G; each q in G has + // window length W_q and evaluation interval E_q + Tumbling(L) for every L that divides every W_q and every E_q + Sliding(L, s) for every L that divides every W_q, + and every s < L that divides L and every E_q // s = L is Tumbling(L) + EH one EH covering the oldest data any q in G reads + +PIECES(w, q): // how many summaries of w one evaluation of q reads + Tumbling(L) W_q / L + Sliding(L, s) W_q / L // 1 when L = W_q: read one completed window + EH number of EH buckets that overlap q's window + +CAN_SHARE(w, f, G): // can one w over f answer every query in G? + if some q in G has PIECES(w, q) > 1 and f is not mergeable then + return false + return w over f, sized for the strictest accuracy in G, meets every q's accuracy + // for EH this includes the bucket-boundary error + +SHARED_WINDOW_SUMMARY(w, f, G): + one window-summary node: w over f, sized for the strictest accuracy in G + for each q in G: + merge node: merges the PIECES(w, q) pieces covering q's window // omitted when PIECES = 1 + estimation node: computes q's estimate from the merged result (or the single piece) +``` + +The divisibility conditions make every summary a query needs a completed one +when the query evaluates; see +[Appendix: Windows and window summaries](#appendix-windows-and-window-summaries). +For example, `quantile_over_time(0.99, latency_ms[5m])` every 1 min +(W = 5 min, E = 1 min) gives `Tumbling(1 min)` with PIECES = 5, and +`Sliding(5 min, 1 min)` with PIECES = 1, which needs no merge node and so works +even for a summary that cannot be merged. + +*Example for Algorithm 2, summary-capability rule* ([Example 2](#example-2-one-summary-for-several-computations--the-summary-capability-rule-in-pass-2)). +Pass 1 gave each of the three queries its own UnivMon. All three have the same +summary input data (`src_ip` from `flows`) and the same 1-min window, and +UnivMon supports all three estimates, so lines 11–15 add one shared UnivMon, +sized for the strictest accuracy requirement. The separate UnivMons stay in +C2 as well. + +```mermaid +flowchart LR + subgraph BEFORE["In C1: one UnivMon per query"] + direction TB + u1["UnivMon · src_ip · 1m"]:::summary --> e1["distinct count"]:::estimate + u2["UnivMon · src_ip · 1m"]:::summary --> e2["entropy"]:::estimate + u3["UnivMon · src_ip · 1m"]:::summary --> e3["L2 norm"]:::estimate + end + subgraph AFTER["Added to C2: one shared UnivMon"] + direction TB + u["UnivMon · src_ip · 1m
sized for strictest accuracy"]:::summary + u --> f1["distinct count"]:::estimate + u --> f2["entropy"]:::estimate + u --> f3["L2 norm"]:::estimate + end + BEFORE -->|"same summary input data
+ same window"| AFTER + classDef data fill:#f1f3f4,stroke:#5f6368,color:#000; + classDef exact fill:#fff,stroke:#5f6368,color:#000; + classDef summary fill:#e8f0fe,stroke:#1a73e8,stroke-width:2px,color:#000; + classDef estimate fill:#e6f4ea,stroke:#188038,color:#000; + classDef pruned fill:#fff,stroke:#d93025,stroke-dasharray:4 3,color:#d93025; + classDef ok fill:#fff,stroke:#188038,stroke-width:2px,color:#000; +``` + +*Example for Algorithm 2, window-composition rule* ([Example 3, Pattern +B](#example-3-pattern-b)). G holds one query, a p99 over 5 min evaluated every +1 min (W = 5 min, E = 1 min). Counting in whole minutes, `WINDOW_SUMMARIES(G)` +returns `Tumbling(1m)`, since 1 min is the only length that divides both 5 and +1, plus `Sliding(5m, 1m)` and one EH. KLL is mergeable, so all three pass +`CAN_SHARE` (line 24) and become options. They all replace the same computation and +therefore conflict, so line 31 adds each one as a separate candidate. + +```mermaid +flowchart TB + Q["G = { q } · quantile_over_time(0.99, latency_ms[5m]) every 1 min · f = KLL"]:::exact + subgraph T["Tumbling(1m) · PIECES = 5"] + direction LR + t1["1-min KLLs"]:::summary --> t2["merge latest 5"]:::summary --> t3["p99"]:::estimate + end + subgraph SL["Sliding(5m, 1m) · PIECES = 1"] + direction LR + s1["5 active 5-min KLLs"]:::summary --> s3["p99 of the window
that just completed"]:::estimate + end + subgraph EH["EH · PIECES = buckets overlapping 5 min"] + direction LR + h1["EH, one KLL per bucket"]:::summary --> h2["merge covering buckets"]:::summary --> h3["p99"]:::estimate + end + Q --> T + Q --> SL + Q --> EH + classDef data fill:#f1f3f4,stroke:#5f6368,color:#000; + classDef exact fill:#fff,stroke:#5f6368,color:#000; + classDef summary fill:#e8f0fe,stroke:#1a73e8,stroke-width:2px,color:#000; + classDef estimate fill:#e6f4ea,stroke:#188038,color:#000; + classDef pruned fill:#fff,stroke:#d93025,stroke-dasharray:4 3,color:#d93025; + classDef ok fill:#fff,stroke:#188038,stroke-width:2px,color:#000; +``` + ### 2. Physical ASAP-aware optimization Physical optimization turns each logical candidate into physical candidates. @@ -417,6 +733,97 @@ Physical operator implementation converts every node to physical operators, for example TopK as a sort followed by a limit, or a KLL node as summary build, merge and quantile estimation operators. +**Algorithm 3: Physical ASAP-aware optimization** + +```text +Input: C2 — CandidateLogicalASAPDAGs from stage 1 + recurrence, predictability, DataWorkload +Output: C3 — CandidatePhysicalASAPDAGs + +1: C3 ← ∅ +2: for each LogicalASAPDAG d in C2 do +3: // Materialization: options per sub-DAG +4: for each sub-DAG s in d do +5: Opt(s) ← { NotMaterialized } +6: for each medium in {memory, disk} do +7: Opt(s) ← Opt(s) ∪ { MatAtIngestion(medium), MatAtQuery(medium) } +8: end for +9: end for +10: +11: for each assignment m in Opt(s_1) × … × Opt(s_n) do +12: // constraint: everything upstream of an ingestion-time node also runs at ingestion time +13: if some s with m(s) = MatAtIngestion has an upstream node u with m(u) ≠ MatAtIngestion then +14: PRUNE(m, "upstream of ingestion-time node not at ingestion time") +15: continue +16: end if +17: // constraint: keep a materialized output as long as any consumer needs it +18: for each s with m(s) ≠ NotMaterialized do +19: retention(s) ← latest time any consumer of s reads it // from recurrence, windows +20: end for +21: // a shared summary is one node, so it is materialized once for all consumers +22: +23: // Physical operator implementation +24: for each node v in d do +25: Impl(v) ← physical operator implementations of v // e.g. TopK → sort + limit +26: end for +27: for each choice (i_1, …, i_k) in Impl(v_1) × … × Impl(v_k) do +28: C3 ← C3 ∪ { BUILDPHYSICAL(d, m, retention, i_1, …, i_k) } +29: end for +30: end for +31: // TODO: parallelism, partitioning, resource management +32: end for +33: return C3 +``` + +`recurrence`, `predictability` and the `DataWorkload` do not prune options +here; they determine the cost of each option, which stage 3 uses to pick the +cheapest plan (the typical outcomes above). + +*Example for Algorithm 3* ([Example 4, Pattern B](#example-4-materialization-of-window-summaries-in-physical-planning)). +The logical candidate is the `Tumbling(1m)` KLL option from Algorithm 2. The +diagram shows three of its materialization assignments, each drawn as a copy +of the DAG with every node colored by its materialization. Line 13 enforces one +rule: a node can run at ingestion time only if every node it reads from also +runs at ingestion time, because otherwise its input does not exist yet when +data arrives. In m₂ the merge node is placed at ingestion time, but the 1-min +KLLs it merges are built only at query time, so there is nothing to merge +while data is arriving, and m₂ is pruned. Assignments m₁ and m₃ are both +valid, and stage 3 chooses between them by cost. + +```mermaid +flowchart TB + subgraph M1["m₁ · ✓ kept"] + direction LR + a0[("latency_ms")]:::data --> a1["KLL build node
1-min tumbling
ingestion · memory
kept 5 min"]:::ingest --> a2["KLL merge node
latest 5
query time · not stored"]:::notmat --> a3["p99 estimation node
query time · not stored"]:::notmat + end + subgraph M2["m₂ · ✗ pruned (line 13)"] + direction LR + b0[("latency_ms")]:::data --> b1["KLL build node
1-min tumbling
query time · stored"]:::qtime -->|"✗ KLLs not built yet
when data arrives"| b2["KLL merge node
latest 5
ingestion"]:::ingest --> b3["p99 estimation node
query time · not stored"]:::notmat + end + subgraph M3["m₃ · ✓ kept"] + direction LR + c0[("latency_ms")]:::data --> c1["KLL build node
1-min tumbling
query time · not stored"]:::notmat --> c2["KLL merge node
latest 5
query time · not stored"]:::notmat --> c3["p99 estimation node
query time · not stored"]:::notmat + end + subgraph P1["BUILDPHYSICAL for m₁ (lines 24–28)"] + direction LR + p1["KLL insert operator
at ingestion"]:::ingest --> p2["KLL merge operator
at query time"]:::notmat --> p3["quantile operator
at query time"]:::notmat + end + subgraph LEG["Legend: when a node runs and whether its output is stored"] + direction LR + l1["MatAtIngestion"]:::ingest ~~~ l2["MatAtQuery"]:::qtime ~~~ l3["NotMaterialized"]:::notmat + end + M1 --> P1 + LEG ~~~ M1 + M1 ~~~ M2 + M2 ~~~ M3 + linkStyle 4 stroke:#d93025,stroke-width:2px,color:#d93025 + classDef data fill:#f1f3f4,stroke:#5f6368,color:#000; + classDef ingest fill:#e8f0fe,stroke:#1a73e8,stroke-width:2px,color:#000; + classDef qtime fill:#fef7e0,stroke:#e37400,stroke-width:2px,color:#000; + classDef notmat fill:#fff,stroke:#5f6368,stroke-dasharray:4 3,color:#000; + style M2 stroke:#d93025,stroke-dasharray:4 3 +``` + ### 3. Plan selection Selection rejects every candidate that misses an accuracy target or a latency @@ -437,6 +844,8 @@ given: it does not choose among summaries or decide what to materialize. ## End-to-end examples +To see examples in DAG Viewer, following the instructions [here](TODO: write an instruction and link here). + Each example's workload is shown as tables. Field names in code font are the fields of [`workload.rs`](https://github.com/ProjectASAP/ASAPPlanner/blob/main/crates/types/src/workload.rs). @@ -507,7 +916,7 @@ flowchart TB classDef estimate fill:#e6f4ea,stroke:#188038,color:#000; ``` -**Stage 1, Pass 1: 3 candidates.** Pass 1 finds local options for each query: +**Stage 1, Pass 1: 3 candidates.** Pass 1 finds per-computation options for each query: * **Q1** has one option, the exact per-series rate and per-`job` sum. Its accuracy requirement is exact, so no summary qualifies. @@ -519,7 +928,7 @@ flowchart TB ```mermaid flowchart LR - subgraph C["Q2's three local options"] + subgraph C["Q2's three per-computation options"] direction TB subgraph E["Exact"] direction LR @@ -748,7 +1157,7 @@ The data workload differs from the shared one in two fields: | `data_ingestion_interval` | not needed for SQL | **Pass 1.** Rewrite rules recognize the three computations, and each gets its -local candidates from the Pass 1 table: exact, a specialized summary, or +per-computation candidates from the Pass 1 table: exact, a specialized summary, or UnivMon. Combined, that is 3 × 3 × 3 = 27 workload candidates. **Pass 2.** All three computations have the same summary input data (`src_ip` @@ -774,7 +1183,7 @@ flowchart LR q3["L2(src_ip)"]:::exact end - subgraph P1["Stage 1, Pass 1 · local candidates per computation"] + subgraph P1["Stage 1, Pass 1 · candidates per computation"] direction TB subgraph D["Distinct"] direction LR @@ -841,6 +1250,7 @@ usually wins because each flow record updates one summary instead of three. This example has two workload patterns that both lead to a shared window summary. + **Pattern A: a batch of sub-interval queries over historical data.** An analyst submits a batch of p99 latency reports over different historical intervals, all executed together at time T. @@ -910,6 +1320,7 @@ flowchart LR classDef estimate fill:#e6f4ea,stroke:#188038,color:#000; ``` + **Pattern B: one repeating query with overlapping windows.** A real-time p99 panel over the last 5 min, refreshed every minute. @@ -1136,6 +1547,55 @@ query in the workload, and a shared summary is costed once. **Unchanged:** frontends, rules, summary families and the deployment's execution. +## Appendix: Windows and window summaries + +The window-composition rule in [Pass 2](#pass-2-asap-aware-common-subexpression-elimination) distinguishes the window a query reads from the +window summary that answers it: + +* A **window** is the time range one query evaluation reads, for example the + last 5 min. Consecutive evaluations of a repeating query read overlapping + windows. Most summaries cannot remove old data, so one summary cannot simply + slide forward with the window. +* A **window summary** keeps summaries so that many windows can be answered. + Three window summaries are considered for now: + * **Sliding window:** summaries over windows of a fixed length L that start + every s (the slide), so several windows are active at once. Each arriving + sample is inserted into every active window that contains it, at the cost + of more ingestion work and memory. A query window of length W is answered + from completed windows: + * **L = W:** each evaluation reads one completed window, with no merge. + For a 5-min window evaluated every 1 min, L = 5 min and s = 1 min, so 5 + windows are active and each sample updates all 5. This works even for + summaries that cannot be merged. + * **L shorter than W:** the query window is covered by W / L + non-overlapping completed windows, which are merged[^sliding-merge]. For + example, a 10-min window evaluated every 1 min merges two 5-min windows + with a 1-min slide. This needs a mergeable summary. + + L must divide W, and s must divide both L and the evaluation interval, so + the windows a query needs have always just completed. + * **Tumbling window:** back-to-back, non-overlapping windows of one fixed + length, each with one summary; a sliding window whose slide equals its + length. A longer query window is answered by + merging the tumbling windows it covers. The tumbling length must divide + both the query window length and the evaluation interval, so that every + query window starts and ends on a tumbling boundary: a 5-min window + evaluated every 1 min uses 1-min tumbling windows and merges exactly 5 of + them. It needs a mergeable summary. + * **Exponential Histogram (EH):** a sequence of EH buckets that covers a + long history. A query window is answered by merging the EH buckets it + covers. Few EH buckets cover a long history, at the cost that an old + query-window boundary may fall inside an EH bucket and is then + approximate. + * An **EH bucket** is one non-overlapping time range of the history with + one summary of the data in it. Unlike tumbling windows, EH buckets are + not all the same length: they grow with age, so recent data sits in + short EH buckets and older data in longer ones. Adjacent EH buckets are + merged into a longer one as they age. + + TODO: evaluate other sliding-window frameworks for sketches as further + window summaries, such as Smooth Histograms[^smooth-histograms], + MicroscopeSketch[^microscope-sketch] and Sliding Sketches[^sliding-sketches]. [^smooth-histograms]: V. Braverman and R. Ostrovsky. [Smooth Histograms for Sliding Windows](https://web.cs.ucla.edu/~rafail/PUBLIC/82.pdf). FOCS 2007. An alternative to EH. [^microscope-sketch]: Y. Wu et al. [MicroscopeSketch: Accurate Sliding Estimation Using Adaptive Zooming](https://yangtonghome.github.io/uploads/MicroscopeSketch_SIGKDD_23_final_paper.pdf). KDD 2023. diff --git a/docs/user_guide_docs/README.md b/docs/user_guide_docs/README.md index e7a845573..c12f697b9 100644 --- a/docs/user_guide_docs/README.md +++ b/docs/user_guide_docs/README.md @@ -4,3 +4,4 @@ User guides show supported commands, inputs, outputs, and ways to verify the result. - [Run and inspect a query](run-a-query.md) +- [See how ASAPPlanner plans a workload: the Stage Viewer](stage-viewer.md) diff --git a/docs/user_guide_docs/stage-viewer.md b/docs/user_guide_docs/stage-viewer.md new file mode 100644 index 000000000..04b79612f --- /dev/null +++ b/docs/user_guide_docs/stage-viewer.md @@ -0,0 +1,225 @@ +# See how ASAPPlanner plans a workload: the Stage Viewer + +The Stage Viewer shows, in a browser, what ASAPPlanner does with one workload +at each of its planning stages, and why it selects the plan it selects. Use it +to check a plan by hand, to compare the alternatives the planner considered, +and to see how the deployment's capabilities and the query requirements +change the choice. + +The stages are those of +[ASAPPlanner's design](../design_docs/proposals/planner-layering.md#asapplanner-system-overview-planning-stages): + +| Stage | What the viewer shows | +| --- | --- | +| 0. Frontends | The logical DAG the queries lower into, one root per query | +| 1. Logical ASAP-aware optimization | The logical ASAP candidates: each choice of exact or approximate computation per query (Pass 1), with shared inputs, shared summaries and tumbling panes (Pass 2) | +| 2. Physical ASAP-aware optimization | The physical candidates of each logical candidate, which differ in what runs at ingestion time and what runs at query time | +| 3. Plan selection | Each physical candidate checked for accuracy, latency and deployment capabilities, priced per second, and the cheapest selected | + +The viewer reads one JSON document per workload, written by the +`stage_pipeline` command. It does not run the plans. + +> **Based on [PR #574](https://github.com/ProjectASAP/ASAPPlanner/pull/574)** +> (branch `stack/509-viewer-stages`), which redesigns the viewer, together with +> the PRs it is stacked on: #613 (`stage_pipeline` examples 2 and 4b, and +> pricing every plan), #610 (deployment inputs in the document) and #609 +> (deployment capabilities). Until they are merged, check out +> `stack/509-viewer-stages` to use the viewer as described here; `main` still +> has the earlier Pre/Post-ASAP viewer. + +## Start the viewer + +Run from the repository root, with Rust/Cargo and Python 3 installed: + +```sh +python3 tools/dag-viewer/server.py +``` + +The server builds `stage_pipeline` (the first build takes a few minutes), +writes the built-in example documents into `tools/dag-viewer/out/`, and +prints the address to open: + +```text +ASAPPlanner Stage Viewer: http://127.0.0.1:8000 +``` + +| Option | Use it to | +| --- | --- | +| `--port 8765` | Serve on another port | +| `--host 0.0.0.0` | Accept connections from other machines (the default accepts only local ones) | +| `--skip-build` | Reuse an already-built `stage_pipeline` from `target/debug/`, or from `$CARGO_TARGET_DIR/debug/` if that variable is set | +| `--regenerate` | Rewrite the example documents, for example after changing the planner | + +If the server runs on a remote machine, forward the port and open +`http://127.0.0.1:8000` locally: + +```sh +ssh -L 8000:127.0.0.1:8000 +``` + +Press `Ctrl+C` in the server's terminal to stop it. + +## Look at an example + +The tabs at the top open the examples of the design document's +[end-to-end examples](../design_docs/proposals/planner-layering.md#end-to-end-examples). +Each tab starts with a sentence that says what the workload is and why its +plan wins. + +| Tab | Workload | What to look for | +| --- | --- | --- | +| 1 · dashboard panels | Two PromQL panels over 1M counter series every 10 s: an exact per-job rate, and the approximate top-10 series per job | The exact query stays exact, the top-k uses one CMS+heap over the raw samples, and both read one shared scan | +| 2 · SQL flow statistics | Distinct sources, entropy and L2 of per-source counts over the last minute of flows | Every UnivMon plan is invalid with the reason "no accuracy model for UnivMon", so an exact plan wins | +| 3a · historical p99 batch | Five p99 reports over 1–5 years, run once | The cheapest plan shares one input scan across the five KLL sketches | +| 3b · live p99 panel | p99 over the last 5 min every minute, 1M series | Keeping KLL panes from ingestion time costs more memory than building the sketch at each evaluation, and rebuilding all panes at query time is over the 200 ms latency bound | +| 4a · monthly p99 reports | Example 3a repeated monthly | The same plan, with its cost amortized over a month | +| 4b · panes that pay off | p99 over the last hour every 10 min, 1k series sampled every second, raw data not kept by the deployment | Maintaining six 10-min KLL panes at ingestion time wins | + +To link to one example, add its name after `#`, for example +`http://127.0.0.1:8000/#example4b`. + +The workload statistics in the examples (series counts, ingestion rates) are +illustrative. Read the costs as a ranking, not as a prediction of real +resource use. + +## Read the page + +### Workload queries + +One box per query, with its language and text and the requirements the +planner must meet: + +- **accuracy**: `exact`, or `ε=…, δ=…` (relative error ε with probability at + least 1 − δ); +- **latency**: the maximum time an evaluation may take at query time; +- **every N s**: how often a repeating query runs. + +The line under the queries counts the candidates at each stage, for example +`1 logical DAG → 88 logical ASAP → 112 physical → 1 selected (40 invalid, 71 +valid but costlier)`. + +### Deployment inputs + +What the deployment told the planner, which Stage 3 uses to accept or reject +plans: + +| Row | Meaning | +| --- | --- | +| Exact aggregates | Exact computations the executor can run (Sum, Count, Min, Max, Increase, Rate) | +| Sketches → estimates | Each sketch the executor can build, with the estimates it can be read for, for example `Kll → Quantile`. A plan that needs anything else is invalid. | +| Ingestion-time maintenance | Whether the deployment can keep state up to date as data arrives | +| Keep query-time results across evaluations | Whether results built at query time can be kept for the next evaluation | +| Memory budget | The most state the deployment can keep, or none | +| Raw data | Whether the deployment keeps raw samples anyway. If it does not, plans that read raw data at query time pay to keep the samples they read. | +| Cost model | The cost model, its unit, and its calibration constants, such as the cost of one CPU operation and of keeping one byte for one second | +| Accuracy model | The model that decides whether a sketch meets a query's accuracy target | + +### Stage 3 · plans by cost + +Every physical plan, cheapest first, with its total cost per second: + +- **✓ selected**: the plan ASAPPlanner chooses. +- **valid · costlier**: meets every requirement, but costs more. +- **✗ invalid**: fails a check, with the reason. Examples: an accuracy target + the sketch cannot guarantee, an evaluation over the latency bound, or a + summary the deployment cannot build. + +Each label names the choice per query, for example `Q1 exact · Q2 CMS+heap · +shared input`, and, for physical plans, what runs at ingestion time, for +example `ingestion time: Kll ×6 panes`. + +Click a plan to show it in the three lanes. For workloads with many plans, +the document may carry only the cheapest ones. The list then says so, for +example `showing 64 of 486 plans, cheapest first; Stage 3 priced 486 and +selected among all of them`. The selection is made over all plans either +way. + +### The three lanes + +| Lane | Shows | +| --- | --- | +| Stage 0 · Logical DAG | The queries as lowered by the frontends | +| Stage 1 · Logical ASAP DAG | The logical candidate of the plan you picked. The **candidate** menu shows any other. | +| Stage 2 · Physical ASAP DAG | The plan you picked, with each node's timing (ingestion time or query time) and Stage 3 cost. The **candidate** menu shows any other. | + +Data flows from the bottom up. Node colors: + +- blue: data sources; +- grey: relational operators (filter, project, aggregate, sort, …); +- purple: summary operators (build a sketch or exact accumulator, merge, + estimate, finalize). + +A thick border marks a query's result. A sketch build shows its configuration +and how many instances it keeps, for example `summary: CmsWithHeap · depth 7 +· heap 100 · width 272` and `instances: one per group`. + +Drag to pan and scroll to zoom. + +### Details + +Click a node to see: + +- its operator and node id; +- for a summary, its parameters and instances; +- which query it answers; +- when it runs and its Stage 3 cost; +- its output schema and coverage (the data and time range a summary + represents); +- the accuracy guarantee it carries; +- the full operator as JSON. + +Click an edge to see the schema and data state it carries. + +## Plan your own queries + +### In the page + +Click **Query editor**, enter PromQL queries one per line, and click +**Plan**: + +| Field | Meaning | +| --- | --- | +| ε | Relative error allowed for every query. Leave it empty for exact queries. | +| δ | Failure probability allowed. Needs ε. | +| sample interval (ms) | How often each series is sampled (default 15000) | + +The editor plans the queries as one batch that runs once. It declares only +the sample interval, so the planner uses default statistics for everything +else. SQL is not available in the editor yet; use the command below with a +built-in example, or a document written elsewhere. + +### From the command line + +`stage_pipeline` writes the document the viewer reads: + +```sh +cargo run -p asap-devtools --bin stage_pipeline -- \ + --promql "topk by (job) (10, sum_over_time(http_requests_total[1m]))" \ + --epsilon 0.01 --delta 0.001 --out plan.json +``` + +| Option | Meaning | +| --- | --- | +| `--promql ` | A query; repeat for several | +| `--epsilon`, `--delta` | Accuracy for every query; without them the queries are exact | +| `--interval-ms ` | Sample interval (default 15000) | +| `--example ` | A built-in example instead of `--promql`: `planner-layering-1`, `-2`, `-3a`, `-3b`, `-4a` or `-4b` | +| `--max-candidates ` | How many plans the document carries, cheapest first (default 64). Stage 3 still prices every plan. | +| `--out ` | Where to write the document | + +Open the document with **Open stage document…**, by dropping the file on the +page, or, if the file is under `tools/dag-viewer/`, with +`http://127.0.0.1:8000/?doc=`. + +## When something goes wrong + +| Symptom | What to do | +| --- | --- | +| `Could not load the planner output: …` with a list of problems | The document is not a valid `asap-stage-pipeline/v1` document. Write it again with the current `stage_pipeline`. | +| `Planning failed: …` in the editor | The message comes from `stage_pipeline`, for example a query the PromQL frontend does not support. The server's terminal shows the full error. | +| `Planning failed` with a network error or `HTTP 501` | The page was opened without the server (for example as a file, or from another web server). Start `server.py` and open the address it prints. | +| The examples do not change after a planner change | Restart the server with `--regenerate`. | +| `stage_pipeline does not exist; start without --skip-build` | Start the server without `--skip-build`, or build it first: `cargo build -p asap-devtools --bin stage_pipeline` | + +For the document format and the viewer's code, see +[`tools/dag-viewer/README.md`](../../tools/dag-viewer/README.md).