From 8430b8ba783d8bc05e81267d298386458033bbc9 Mon Sep 17 00:00:00 2001 From: Beto Dealmeida Date: Mon, 5 Oct 2026 10:06:00 -0700 Subject: [PATCH 1/5] perf: instrument and speed up deployment fingerprints A shared.main dry-run for semantic-shared PR #255 spent 340 seconds building fingerprints for 6,663 downstream impacts. Log time spent loading targets and external ancestors, converting specs, finding ancestor closures, and building current and proposed fingerprint graphs. Avoid deep-copying external target specs between graph snapshots. Clear their current values from the shared cache before calculating proposed values so changed ancestors still produce fresh hashes. Reuse parsed query ASTs when rendered specs are fingerprinted instead of parsing the same SQL on each pass. Add a repeatable 6,000-metric benchmark and regression tests. The synthetic run fell from about 4.0 to 2.1 seconds, with SQL parses falling from 18,000 to 6,000. Database access is mocked, so the production speedup is not yet measured. Tests: 30 deployment fingerprint tests and 44 model fingerprint tests passed. --- .../internal/deployment/fingerprints.py | 74 ++++++++++++++-- .../semantic_fingerprints/v1.py | 5 ++ .../scripts/benchmark_fingerprints.py | 88 +++++++++++++++++++ .../internal/deployment/fingerprints_test.py | 60 ++++++++++++- 4 files changed, 220 insertions(+), 7 deletions(-) create mode 100644 datajunction-server/scripts/benchmark_fingerprints.py diff --git a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py index b78ad720c..dbb739f30 100644 --- a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py +++ b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py @@ -1,4 +1,5 @@ import logging +import time from collections.abc import Callable, Iterable from heapq import heappop, heappush @@ -289,11 +290,18 @@ async def _load_external_specs( parent_cache: ParentCandidateCache, ) -> dict[str, NodeSpec]: """Transitively load names unresolved within the deployment's own specs.""" + started = time.perf_counter() + db_elapsed = 0.0 + conversion_elapsed = 0.0 + parent_elapsed = 0.0 + rounds = 0 known_names = set(known_names) external_specs: dict[str, NodeSpec] = {} frontier = sorted(set(initial_frontier) - known_names) while frontier: + rounds += 1 known_names.update(frontier) + step_started = time.perf_counter() nodes = await Node.get_by_names( session, frontier, @@ -305,8 +313,12 @@ async def _load_external_specs( selectinload(Node.owners), ], ) + db_elapsed += time.perf_counter() - step_started + step_started = time.perf_counter() pending = [await node.to_spec(session) for node in nodes] + conversion_elapsed += time.perf_counter() - step_started external_specs.update({spec.rendered_name: spec for spec in pending}) + step_started = time.perf_counter() candidates: set[str] = set() for spec in pending: try: @@ -319,6 +331,18 @@ async def _load_external_specs( exc, ) frontier = sorted(candidates - known_names) + parent_elapsed += time.perf_counter() - step_started + if rounds: + logger.info( + "External fingerprint specs: %d nodes in %d rounds; " + "db=%0.fms, conversion=%0.fms, parents=%0.fms, total=%0.fms", + len(external_specs), + rounds, + db_elapsed * 1000, + conversion_elapsed * 1000, + parent_elapsed * 1000, + (time.perf_counter() - started) * 1000, + ) return external_specs @@ -389,6 +413,7 @@ def _compute_merkle_fingerprints( shared_fingerprints: dict[int, SemanticFingerprintValue] | None = None, ) -> FingerprintMap: """`shared_fingerprints` is a cache keyed by id(spec), shared across snapshots.""" + started = time.perf_counter() parent_cache = parent_cache if parent_cache is not None else {} shared_fingerprints = shared_fingerprints if shared_fingerprints is not None else {} names_to_process = sorted(specs) if only_names is None else sorted(only_names) @@ -433,6 +458,7 @@ def _compute_merkle_fingerprints( if name not in ignored_parse_errors: logger.warning("Fingerprint unavailable for %s: %s", name, exc) + graph_elapsed = time.perf_counter() - started components = strongly_connected_components(graph) component_by_name = { name: component_index @@ -465,6 +491,7 @@ def _compute_merkle_fingerprints( if remaining == 0: heappush(ready, component_index) + scc_elapsed = time.perf_counter() - started - graph_elapsed processed = 0 while ready: component_index = heappop(ready) @@ -555,6 +582,16 @@ def _compute_merkle_fingerprints( if processed != len(components): # pragma: no cover raise RuntimeError("SCC condensation graph contains a cycle") + logger.info( + "Semantic fingerprint graph: %d nodes, %d cached, %d components; " + "parents=%0.fms, components=%0.fms, hashing=%0.fms", + len(names_to_process), + len(cached_results), + len(components), + graph_elapsed * 1000, + scc_elapsed * 1000, + (time.perf_counter() - started - graph_elapsed - scc_elapsed) * 1000, + ) return { name: component_results[component_by_name[name]][name] for name in names_to_process @@ -637,6 +674,8 @@ async def build_deployment_fingerprints( pre_parsed_queries: dict[str, ReusableQuery] | None = None, ) -> tuple[FingerprintMap, FingerprintMap]: """`only_proposed_names` limits which nodes get a fresh proposed hash.""" + started = time.perf_counter() + timings: dict[str, float] = {} proposed_specs = list(proposed_specs) additional_target_names = set(additional_target_names) deleted_names = {spec.rendered_name for spec in deleted_specs} @@ -652,11 +691,13 @@ async def build_deployment_fingerprints( deleted_names, unchanged_names, ) + timings["setup"] = time.perf_counter() - started target_specs: dict[str, NodeSpec] = {} target_names_to_load = ( additional_target_names - existing_specs.keys() - submitted_names ) if target_names_to_load: + step_started = time.perf_counter() target_nodes = await Node.get_by_names( session, sorted(target_names_to_load), @@ -668,11 +709,15 @@ async def build_deployment_fingerprints( selectinload(Node.owners), ], ) + timings["target_db"] = time.perf_counter() - step_started + step_started = time.perf_counter() target_specs = { spec.rendered_name: spec for spec in [await node.to_spec(session) for node in target_nodes] } + timings["target_to_spec"] = time.perf_counter() - step_started + step_started = time.perf_counter() parent_cache: ParentCandidateCache = {} # Reuse ASTs propagate_impact already parsed. current_seed_specs = {**existing_specs, **target_specs} @@ -693,6 +738,8 @@ async def build_deployment_fingerprints( proposed, parent_cache, ) + timings["ancestor_closure"] = time.perf_counter() - step_started + step_started = time.perf_counter() external = await _load_external_specs( session, current_unresolved | proposed_unresolved, @@ -700,13 +747,10 @@ async def build_deployment_fingerprints( ignored_parse_errors=deleted_names, parent_cache=parent_cache, ) - # target_specs need a separate copy per side so the id()-keyed - # shared_fingerprints cache can't cross-reuse a stale ancestor chain. + timings["external_load"] = time.perf_counter() - step_started + step_started = time.perf_counter() current_external = {**external, **target_specs} - proposed_external = { - **external, - **{name: spec.model_copy(deep=True) for name, spec in target_specs.items()}, - } + proposed_external = {**external, **target_specs} shared_fingerprints: dict[int, SemanticFingerprintValue] = {} current_graph = SemanticFingerprintGraph( {**current_external, **existing_specs}, @@ -721,9 +765,27 @@ async def build_deployment_fingerprints( version=version, shared_fingerprints=shared_fingerprints, ) + timings["snapshot_setup"] = time.perf_counter() - step_started + step_started = time.perf_counter() # unchanged_names gives no-op nodes a current fingerprint to fall back on. current = current_graph.fingerprints( deleted_names | additional_target_names | unchanged_names, ) + timings["current_hashes"] = time.perf_counter() - step_started + # A target can have different ancestors in the proposed graph. Drop its + # current hash from the shared cache instead of deep-copying every target. + for spec in target_specs.values(): + shared_fingerprints.pop(id(spec), None) + step_started = time.perf_counter() proposed_hashes = proposed_graph.fingerprints(proposed_target_names) + timings["proposed_hashes"] = time.perf_counter() - step_started + logger.info( + "Deployment fingerprints: %d targets (%d loaded), %d external ancestors; " + "total=%.0fms, %s", + len(additional_target_names), + len(target_specs), + len(external), + (time.perf_counter() - started) * 1000, + ", ".join(f"{name}={elapsed * 1000:.0f}ms" for name, elapsed in timings.items()), + ) return current, proposed_hashes diff --git a/datajunction-server/datajunction_server/semantic_fingerprints/v1.py b/datajunction-server/datajunction_server/semantic_fingerprints/v1.py index 58ddcdc9b..88ddf745b 100644 --- a/datajunction-server/datajunction_server/semantic_fingerprints/v1.py +++ b/datajunction-server/datajunction_server/semantic_fingerprints/v1.py @@ -69,6 +69,11 @@ def build_fingerprint( """Build a version 1 semantic fingerprint.""" fingerprint_fields = semantic_fields(type(spec)) rendered = spec.rendered_spec() + # rendered_spec() reconstructs the model and drops private attributes. + # Both specs have the same rendered query, so reuse its parsed AST across + # current and proposed fingerprints instead of parsing each copy again. + if rendered.rendered_query is not None: + rendered._query_ast = spec.query_ast fields = { field: normalize_field( rendered, diff --git a/datajunction-server/scripts/benchmark_fingerprints.py b/datajunction-server/scripts/benchmark_fingerprints.py new file mode 100644 index 000000000..66fd10f3e --- /dev/null +++ b/datajunction-server/scripts/benchmark_fingerprints.py @@ -0,0 +1,88 @@ +"""Profile fingerprinting many downstream metrics without a database. + +Run from the DJ repository with the server package installed: + python datajunction-server/scripts/benchmark_fingerprints.py 6000 --profile + +The mocked node lookup keeps this focused on Python work. It does not model +production database loading or the complexity of real metric definitions. +""" + +import argparse +import asyncio +import cProfile +import logging +import pstats +import time +from unittest.mock import AsyncMock, MagicMock, patch + +from datajunction_server.internal.deployment.fingerprints import ( + build_deployment_fingerprints, +) +from datajunction_server.models.deployment import MetricSpec, SourceSpec + + +async def run(count: int, profile: bool) -> None: + """Build current and proposed fingerprints for ``count`` external metrics.""" + before = SourceSpec( + namespace="ns", + name="source", + catalog="c", + schema_="s", + table="before", + ) + after = before.model_copy(update={"table": "after"}) + metrics = [ + MetricSpec( + namespace="other", + name=f"metric_{index}", + query="SELECT COUNT(*) FROM ns.source", + ) + for index in range(count) + ] + nodes = [] + for metric in metrics: + node = MagicMock() + node.to_spec = AsyncMock(return_value=metric) + nodes.append(node) + + profiler = cProfile.Profile() if profile else None + started = time.perf_counter() + with patch( + "datajunction_server.internal.deployment.fingerprints.Node.get_by_names", + AsyncMock(return_value=nodes), + ): + if profiler is not None: + profiler.enable() + current, proposed = await build_deployment_fingerprints( + MagicMock(), + {before.rendered_name: before}, + [after], + [], + additional_target_names=[metric.rendered_name for metric in metrics], + ) + if profiler is not None: + profiler.disable() + + print( + f"count={count} elapsed={time.perf_counter() - started:.3f}s " + f"current={len(current)} proposed={len(proposed)}", + ) + if profiler is not None: + pstats.Stats(profiler).strip_dirs().sort_stats("cumulative").print_stats(35) + + +def main() -> None: + """Parse CLI arguments and run the benchmark.""" + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("count", type=int, help="number of downstream metrics") + parser.add_argument("--profile", action="store_true", help="show cProfile output") + args = parser.parse_args() + logging.basicConfig(level=logging.WARNING, format="%(message)s") + logging.getLogger("datajunction_server.internal.deployment.fingerprints").setLevel( + logging.INFO, + ) + asyncio.run(run(args.count, args.profile)) + + +if __name__ == "__main__": + main() diff --git a/datajunction-server/tests/internal/deployment/fingerprints_test.py b/datajunction-server/tests/internal/deployment/fingerprints_test.py index 65a254d72..74dca8cb3 100644 --- a/datajunction-server/tests/internal/deployment/fingerprints_test.py +++ b/datajunction-server/tests/internal/deployment/fingerprints_test.py @@ -30,7 +30,11 @@ UNKNOWN_SEMANTIC_FINGERPRINT, SemanticFingerprint, ) -from datajunction_server.semantic_fingerprints.engine import local_node_fingerprint +from datajunction_server.semantic_fingerprints.engine import ( + compose_node_fingerprint, + local_node_fingerprint, +) +from datajunction_server.sql.parsing.backends.antlr4 import parse def source_spec( @@ -594,6 +598,42 @@ async def test_build_deployment_fingerprints_only_proposed_names_reuses_unchange assert scoped == {downstream.rendered_name: full[downstream.rendered_name]} +@pytest.mark.asyncio +async def test_external_target_is_rehashed_when_its_ancestor_changes(): + """An external target needs separate current and proposed graph hashes.""" + current_source = source_spec("source", table="before") + proposed_source = source_spec("source", table="after") + downstream = transform_spec( + "downstream", + "SELECT * FROM ns.source", + namespace="other", + ) + node = MagicMock() + node.to_spec = AsyncMock(return_value=downstream) + + with patch( + "datajunction_server.internal.deployment.fingerprints.Node.get_by_names", + AsyncMock(return_value=[node]), + ): + current, proposed = await build_deployment_fingerprints( + MagicMock(), + {current_source.rendered_name: current_source}, + [proposed_source], + [], + additional_target_names=[downstream.rendered_name], + ) + + expected_current = SemanticFingerprintGraph( + spec_map(current_source, downstream), + ).fingerprint(downstream.rendered_name) + expected_proposed = SemanticFingerprintGraph( + spec_map(proposed_source, downstream), + ).fingerprint(downstream.rendered_name) + assert current[downstream.rendered_name] == expected_current + assert proposed[downstream.rendered_name] == expected_proposed + assert expected_current != expected_proposed + + @pytest.mark.asyncio async def test_build_deployment_fingerprints_without_external_parents(): source = source_spec("source", table="table") @@ -808,3 +848,21 @@ def test_semantic_fingerprint_graph_caches_whole_graph_result(): # A subsequent scoped call also reuses the already-cached result. scoped = graph.fingerprints({transform.rendered_name}) assert scoped[transform.rendered_name] == first[transform.rendered_name] + + +def test_fingerprinting_reuses_parsed_query_across_rendered_copies(): + metric = MetricSpec( + namespace="ns", + name="metric", + query="SELECT COUNT(*) FROM ${prefix}source", + ) + + with patch( + "datajunction_server.sql.parsing.backends.antlr4.parse", + wraps=parse, + ) as parse_query: + first = compose_node_fingerprint(metric, parent_fingerprints=[]) + second = compose_node_fingerprint(metric, parent_fingerprints=[]) + + assert first == second + assert parse_query.call_count == 1 From 251d1f1ccc0bae6ef68d56debd6007a6246c381a Mon Sep 17 00:00:00 2001 From: Beto Dealmeida Date: Mon, 5 Oct 2026 10:09:24 -0700 Subject: [PATCH 2/5] chore: keep one-off fingerprint benchmark out of tree The benchmark uses mocked database access and a synthetic graph. Keep its measurements in the PR description while retaining the production timing logs and regression tests in the source tree. --- .../scripts/benchmark_fingerprints.py | 88 ------------------- 1 file changed, 88 deletions(-) delete mode 100644 datajunction-server/scripts/benchmark_fingerprints.py diff --git a/datajunction-server/scripts/benchmark_fingerprints.py b/datajunction-server/scripts/benchmark_fingerprints.py deleted file mode 100644 index 66fd10f3e..000000000 --- a/datajunction-server/scripts/benchmark_fingerprints.py +++ /dev/null @@ -1,88 +0,0 @@ -"""Profile fingerprinting many downstream metrics without a database. - -Run from the DJ repository with the server package installed: - python datajunction-server/scripts/benchmark_fingerprints.py 6000 --profile - -The mocked node lookup keeps this focused on Python work. It does not model -production database loading or the complexity of real metric definitions. -""" - -import argparse -import asyncio -import cProfile -import logging -import pstats -import time -from unittest.mock import AsyncMock, MagicMock, patch - -from datajunction_server.internal.deployment.fingerprints import ( - build_deployment_fingerprints, -) -from datajunction_server.models.deployment import MetricSpec, SourceSpec - - -async def run(count: int, profile: bool) -> None: - """Build current and proposed fingerprints for ``count`` external metrics.""" - before = SourceSpec( - namespace="ns", - name="source", - catalog="c", - schema_="s", - table="before", - ) - after = before.model_copy(update={"table": "after"}) - metrics = [ - MetricSpec( - namespace="other", - name=f"metric_{index}", - query="SELECT COUNT(*) FROM ns.source", - ) - for index in range(count) - ] - nodes = [] - for metric in metrics: - node = MagicMock() - node.to_spec = AsyncMock(return_value=metric) - nodes.append(node) - - profiler = cProfile.Profile() if profile else None - started = time.perf_counter() - with patch( - "datajunction_server.internal.deployment.fingerprints.Node.get_by_names", - AsyncMock(return_value=nodes), - ): - if profiler is not None: - profiler.enable() - current, proposed = await build_deployment_fingerprints( - MagicMock(), - {before.rendered_name: before}, - [after], - [], - additional_target_names=[metric.rendered_name for metric in metrics], - ) - if profiler is not None: - profiler.disable() - - print( - f"count={count} elapsed={time.perf_counter() - started:.3f}s " - f"current={len(current)} proposed={len(proposed)}", - ) - if profiler is not None: - pstats.Stats(profiler).strip_dirs().sort_stats("cumulative").print_stats(35) - - -def main() -> None: - """Parse CLI arguments and run the benchmark.""" - parser = argparse.ArgumentParser(description=__doc__) - parser.add_argument("count", type=int, help="number of downstream metrics") - parser.add_argument("--profile", action="store_true", help="show cProfile output") - args = parser.parse_args() - logging.basicConfig(level=logging.WARNING, format="%(message)s") - logging.getLogger("datajunction_server.internal.deployment.fingerprints").setLevel( - logging.INFO, - ) - asyncio.run(run(args.count, args.profile)) - - -if __name__ == "__main__": - main() From b1751bbc3d2df43e80c5cdd814216f2eeaa1ae50 Mon Sep 17 00:00:00 2001 From: Beto Dealmeida Date: Mon, 5 Oct 2026 14:59:21 -0700 Subject: [PATCH 3/5] docs: explain fingerprint cache invalidation Describe why a downstream target must be rehashed against proposed ancestors when the shared cache is keyed by spec identity. --- .../datajunction_server/internal/deployment/fingerprints.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py index dbb739f30..2988a7963 100644 --- a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py +++ b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py @@ -772,8 +772,9 @@ async def build_deployment_fingerprints( deleted_names | additional_target_names | unchanged_names, ) timings["current_hashes"] = time.perf_counter() - step_started - # A target can have different ancestors in the proposed graph. Drop its - # current hash from the shared cache instead of deep-copying every target. + # The shared cache is keyed only by spec identity. A target can keep the + # same spec object but have a new hash when its proposed ancestors change, + # so evict its current hash before evaluating the proposed graph. for spec in target_specs.values(): shared_fingerprints.pop(id(spec), None) step_started = time.perf_counter() From 6cb7b11028e0b9eb1b96c8fbf4bf78e926243811 Mon Sep 17 00:00:00 2001 From: Beto Dealmeida Date: Mon, 5 Oct 2026 15:01:16 -0700 Subject: [PATCH 4/5] Cleanup --- .../internal/deployment/fingerprints.py | 7 ++++--- .../datajunction_server/semantic_fingerprints/v1.py | 1 + 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py index 2988a7963..2f405e3a2 100644 --- a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py +++ b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py @@ -291,13 +291,12 @@ async def _load_external_specs( ) -> dict[str, NodeSpec]: """Transitively load names unresolved within the deployment's own specs.""" started = time.perf_counter() - db_elapsed = 0.0 - conversion_elapsed = 0.0 - parent_elapsed = 0.0 + db_elapsed = conversion_elapsed = parent_elapsed = 0.0 rounds = 0 known_names = set(known_names) external_specs: dict[str, NodeSpec] = {} frontier = sorted(set(initial_frontier) - known_names) + while frontier: rounds += 1 known_names.update(frontier) @@ -332,6 +331,7 @@ async def _load_external_specs( ) frontier = sorted(candidates - known_names) parent_elapsed += time.perf_counter() - step_started + if rounds: logger.info( "External fingerprint specs: %d nodes in %d rounds; " @@ -343,6 +343,7 @@ async def _load_external_specs( parent_elapsed * 1000, (time.perf_counter() - started) * 1000, ) + return external_specs diff --git a/datajunction-server/datajunction_server/semantic_fingerprints/v1.py b/datajunction-server/datajunction_server/semantic_fingerprints/v1.py index 88ddf745b..1704e533b 100644 --- a/datajunction-server/datajunction_server/semantic_fingerprints/v1.py +++ b/datajunction-server/datajunction_server/semantic_fingerprints/v1.py @@ -74,6 +74,7 @@ def build_fingerprint( # current and proposed fingerprints instead of parsing each copy again. if rendered.rendered_query is not None: rendered._query_ast = spec.query_ast + fields = { field: normalize_field( rendered, From d5012b54d4fc6b046bef38ec6a83b85adae212f9 Mon Sep 17 00:00:00 2001 From: Beto Dealmeida Date: Mon, 5 Oct 2026 15:09:33 -0700 Subject: [PATCH 5/5] style: format fingerprint timing log --- .../datajunction_server/internal/deployment/fingerprints.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py index 2f405e3a2..ba4e3e055 100644 --- a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py +++ b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py @@ -788,6 +788,8 @@ async def build_deployment_fingerprints( len(target_specs), len(external), (time.perf_counter() - started) * 1000, - ", ".join(f"{name}={elapsed * 1000:.0f}ms" for name, elapsed in timings.items()), + ", ".join( + f"{name}={elapsed * 1000:.0f}ms" for name, elapsed in timings.items() + ), ) return current, proposed_hashes