Repository navigation
Design and merge plan: operator output port result cache (MVP) #5880
Replies: 6 comments 1 reply
|
@Xiao-zhen-Liu Thanks for the great summary. Please also describe our plan to manage the lifecycle of the cached results. |
|
@chenlica, here is the plan for managing the lifecycle of cached results. Lifecycle of a cached result Creation. An entry is written when a port's result is materialized during a run, as part of saving that result. It records the cache key of the port's upstream computation, the result's storage location, and the source execution. Reuse / invalidation. On a later run, a port whose upstream computation is unchanged matches by cache key and reads the stored result instead of recomputing; entries persist across runs. A changed upstream (parameter, wiring, schema, operator version) yields a new cache key and a new entry; old entries are never overwritten, they just stop being matched. Where results live. Materialized results are stored in the shared result storage (Iceberg via the REST catalog plus object store), addressed by workflow/execution/operator/port. This storage is process-wide and independent of any computing unit, so cached results survive a computing unit being destroyed. (This supersedes the earlier assumption that results are tied to the computing unit; after the storage-from-compute separation, they are not.) Removal. Cleanup is driven by the workflow lifecycle, not the computing unit. The direction is to reuse
Either way, cache entries are also removed when the workflow is deleted (the table's foreign key to Out of scope (MVP). No size-bounded or cost-based eviction; that is the cost-model work and is intentionally not in this merge. |
|
Does the cache need a new operator state? In the #6729 review, @Yicong-Huang asked whether the CACHE_REUSED state that PR introduced needs to exist at all. Answering here since it is a design decision rather than a PR detail. It does not. The state was never the goal, only a means: the goal is showing the user which operators and ports were reused, and nobody disagrees with that. The question is only how to carry that information. A state is the wrong carrier. Checking every consumer: nothing behaves differently for a reused operator than for a completed one, and every place that checks state expects a finished operator to be COMPLETED:
So a reused operator reports COMPLETED, and the reuse information travels as one boolean:
#6729 is reworked down to exactly this: the flag and the code that carries it, 6 files. The diff there is the concrete proposal. |
|
Design update: planning cache reuse before scheduling The first three PRs of the plan above are merged: the cache key (#5966), the
Nothing sets a run's matched results until the lookup PR, so all four leave every run on main unchanged. Below is the design they introduce, and what changed from the post above. Important Assumption: full reuse. The design always reuses a usable saved result. Every retained operator that reads a port with a usable saved result reads that saved result, and an operator is skipped whenever usable saved results cover every port the run needs from it. There is no cost-based choice between reading a saved result and computing it again, even where computing would be cheaper. That choice (cost-based reuse planning) stays future work, as in the post above. Deciding what a run skipsA run's matched results are the saved results found for its output ports (
The example below applies the rules to a join. A, B and G are viewed; B and D have saved results. B and D are skipped, and C-Build and E read their saved results. A runs because its viewed result has no saved copy, and the link from A into the skipped B is skipped.
Two changes from the post above:
A saved result is unusable, and the run just does not reuse it, when its port is not in the plan or its location is not that port's in this workflow and this run's warehouse. A plan with a loop reuses nothing. An unusable result never fails a run. Scheduling the retained part
A schedule for the example (the generated regions are the ones the search would typically pick):
With no matched results, the scheduler runs today's code: the whole plan goes to the generator. |
|
Design discussion of October 1: the open PRs are on hold On October 1, @chenlica, @Yicong-Huang, @parshimers, @sshiv012 and I went over the design update above. The outcome affects how the work goes on, so I am posting it here. Under full reuse (what the open PRs build). We discussed what a skipped operator and a skip region should do, and whether they add new region states. As I understood it, @Yicong-Huang's main requirement was that skipping should not add new region states, which the current design already meets. Beyond full reuse. The main question was whether the current two-layer design can be extended past full reuse in a way that is compatible with other work. In the two-layer design, the skeleton decides what runs and then
In my view, both designs explore the same space of plans. We agreed that the two-layer design cannot be extended into the one-layer form for the general setting. Outcome. There was no consensus that the one-layer alternative has been investigated enough to rule it out. Until it has, #8731, #8732, #8733 and #8734 are on hold, and the rest of the series waits. I will find out whether the one-layer approach works and how much complexity it adds to the code base, and whether the two-layer design can be generalized in a way that is compatible with the other work. I will discuss both with @chenlica and post the findings here. Background on how |




Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Design and merge plan: operator output port result cache (MVP)
In current Texera main, the engine runs a workflow from the start every time, even when the user changed only one operator near the end. This proposal adds a result cache so that, on a re-run, an output port whose upstream computation logic is unchanged reads its saved result instead of recomputing it. The code is written and working on a prototype branch. This post describes the design and the plan to bring it into
mainas small PRs, so anyone can raise concerns before the PRs go up.Matching results across executions
Each output port has a cache key built from its upstream operators, their parameters, their output schemas, and the wiring between them. Two ports with the same cache key produce the same result (output port equivalence). When a run saves a port's result, we record
(workflow, port, cache key) -> result location. On a later run, a port whose cache key has a recorded result is a matched port, and that result can be reused. Any edit upstream of a port changes its cache key, so its old result is no longer matched.Scope (MVP)
In scope: reuse the saved result at every matched port (full reuse), match by cache key, invalidate entries that no longer match after an edit, and read, write, and clear the cache from the UI.
Out of scope (future work, not in these PRs): choosing per port whether reuse is cheaper than recompute (cost-based reuse planning), and removing results under storage limits (eviction). The merged code always reuses a matched port's result.
How it fits the current system
Current main (figure below): the Workflow Compiler builds a physical plan,
CostBasedScheduleGeneratorbuilds a schedule of regions, and the executor runs the regions on workers that read and write tables in storage.With the cache MVP (figure below), the engine includes additional modules:
CostBasedScheduleGenerator): schedule the run-skeleton as it does today. The removed part becomes regions that are skipped, and operators that read from it use the saved result locations.The executor saves results to the cached-result storage as ports finish.
Nothing changes when there are no matched ports
On the first run of a workflow, or any run right after an upstream edit, there are no matched ports: the run-skeleton is the whole plan and the schedule is the same as today. The cache changes behavior only when a matched port exists, so the code can land inactive and turn on once results are saved.
Merge plan: five PRs
In dependency order; each has its own issue:
PRs 1 and 2 are independent. PR 3 needs 1 and 2; PR 4 needs 3; PR 5 needs 4.
All reactions