Skip to content

Commit e10b797

Browse files
perf: remove orphan tracker as it's causing significant performance degradation
1 parent 8f5e1ae commit e10b797

7 files changed

Lines changed: 57 additions & 129 deletions

File tree

‎src/dve/core_engine/backends/base/rules.py‎

Lines changed: 35 additions & 73 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
)
1818
from dve.core_engine.backends.base.core import get_entity_type
1919
from dve.core_engine.backends.exceptions import render_error
20-
from dve.core_engine.backends.metadata.reporting import ReportingConfig
2120
from dve.core_engine.backends.metadata.rules import (
2221
AbstractStep,
2322
Aggregation,
@@ -36,7 +35,6 @@
3635
Notification,
3736
OneToOneJoin,
3837
OrphanIdentification,
39-
OrphanRemoval,
4038
ParentMetadata,
4139
RenameEntity,
4240
Rule,
@@ -47,7 +45,6 @@
4745
)
4846
from dve.core_engine.backends.types import Entities, EntityType, StageSuccessful
4947
from dve.core_engine.configuration.v1.hierarchy import EntityHierarchy, HierarchyNode
50-
from dve.core_engine.constants import ORPHANED_RECORD_ENTITY_NAME
5148
from dve.core_engine.exceptions import CriticalProcessingError
5249
from dve.core_engine.loggers import get_logger
5350
from dve.core_engine.message import FeedbackMessage
@@ -317,7 +314,7 @@ def join_header(self, entities: Entities, *, config: HeaderJoin) -> Messages:
317314
@abstractmethod
318315
def identify_orphans(
319316
self, entities: Entities, *, config: OrphanIdentification
320-
) -> tuple[Messages, int]:
317+
) -> Iterable:
321318
"""Identify records in an entity which don't have at least one corresponding
322319
match in the target. A new boolean column will be added to `entity` ('IsOrphaned')
323320
indicating whether the condition matched.
@@ -330,18 +327,6 @@ def identify_orphans(
330327
"""
331328
raise NotImplementedError
332329

333-
@abstractmethod
334-
def remove_orphans(self, entities: Entities, *, config: OrphanRemoval) -> Iterator:
335-
"""
336-
Remove orphaned records from an entity based on the orphans found in
337-
identify_orphans method. Returns a generator objects with the records removed
338-
for generating feedback messages from.
339-
340-
This may not be implemented by some backends.
341-
342-
"""
343-
raise NotImplementedError
344-
345330
@abstractmethod
346331
def check_mandatory_group(self, entities: Entities, *, config: GroupIdentification) -> Iterator:
347332
"""
@@ -395,83 +380,60 @@ def identify_and_remove_orphans(
395380
Processes recursively: removes orphans at each level, then processes children.
396381
"""
397382

398-
def process_node(node: HierarchyNode):
383+
def process_node(node: HierarchyNode) -> bool:
399384
"""Identify orphans and remove in a given node"""
400385
issues_found: bool = False
401386
if node.parent_entity is None:
402387
return issues_found
403388

404-
self.logger.info(f"Identifying orphans in {node.entity_name}")
389+
self.logger.info(f"Checking for orphan records in {node.entity_name}")
405390

