From 4754e52616af203121b552c8bd1d737f8db8b8b2 Mon Sep 17 00:00:00 2001 From: Tim Paine <3105306+timkpaine@users.noreply.github.com> Date: Tue, 29 Sep 2026 17:34:02 -0400 Subject: [PATCH 1/2] Reuse graph node results within one GraphEvaluator evaluation When a node's __call__ calls one of its declared dependencies again, return the result already computed for that graph node instead of evaluating it again, independent of cacheable. Results are released when the evaluation finishes, and calls that are not graph nodes are not retained. Signed-off-by: Tim Paine <3105306+timkpaine@users.noreply.github.com> --- ccflow/evaluators/common.py | 14 ++++++-- ccflow/tests/evaluators/test_common.py | 46 ++++++++++++++++++++++++++ docs/wiki/how-to/Cache-Results.md | 4 ++- docs/wiki/reference/Built-in-Models.md | 34 +++++++++---------- 4 files changed, 78 insertions(+), 20 deletions(-) diff --git a/ccflow/evaluators/common.py b/ccflow/evaluators/common.py index f1dcfbac..8655a78b 100644 --- a/ccflow/evaluators/common.py +++ b/ccflow/evaluators/common.py @@ -486,10 +486,14 @@ def get_dependency_graph(evaluation_context: ModelEvaluationContext) -> Callable class GraphEvaluator(EvaluatorBase): """Evaluator that evaluates the dependency graph of callable models in topologically sorted order. - It is suggested to combine it with a caching evaluator. + Each graph node is evaluated once per graph evaluation. When a node's ``__call__`` calls one of its declared + dependencies again, the result already computed for that graph node is reused, whether or not results are + cacheable. Those results are released when the graph evaluation finishes. Calls that are not graph nodes are + evaluated normally and not retained; combine with a caching evaluator to reuse them. """ _is_evaluating: bool = PrivateAttr(False) + _node_results: dict[bytes, ResultType] = PrivateAttr(default_factory=dict) def is_transparent(self, context: ModelEvaluationContext) -> bool: return True @@ -499,8 +503,12 @@ def __call__(self, context: ModelEvaluationContext) -> ResultType: import graphlib # If we are evaluating deps, or if we have already started using the graph evaluator further up the call tree, - # do not apply it any further + # do not apply it any further beyond reusing results for nodes of the graph being evaluated. if self._is_evaluating: + if self._node_results: + key = _effective_evaluation_key(context) + if key in self._node_results: + return self._node_results[key] return context() self._is_evaluating = True root_result = None @@ -510,8 +518,10 @@ def __call__(self, context: ModelEvaluationContext) -> ResultType: for key in ts.static_order(): evaluation_context = graph.ids[key] result = evaluation_context() + self._node_results[key] = result if key == graph.root_id: root_result = result finally: self._is_evaluating = False + self._node_results = {} return root_result diff --git a/ccflow/tests/evaluators/test_common.py b/ccflow/tests/evaluators/test_common.py index 52990e6e..627ac75d 100644 --- a/ccflow/tests/evaluators/test_common.py +++ b/ccflow/tests/evaluators/test_common.py @@ -6,6 +6,7 @@ import pyarrow as pa from ccflow import ( + CallableModel, DateContext, DateRangeContext, Evaluator, @@ -14,6 +15,7 @@ FlowContext, FlowOptionsOverride, FromContext, + GenericResult, ModelEvaluationContext, NullContext, TransparentModelEvaluationContext, @@ -734,6 +736,50 @@ def test_graph_evaluator_cache(self): self.assertEqual(len(captured.records), (4 + 4) * 3) + def test_graph_evaluator_reuses_node_results_without_cache(self): + """Dependencies called again inside __call__ reuse the graph node result, even when nothing is cacheable.""" + n0 = NodeModel(meta={"name": "n0"}, run_deps=True) + n1 = NodeModel(meta={"name": "n1"}, deps_model=[n0], run_deps=True) + n2 = NodeModel(meta={"name": "n2"}, deps_model=[n0], run_deps=True) + root = NodeModel(meta={"name": "n3"}, deps_model=[n1, n2], run_deps=True) + context = DateContext(date=date(2022, 1, 1)) + + NodeModel._calls = [] + NodeModel._deps_calls = [] + with FlowOptionsOverride(options={"evaluator": GraphEvaluator(), "cacheable": False}): + root(context) + first_calls = list(NodeModel._calls) + root(context) + + self.assertEqual( + sorted(first_calls), sorted([("n0", date(2022, 1, 1)), ("n1", date(2022, 1, 1)), ("n2", date(2022, 1, 1)), ("n3", date(2022, 1, 1))]) + ) + # Node results are released after each graph evaluation, so the second evaluation runs every node again. + self.assertEqual(len(NodeModel._calls), 8) + + def test_graph_evaluator_does_not_retain_calls_outside_the_graph(self): + calls = [] + + class Inner(CallableModel): + @Flow.call + def __call__(self, context: DateContext) -> GenericResult: + calls.append(context.date) + return GenericResult(value=True) + + class Outer(CallableModel): + inner: Inner + + @Flow.call + def __call__(self, context: DateContext) -> GenericResult: + self.inner(context) + self.inner(context) + return GenericResult(value=True) + + with FlowOptionsOverride(options={"evaluator": GraphEvaluator(), "cacheable": False}): + Outer(inner=Inner())(DateContext(date=date(2022, 1, 1))) + + self.assertEqual(calls, [date(2022, 1, 1), date(2022, 1, 1)]) + def test_graph_evaluator_circular(self): root = CircularModel() context = DateContext(date=date(2022, 1, 1)) diff --git a/docs/wiki/how-to/Cache-Results.md b/docs/wiki/how-to/Cache-Results.md index 2574ac9b..59a7d181 100644 --- a/docs/wiki/how-to/Cache-Results.md +++ b/docs/wiki/how-to/Cache-Results.md @@ -77,7 +77,7 @@ Wrapping a model in *transparent* evaluators (logging, timing) does not change i ## Evaluate a dependency graph -To evaluate steps in an optimal order rather than Python's call order, declare dependencies explicitly with `@Flow.deps` and use the `GraphEvaluator` (together with the cache, since graph nodes still run their `__call__` bodies): +To evaluate steps in an optimal order rather than Python's call order, declare dependencies explicitly with `@Flow.deps` and use the `GraphEvaluator`: ```python class FibonacciDepsModel(FibonacciModel): @@ -107,6 +107,8 @@ with FlowOptionsOverride(options={"cacheable": True, "evaluator": evaluator}): Note the topological order (0, 1, 2, 3, 4), and that each node runs once. This is also the foundation for distributed evaluation. +Within one graph evaluation, when a node's `__call__` calls one of its declared dependencies again (as `FibonacciModel` does), the `GraphEvaluator` returns the result it already computed for that node instead of running it again, even when `cacheable` is off. Those results are released when the evaluation finishes. Add a caching evaluator to reuse results across evaluations, or for calls that are not declared as dependencies. + ## Write a custom evaluator No library can provide every execution strategy, so evaluators are extensible. An evaluator is a model that takes a `ModelEvaluationContext` (which carries the model, context, function, and options) and returns a result. Override `is_transparent` to return `True` if it does not change the result (so caching ignores it): diff --git a/docs/wiki/reference/Built-in-Models.md b/docs/wiki/reference/Built-in-Models.md index cc7902c4..092b6232 100644 --- a/docs/wiki/reference/Built-in-Models.md +++ b/docs/wiki/reference/Built-in-Models.md @@ -56,23 +56,23 @@ Publishers (`ccflow.publishers`) are models that write or send data. A common in Evaluators control *how* a `CallableModel` runs. Set one through `FlowOptions` (see [Core Types](Core-Types#flow-options)). Usage is in [Cache Results](Cache-Results) and [Retry on Failure](Retry-on-Failure). -| Name | Path | Description | -| :---------------------------------- | :------------------ | :----------------------------------------------------------------- | -| `LazyEvaluator` | `ccflow.evaluators` | Runs the callable only once an attribute of the result is queried. | -| `LoggingEvaluator` | `ccflow.evaluators` | Logs information about evaluating the callable (the default). | -| `MemoryCacheEvaluator` | `ccflow.evaluators` | Caches results in memory. | -| `MultiEvaluator` | `ccflow.evaluators` | Combines multiple evaluators. | -| `GraphEvaluator` | `ccflow.evaluators` | Evaluates the dependency graph in topologically sorted order. | -| `RetryEvaluator` | `ccflow.evaluators` | Retries evaluation on failure with exponential backoff and jitter. | -| `ChunkedDateRangeEvaluator` | *Coming Soon!* | | -| `ChunkedDateRangeResultsAggregator` | *Coming Soon!* | | -| `DependencyTrackingEvaluator` | *Coming Soon!* | | -| `DiskCacheEvaluator` | *Coming Soon!* | | -| `ParquetCacheEvaluator` | *Coming Soon!* | | -| `RayChunkedDateRangeEvaluator` | *Coming Soon!* | | -| `RayCacheEvaluator` | *Coming Soon!* | | -| `RayGraphEvaluator` | *Coming Soon!* | | -| `RayDelayedDistributedEvaluator` | *Coming Soon!* | | +| Name | Path | Description | +| :---------------------------------- | :------------------ | :--------------------------------------------------------------------------------------------------- | +| `LazyEvaluator` | `ccflow.evaluators` | Runs the callable only once an attribute of the result is queried. | +| `LoggingEvaluator` | `ccflow.evaluators` | Logs information about evaluating the callable (the default). | +| `MemoryCacheEvaluator` | `ccflow.evaluators` | Caches results in memory. | +| `MultiEvaluator` | `ccflow.evaluators` | Combines multiple evaluators. | +| `GraphEvaluator` | `ccflow.evaluators` | Evaluates the dependency graph in topologically sorted order, running each node once per evaluation. | +| `RetryEvaluator` | `ccflow.evaluators` | Retries evaluation on failure with exponential backoff and jitter. | +| `ChunkedDateRangeEvaluator` | *Coming Soon!* | | +| `ChunkedDateRangeResultsAggregator` | *Coming Soon!* | | +| `DependencyTrackingEvaluator` | *Coming Soon!* | | +| `DiskCacheEvaluator` | *Coming Soon!* | | +| `ParquetCacheEvaluator` | *Coming Soon!* | | +| `RayChunkedDateRangeEvaluator` | *Coming Soon!* | | +| `RayCacheEvaluator` | *Coming Soon!* | | +| `RayGraphEvaluator` | *Coming Soon!* | | +| `RayDelayedDistributedEvaluator` | *Coming Soon!* | | ### How cache keys are built From 3a117759e7f9354c7a56c75cce6525c7d6b2c044 Mon Sep 17 00:00:00 2001 From: Tim Paine <3105306+timkpaine@users.noreply.github.com> Date: Tue, 29 Sep 2026 18:44:19 -0400 Subject: [PATCH 2/2] Do not reuse volatile graph node results Volatile nodes are still evaluated in the graph pre-pass, but their results are not stored, so each consumer that calls them recomputes. Signed-off-by: Tim Paine <3105306+timkpaine@users.noreply.github.com> --- ccflow/evaluators/common.py | 7 +++++-- ccflow/tests/evaluators/test_common.py | 20 ++++++++++++++++++++ docs/wiki/how-to/Cache-Results.md | 2 +- 3 files changed, 26 insertions(+), 3 deletions(-) diff --git a/ccflow/evaluators/common.py b/ccflow/evaluators/common.py index 8655a78b..3228d0d6 100644 --- a/ccflow/evaluators/common.py +++ b/ccflow/evaluators/common.py @@ -488,7 +488,7 @@ class GraphEvaluator(EvaluatorBase): Each graph node is evaluated once per graph evaluation. When a node's ``__call__`` calls one of its declared dependencies again, the result already computed for that graph node is reused, whether or not results are - cacheable. Those results are released when the graph evaluation finishes. Calls that are not graph nodes are + cacheable. Volatile nodes are never reused and recompute on every call. Those results are released when the graph evaluation finishes. Calls that are not graph nodes are evaluated normally and not retained; combine with a caching evaluator to reuse them. """ @@ -518,7 +518,10 @@ def __call__(self, context: ModelEvaluationContext) -> ResultType: for key in ts.static_order(): evaluation_context = graph.ids[key] result = evaluation_context() - self._node_results[key] = result + # Volatile nodes always recompute, so never share their result with nested calls. + inner, _, _ = _flatten_cache_key_context(evaluation_context) + if not inner.options.get("volatile"): + self._node_results[key] = result if key == graph.root_id: root_result = result finally: diff --git a/ccflow/tests/evaluators/test_common.py b/ccflow/tests/evaluators/test_common.py index 627ac75d..16826c7d 100644 --- a/ccflow/tests/evaluators/test_common.py +++ b/ccflow/tests/evaluators/test_common.py @@ -757,6 +757,26 @@ def test_graph_evaluator_reuses_node_results_without_cache(self): # Node results are released after each graph evaluation, so the second evaluation runs every node again. self.assertEqual(len(NodeModel._calls), 8) + def test_graph_evaluator_recomputes_volatile_nodes_for_each_consumer(self): + n0 = NodeModel(meta={"name": "n0"}) + n1 = NodeModel(meta={"name": "n1"}, deps_model=[n0], run_deps=True) + n2 = NodeModel(meta={"name": "n2"}, deps_model=[n0], run_deps=True) + root = NodeModel(meta={"name": "n3"}, deps_model=[n1, n2]) + context = DateContext(date=date(2022, 1, 1)) + + NodeModel._calls = [] + NodeModel._deps_calls = [] + with ( + FlowOptionsOverride(options={"evaluator": GraphEvaluator(), "cacheable": False}), + FlowOptionsOverride(options={"volatile": True}, models=(n0,)), + ): + root(context) + + # Once in the graph pre-pass, then once more for each consumer that calls it. + self.assertEqual(NodeModel._calls.count(("n0", date(2022, 1, 1))), 3) + self.assertEqual(NodeModel._calls.count(("n1", date(2022, 1, 1))), 1) + self.assertEqual(NodeModel._calls.count(("n2", date(2022, 1, 1))), 1) + def test_graph_evaluator_does_not_retain_calls_outside_the_graph(self): calls = [] diff --git a/docs/wiki/how-to/Cache-Results.md b/docs/wiki/how-to/Cache-Results.md index 59a7d181..b51da754 100644 --- a/docs/wiki/how-to/Cache-Results.md +++ b/docs/wiki/how-to/Cache-Results.md @@ -107,7 +107,7 @@ with FlowOptionsOverride(options={"cacheable": True, "evaluator": evaluator}): Note the topological order (0, 1, 2, 3, 4), and that each node runs once. This is also the foundation for distributed evaluation. -Within one graph evaluation, when a node's `__call__` calls one of its declared dependencies again (as `FibonacciModel` does), the `GraphEvaluator` returns the result it already computed for that node instead of running it again, even when `cacheable` is off. Those results are released when the evaluation finishes. Add a caching evaluator to reuse results across evaluations, or for calls that are not declared as dependencies. +Within one graph evaluation, when a node's `__call__` calls one of its declared dependencies again (as `FibonacciModel` does), the `GraphEvaluator` returns the result it already computed for that node instead of running it again, even when `cacheable` is off. Nodes marked `volatile` are never reused and recompute on every call. Those results are released when the evaluation finishes. Add a caching evaluator to reuse results across evaluations, or for calls that are not declared as dependencies. ## Write a custom evaluator