diff --git a/ccflow/evaluators/common.py b/ccflow/evaluators/common.py index f1dcfbac..3228d0d6 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. 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. """ _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,13 @@ def __call__(self, context: ModelEvaluationContext) -> ResultType: for key in ts.static_order(): evaluation_context = graph.ids[key] result = evaluation_context() + # 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: 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..16826c7d 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,70 @@ 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_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 = [] + + 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..b51da754 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. 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 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