406391
join_expr = " AND ".join(
407392
f"{node.parent_entity}.{k} = {node.entity_name}.{v}"
408393
for k, v in node.join_fields.items()
409394
)
410-
411-
_, no_orphs = self.identify_orphans(
412-
entities=entities,
413-
config=OrphanIdentification(
414-
id=list(node.join_fields.values())[0],
415-
entity_name=node.entity_name,
416-
target_name=node.parent_entity,
417-
join_condition=join_expr,
418-
),
419-
)
420-
421-
if no_orphs > 0:
422-
self.logger.info(f"Removing records with missing parent from {node.entity_name}")
423-
issues_found = True
424-
location = list(node.join_fields.values())[0]
425-
with BackgroundMessageWriter(
426-
working_directory=working_directory,
427-
dve_stage=self.__stage_name__,
428-
key_fields=key_fields,
429-
logger=self.logger,
430-
) as msg_writer:
431-
_orph_records = self.remove_orphans(
432-
entities=entities,
433-
config=OrphanRemoval(
434-
entity_name=node.entity_name,
435-
reporting=ReportingConfig(
436-
emit="record_failure",
437-
code=node.missing_parent_id_error_code,
438-
message=node.missing_parent_id_error_message,
439-
location=location,
440-
),
395+
location = list(node.join_fields.values())[0]
396+
with BackgroundMessageWriter(
397+
working_directory=working_directory,
398+
dve_stage=self.__stage_name__,
399+
key_fields=key_fields,
400+
logger=self.logger,
401+
) as msg_writer:
402+
_orph_records = self.identify_orphans(
403+
entities=entities,
404+
config=OrphanIdentification(
405+
id=list(node.join_fields.values())[0],
406+
entity_name=node.entity_name,
407+
target_name=node.parent_entity,
408+
join_condition=join_expr,
409+
),
410+
)
411+
_messages = [
412+
FeedbackMessage(
413+
entity=node.entity_name,
414+
record=record, # type: ignore
415+
error_location=location,
416+
error_message=template_object(
417+
node.missing_parent_id_error_message, record
441418
),
419+
failure_type="record",
420+
error_type="record",
421+
error_code=node.missing_parent_id_error_code,
422+
reporting_field=location,
423+
category="Parent Missing",
442424
)
443-
# moved to batch the write - risky if large number of
444-
msg_writer.write_queue.put(
445-
[
446-
FeedbackMessage(
447-
entity=node.entity_name,
448-
record=record, # type: ignore
449-
error_location=location,
450-
error_message=template_object(
451-
node.missing_parent_id_error_message, record
452-
),
453-
failure_type="record",
454-
error_type="record",
455-
error_code=node.missing_parent_id_error_code,
456-
reporting_field=location,
457-
category="Parent Missing",
458-
)
459-
for record in _orph_records
460-
]
461-
)
425+
for record in _orph_records
426+
]
427+
msg_writer.write_queue.put(_messages)
462428

463-
return issues_found
429+
return len(_messages) > 0
464430

465431
entity_issues_found: dict[EntityName, bool] = {}
466432

467433
for tree in entity_hierarchy.entity_trees.values():
468434
for node in tree.iterate_root_down():
469435
entity_issues_found[node.entity_name] = process_node(node)
470436

471-
_orph_rel = entities.get(ORPHANED_RECORD_ENTITY_NAME)
472-
if _orph_rel is not None:
473-
del entities[ORPHANED_RECORD_ENTITY_NAME]
474-
475437
entities.update(entities)
476438

477439
return [], entity_issues_found

‎src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33

44
"""Helper objects for duckdb data contract implementation"""
55

6+
import itertools
67
from collections.abc import Generator, Iterator
78
from dataclasses import is_dataclass
89
from datetime import date, datetime, time
@@ -99,8 +100,11 @@ def __call__(self):
99100

100101
def table_exists(connection: DuckDBPyConnection, table_name: str) -> bool:
101102
"""check if a table exists in a given DuckDBPyConnection"""
102-
return table_name in map(lambda x: x[0], connection.sql("SHOW TABLES").fetchall())
103+
return table_name in get_all_existing_ddb_tables(connection)
103104

105+
def get_all_existing_ddb_tables(connection: DuckDBPyConnection) -> tuple[str]:
106+
"""Fetch all tables available ina given duckdb connection"""
107+
return tuple(itertools.chain.from_iterable(connection.sql("SHOW TABLES").fetchall()))
104108

105109
def relation_is_empty(relation: DuckDBPyRelation) -> bool:
106110
"""Check if a duckdb relation is empty"""

‎src/dve/core_engine/backends/implementations/duckdb/rules.py‎

Lines changed: 13 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
"""Business rule definitions for duckdb backend"""
22
# pylint: disable=R0801
3-
from collections.abc import Callable, Iterator
3+
from collections.abc import Callable, Iterable, Iterator
44
from typing import get_type_hints
55
from uuid import uuid4
66

@@ -53,11 +53,10 @@
5353
Notification,
5454
OneToOneJoin,
5555
OrphanIdentification,
56-
OrphanRemoval,
5756
SemiJoin,
5857
TableUnion,
5958
)
60-
from dve.core_engine.constants import ORPHANED_RECORD_ENTITY_NAME, RECORD_INDEX_COLUMN_NAME
59+
from dve.core_engine.constants import RECORD_INDEX_COLUMN_NAME
6160
from dve.core_engine.functions import implementations as functions
6261
from dve.core_engine.message import FeedbackMessage
6362
from dve.core_engine.templating import template_object
@@ -385,7 +384,7 @@ def identify_orphans(
385384
entities: DuckDBEntities,
386385
*,
387386
config: OrphanIdentification,
388-
) -> tuple[Messages, int]:
387+
) -> Iterable:
389388
"""Identify records in an entity which don't have at least one corresponding
390389
match in the target. A new boolean column will be added to `entity` ('IsOrphaned')
391390
indicating whether the condition matched.
@@ -401,56 +400,34 @@ def identify_orphans(
401400

402401
if relation_is_empty(source_rel):
403402
self.logger.info(f"{config.entity_name} is empty. Skipping orphan check.")
404-
return [], 0
403+
return []
405404

406405
match_name = f"matched_{uuid4().hex}"
407406
target_rel = target_rel.select(
408407
StarExpression(exclude=[]), ConstantExpression(1).alias(match_name)
409408
).set_alias(config.target_name)
410409

411-
pk, _fk = config.join_condition.split("=")
412-
413410
orphaned_rel: DuckDBPyRelation = (
414411
source_rel.join(target_rel, condition=config.join_condition, how="left")
415412
.aggregate(
416413
f"{config.entity_name}.{RECORD_INDEX_COLUMN_NAME}, {config.entity_name}.{config.id}, coalesce(count({match_name}), 0)==0 AS IsOrphaned" # pylint: disable=C0301
417414
)
418415
.filter("IsOrphaned")
419-
.select(
420-
RECORD_INDEX_COLUMN_NAME,
421-
ConstantExpression(config.entity_name).alias("entity_name"),
422-
ConstantExpression(pk.strip().rsplit(".")[1]).alias("pk"),
423-
ColumnExpression(config.id).alias("pk_value"), # type: ignore
424-
)
425-
.unique("*")
416+
.select(RECORD_INDEX_COLUMN_NAME)
417+
.set_alias("orphan")
426418
)
427-
_orph_records: tuple[int] = orphaned_rel.count(RECORD_INDEX_COLUMN_NAME).fetchone() # type: ignore # pylint: disable=C0301
428-
if _orph_records:
429-
_no_orphans = _orph_records[0]
430-
if entities.get(ORPHANED_RECORD_ENTITY_NAME) is not None:
431-
entities[ORPHANED_RECORD_ENTITY_NAME] = entities[ORPHANED_RECORD_ENTITY_NAME].union(
432-
orphaned_rel
433-
)
434-
else:
435-
entities[ORPHANED_RECORD_ENTITY_NAME] = orphaned_rel
436-
else:
437-
_no_orphans = 0
438-
self.logger.info(f"Found {_no_orphans} orphaned records in {config.entity_name}.")
439419

440-
return [], _no_orphans
420+
if relation_is_empty(orphaned_rel):
421+
self.logger.info(
422+
f"Found 0 orphan records between {config.entity_name} and {config.target_name}"
423+
)
424+
return []
441425

442-
def remove_orphans(self, entities: DuckDBEntities, *, config: OrphanRemoval) -> Iterator:
443-
"""Method to remove identified orphans in the orphan tracker entity."""
444-
orphan_rel = (
445-
entities[ORPHANED_RECORD_ENTITY_NAME]
446-
.filter(f"entity_name = '{config.entity_name}'")
447-
.set_alias("orphan")
448-
)
449426
message_rel = (
450427
entities[config.entity_name]
451428
.set_alias(config.entity_name)
452429
.join(
453-
orphan_rel,
430+
orphaned_rel,
454431
f"{config.entity_name}.{RECORD_INDEX_COLUMN_NAME} = orphan.{RECORD_INDEX_COLUMN_NAME}", # pylint: disable=C0301
455432
"semi",
456433
)
@@ -459,7 +436,7 @@ def remove_orphans(self, entities: DuckDBEntities, *, config: OrphanRemoval) ->
459436
entities[config.entity_name]
460437
.set_alias(config.entity_name)
461438
.join(
462-
orphan_rel,
439+
orphaned_rel,
463440
f"{config.entity_name}.{RECORD_INDEX_COLUMN_NAME} = orphan.{RECORD_INDEX_COLUMN_NAME}", # pylint: disable=C0301
464441
"anti",
465442
)

‎src/dve/core_engine/backends/implementations/spark/rules.py‎

Lines changed: 0 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,6 @@
4343
Notification,
4444
OneToOneJoin,
4545
OrphanIdentification,
46-
OrphanRemoval,
4746
SelectColumns,
4847
SemiJoin,
4948
TableUnion,
@@ -378,15 +377,6 @@ def identify_orphans(
378377
entities[config.new_entity_name or config.entity_name] = result
379378
return [], 0
380379

381-
def remove_orphans(
382-
self,
383-
entities: SparkEntities,
384-
*,
385-
config: OrphanRemoval,
386-
) -> Iterator:
387-
# TODO - implement for spark
388-
raise NotImplementedError
389-
390380
def check_mandatory_group(
391381
self, entities: SparkEntities, *, config: GroupIdentification
392382
) -> Iterator:

‎src/dve/core_engine/constants.py‎

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,3 @@
66
CONTRACT_ERROR_VALUE_FIELD_NAME: str = "__error_value"
77
"""The name of the field that can be used to extract the field value that caused
88
a pydantic validation error"""
9-
10-
ORPHANED_RECORD_ENTITY_NAME: str = "orphaned_record_tracker"
11-
"""Name of entity to keep track of records where there is a missing parent record"""

‎tests/features/steps/steps_post_pipeline.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -129,4 +129,5 @@ def check_entity_row_counts(context: Context):
129129
record = row.as_dict()
130130
entity_name = record["entity_name"]
131131
expected_count = int(record["row_count"])
132-
assert expected_count == read_output_parquet(processing_loc, entity_name, "business_rules").shape[0]
132+
output_df = read_output_parquet(processing_loc, entity_name, "business_rules")
133+
assert expected_count == output_df.shape[0], output_df

‎tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_rules.py‎

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,6 @@
4040
SemiJoin,
4141
TableUnion,
4242
)
43-
from dve.core_engine.constants import ORPHANED_RECORD_ENTITY_NAME
4443
from dve.core_engine.configuration.v1.hierarchy import (
4544
EntityHierarchy, HierarchyNode
4645
)
@@ -622,7 +621,7 @@ def test_identify_orphan_record_single_entity(self):
622621
)
623622

624623
rules = DuckDBStepImplementations(connection=cnn)
625-
_msgs = rules.identify_orphans(
624+
msgs = rules.identify_orphans(
626625
mod_entities.entities,
627626
config=OrphanIdentification(
628627
id="flight_id",
@@ -631,9 +630,7 @@ def test_identify_orphan_record_single_entity(self):
631630
join_condition="passengers.flight_id = flights.flight_id"
632631
)
633632
)
634-
result = mod_entities[ORPHANED_RECORD_ENTITY_NAME]
635-
assert result.count("*").fetchone()[0] == 1 # type: ignore
636-
assert result.select("entity_name").unique("*").count("*").fetchone()[0] == 1 # type: ignore
633+
assert len(list(msgs)) == 1
637634

638635
def test_identify_and_remove_orphans(self):
639636
with duckdb.connect() as cnn:

0 commit comments

Comments
 (0)