From db63b1cb79d89808611c0664e5192a611a99ddeb Mon Sep 17 00:00:00 2001 From: georgeRobertson <50412379+georgeRobertson@users.noreply.github.com> Date: Mon, 28 Sep 2026 12:24:50 +0100 Subject: [PATCH 1/3] perf: only check for orphaned records on entities with at least one record --- .../backends/implementations/duckdb/rules.py | 17 +++++++++++------ 1 file changed, 11 insertions(+), 6 deletions(-) diff --git a/src/dve/core_engine/backends/implementations/duckdb/rules.py b/src/dve/core_engine/backends/implementations/duckdb/rules.py index dbbbc54..9e58a1f 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/rules.py +++ b/src/dve/core_engine/backends/implementations/duckdb/rules.py @@ -29,6 +29,7 @@ duckdb_write_parquet, get_all_registered_udfs, get_duckdb_type_from_annotation, + relation_is_empty, ) from dve.core_engine.backends.implementations.duckdb.types import ( DuckDBEntities, @@ -396,6 +397,10 @@ def identify_orphans( target_rel: DuckDBPyRelation = entities[config.target_name] target_rel = target_rel.set_alias(config.target_name) + if relation_is_empty(source_rel): + self.logger.info(f"{config.entity_name} is empty. Skipping orphan check.") + return [], 0 + match_name = f"matched_{uuid4().hex}" target_rel = target_rel.select( StarExpression(exclude=[]), ConstantExpression(1).alias(match_name) @@ -420,16 +425,16 @@ def identify_orphans( _orph_records: tuple[int] = orphaned_rel.count(RECORD_INDEX_COLUMN_NAME).fetchone() # type: ignore # pylint: disable=C0301 if _orph_records: _no_orphans = _orph_records[0] + if entities.get(ORPHANED_RECORD_ENTITY_NAME) is not None: + entities[ORPHANED_RECORD_ENTITY_NAME] = entities[ORPHANED_RECORD_ENTITY_NAME].union( + orphaned_rel + ) + else: + entities[ORPHANED_RECORD_ENTITY_NAME] = orphaned_rel else: _no_orphans = 0 self.logger.info(f"Found {_no_orphans} orphaned records in {config.entity_name}.") - if entities.get(ORPHANED_RECORD_ENTITY_NAME) is not None: - entities[ORPHANED_RECORD_ENTITY_NAME] = entities[ORPHANED_RECORD_ENTITY_NAME].union( - orphaned_rel - ) - else: - entities[ORPHANED_RECORD_ENTITY_NAME] = orphaned_rel return [], _no_orphans def remove_orphans(self, entities: DuckDBEntities, *, config: OrphanRemoval) -> Iterator: From 872972b7d23ab7e9fc70371ecef37800353b7a71 Mon Sep 17 00:00:00 2001 From: georgeRobertson <50412379+georgeRobertson@users.noreply.github.com> Date: Mon, 28 Sep 2026 15:13:16 +0100 Subject: [PATCH 2/3] fix: add the ability to cache entities in duckdb to hopefully improve performance on orphan checks --- .../backends/implementations/duckdb/rules.py | 61 +++++++++++++++++-- 1 file changed, 57 insertions(+), 4 deletions(-) diff --git a/src/dve/core_engine/backends/implementations/duckdb/rules.py b/src/dve/core_engine/backends/implementations/duckdb/rules.py index 9e58a1f..7753393 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/rules.py +++ b/src/dve/core_engine/backends/implementations/duckdb/rules.py @@ -1,7 +1,7 @@ """Business rule definitions for duckdb backend""" from collections.abc import Callable, Iterator -from typing import get_type_hints +from typing import get_type_hints, Optional from uuid import uuid4 from duckdb import ( @@ -56,11 +56,12 @@ SemiJoin, TableUnion, ) +from dve.core_engine.configuration.v1.hierarchy import EntityHierarchy from dve.core_engine.constants import ORPHANED_RECORD_ENTITY_NAME, RECORD_INDEX_COLUMN_NAME from dve.core_engine.functions import implementations as functions from dve.core_engine.message import FeedbackMessage from dve.core_engine.templating import template_object -from dve.core_engine.type_hints import Messages +from dve.core_engine.type_hints import URI, EntityName, Messages @duckdb_record_index @@ -73,6 +74,7 @@ class DuckDBStepImplementations(BaseStepImplementations[DuckDBPyRelation]): def __init__(self, connection: DuckDBPyConnection, **kwargs): self._connection = connection self.registered_functions = get_all_registered_udfs(self._connection) + self._temp_tables: dict[EntityName, str] = {} super().__init__(**kwargs) @property @@ -113,6 +115,29 @@ def register_udfs( # type: ignore connection.sql(_sql) return cls(connection=connection, **kwargs) + def _materialise_temp_tables_from_entity( + self, + entities: DuckDBEntities, + entity_names: list[EntityName], + refresh: bool = False + ): + """Materialise an entity into a temporary duckdb table.""" + for entity_name in entity_names: + temp_name = f"temp_{entity_name}" + if entity_name in self._temp_tables: + if not refresh: + continue + self.connection.unregister(temp_name) + del self._temp_tables[entity_name] + + self.connection.register(temp_name, entities[entity_name]) + self._temp_tables[entity_name] = temp_name + + def _drop_temp_tables(self) -> None: + for temp_name in self._temp_tables.values(): + self.connection.unregister(temp_name) + self._temp_tables.clear() + def add(self, entities: DuckDBEntities, *, config: ColumnAddition) -> Messages: """A transformation step which adds a column to an entity.""" entity: DuckDBPyRelation = entities[config.entity_name] @@ -392,9 +417,14 @@ def identify_orphans( logical OR of its current value and the value it would have been set to otherwise. """ - source_rel: DuckDBPyRelation = entities[config.entity_name] + self._materialise_temp_tables_from_entity( + entities, + [config.entity_name, config.target_name] + ) + + source_rel: DuckDBPyRelation = self.connection.table(self._temp_tables[config.entity_name]) source_rel = source_rel.set_alias(config.entity_name) - target_rel: DuckDBPyRelation = entities[config.target_name] + target_rel: DuckDBPyRelation = self.connection.table(self._temp_tables[config.target_name]) target_rel = target_rel.set_alias(config.target_name) if relation_is_empty(source_rel): @@ -465,8 +495,31 @@ def remove_orphans(self, entities: DuckDBEntities, *, config: OrphanRemoval) -> entities[config.entity_name] = filtered_rel + self._materialise_temp_tables_from_entity(entities, [config.entity_name], refresh=True) + return duckdb_rel_to_dictionaries(message_rel) + def identify_and_remove_orphans( + self, + working_directory: URI, + entities: DuckDBEntities, + entity_hierarchy: EntityHierarchy, + key_fields: Optional[dict[str, list[str]]] = None, + ) -> tuple[Messages, dict[EntityName, bool]]: + """ + Identifies and removes orphan records by traversing the EntityHierarchy object. + An orphan is a child record whose parent FK does not exist in the parent entity. + Processes recursively: removes orphans at each level, then processes children. + """ + _msgs, entity_issues_found = super().identify_and_remove_orphans( + working_directory, + entities, + entity_hierarchy, + key_fields, + ) + self._drop_temp_tables() + return _msgs, entity_issues_found + def check_mandatory_group( self, entities: DuckDBEntities, *, config: GroupIdentification ) -> Iterator: From 78a261b23b0c02af83aef5e1227645639c32c50d Mon Sep 17 00:00:00 2001 From: georgeRobertson <50412379+georgeRobertson@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:12:16 +0100 Subject: [PATCH 3/3] Revert "fix: add the ability to cache entities in duckdb to hopefully improve performance on orphan checks" This reverts commit 872972b7d23ab7e9fc70371ecef37800353b7a71. --- .../backends/implementations/duckdb/rules.py | 61 ++----------------- 1 file changed, 4 insertions(+), 57 deletions(-) diff --git a/src/dve/core_engine/backends/implementations/duckdb/rules.py b/src/dve/core_engine/backends/implementations/duckdb/rules.py index 7753393..9e58a1f 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/rules.py +++ b/src/dve/core_engine/backends/implementations/duckdb/rules.py @@ -1,7 +1,7 @@ """Business rule definitions for duckdb backend""" from collections.abc import Callable, Iterator -from typing import get_type_hints, Optional +from typing import get_type_hints from uuid import uuid4 from duckdb import ( @@ -56,12 +56,11 @@ SemiJoin, TableUnion, ) -from dve.core_engine.configuration.v1.hierarchy import EntityHierarchy from dve.core_engine.constants import ORPHANED_RECORD_ENTITY_NAME, RECORD_INDEX_COLUMN_NAME from dve.core_engine.functions import implementations as functions from dve.core_engine.message import FeedbackMessage from dve.core_engine.templating import template_object -from dve.core_engine.type_hints import URI, EntityName, Messages +from dve.core_engine.type_hints import Messages @duckdb_record_index @@ -74,7 +73,6 @@ class DuckDBStepImplementations(BaseStepImplementations[DuckDBPyRelation]): def __init__(self, connection: DuckDBPyConnection, **kwargs): self._connection = connection self.registered_functions = get_all_registered_udfs(self._connection) - self._temp_tables: dict[EntityName, str] = {} super().__init__(**kwargs) @property @@ -115,29 +113,6 @@ def register_udfs( # type: ignore connection.sql(_sql) return cls(connection=connection, **kwargs) - def _materialise_temp_tables_from_entity( - self, - entities: DuckDBEntities, - entity_names: list[EntityName], - refresh: bool = False - ): - """Materialise an entity into a temporary duckdb table.""" - for entity_name in entity_names: - temp_name = f"temp_{entity_name}" - if entity_name in self._temp_tables: - if not refresh: - continue - self.connection.unregister(temp_name) - del self._temp_tables[entity_name] - - self.connection.register(temp_name, entities[entity_name]) - self._temp_tables[entity_name] = temp_name - - def _drop_temp_tables(self) -> None: - for temp_name in self._temp_tables.values(): - self.connection.unregister(temp_name) - self._temp_tables.clear() - def add(self, entities: DuckDBEntities, *, config: ColumnAddition) -> Messages: """A transformation step which adds a column to an entity.""" entity: DuckDBPyRelation = entities[config.entity_name] @@ -417,14 +392,9 @@ def identify_orphans( logical OR of its current value and the value it would have been set to otherwise. """ - self._materialise_temp_tables_from_entity( - entities, - [config.entity_name, config.target_name] - ) - - source_rel: DuckDBPyRelation = self.connection.table(self._temp_tables[config.entity_name]) + source_rel: DuckDBPyRelation = entities[config.entity_name] source_rel = source_rel.set_alias(config.entity_name) - target_rel: DuckDBPyRelation = self.connection.table(self._temp_tables[config.target_name]) + target_rel: DuckDBPyRelation = entities[config.target_name] target_rel = target_rel.set_alias(config.target_name) if relation_is_empty(source_rel): @@ -495,31 +465,8 @@ def remove_orphans(self, entities: DuckDBEntities, *, config: OrphanRemoval) -> entities[config.entity_name] = filtered_rel - self._materialise_temp_tables_from_entity(entities, [config.entity_name], refresh=True) - return duckdb_rel_to_dictionaries(message_rel) - def identify_and_remove_orphans( - self, - working_directory: URI, - entities: DuckDBEntities, - entity_hierarchy: EntityHierarchy, - key_fields: Optional[dict[str, list[str]]] = None, - ) -> tuple[Messages, dict[EntityName, bool]]: - """ - Identifies and removes orphan records by traversing the EntityHierarchy object. - An orphan is a child record whose parent FK does not exist in the parent entity. - Processes recursively: removes orphans at each level, then processes children. - """ - _msgs, entity_issues_found = super().identify_and_remove_orphans( - working_directory, - entities, - entity_hierarchy, - key_fields, - ) - self._drop_temp_tables() - return _msgs, entity_issues_found - def check_mandatory_group( self, entities: DuckDBEntities, *, config: GroupIdentification ) -> Iterator: