Skip to content

Commit 40661d0

Browse files
Gayathri Srividya RajavarapuGayathri Srividya Rajavarapu
authored andcommitted
fix dynamic overwrite null-row regression after spec evolution
1 parent 8a05f44 commit 40661d0

2 files changed

Lines changed: 14 additions & 17 deletions

File tree

pyiceberg/table/__init__.py

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -622,6 +622,11 @@ def dynamic_partition_overwrite(
622622
# _compute_deletes uses the right predicate when evaluating each manifest file.
623623
source_id_to_pos = {field.source_id: pos for pos, field in enumerate(current_spec.fields)}
624624
source_id_to_col = {field.source_id: schema.find_field(field.source_id).name for field in current_spec.fields}
625+
exact_delete_filter = self._build_partition_predicate(
626+
partition_records=partitions_to_overwrite,
627+
spec=current_spec,
628+
schema=schema,
629+
)
625630

626631
per_spec_predicates: dict[int, BooleanExpression] = {}
627632
for spec_id, hist_spec in all_specs.items():
@@ -645,24 +650,21 @@ def dynamic_partition_overwrite(
645650

646651
per_spec_predicates[spec_id] = Or(*per_record_exprs) if len(per_record_exprs) > 1 else per_record_exprs[0]
647652

648-
preds = list(per_spec_predicates.values())
649-
global_delete_filter = Or(*preds) if len(preds) > 1 else preds[0] if preds else AlwaysFalse()
650-
651653
# Open the delete snapshot and set per-spec predicates before committing.
652654
# This mirrors Transaction.delete() but injects per_spec_predicates so that
653655
# _compute_deletes uses the right predicate for each historical spec.
654656
from pyiceberg.io.pyarrow import ArrowScan, _dataframe_to_data_files, _expression_to_complementary_pyarrow
655657

656658
with self.update_snapshot(snapshot_properties=snapshot_properties, branch=branch).delete() as delete_snapshot:
657659
delete_snapshot._per_spec_predicates = per_spec_predicates
658-
delete_snapshot.delete_by_predicate(global_delete_filter)
660+
delete_snapshot.delete_by_predicate(exact_delete_filter)
659661

660662
# Handle partial-match files that need to be rewritten (copy-on-write).
661663
if delete_snapshot.rewrites_needed is True:
662-
bound_delete_filter = bind(self.table_metadata.schema(), global_delete_filter, case_sensitive=True)
664+
bound_delete_filter = bind(self.table_metadata.schema(), exact_delete_filter, case_sensitive=True)
663665
preserve_row_filter = _expression_to_complementary_pyarrow(bound_delete_filter, self.table_metadata.schema())
664666

665-
file_scan = self._scan(row_filter=global_delete_filter)
667+
file_scan = self._scan(row_filter=exact_delete_filter)
666668
if branch is not None:
667669
file_scan = file_scan.use_ref(branch)
668670

pyiceberg/table/update/snapshot.py

Lines changed: 6 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -362,7 +362,8 @@ def fetch_manifest_entry(self, manifest: ManifestFile, discard_deleted: bool = T
362362

363363
def _build_partition_projection(self, spec_id: int) -> BooleanExpression:
364364
project = inclusive_projection(self.schema(), self.spec(spec_id), self._case_sensitive)
365-
return project(self._predicate)
365+
predicate = self._per_spec_predicates.get(spec_id, self._predicate)
366+
return project(predicate)
366367

367368
@cached_property
368369
def partition_filters(self) -> KeyDefaultDict[int, BooleanExpression]:
@@ -434,13 +435,10 @@ def _copy_with_new_status(entry: ManifestEntry, status: ManifestEntryStatus) ->
434435

435436
manifest_evaluators: dict[int, Callable[[ManifestFile], bool]] = KeyDefaultDict(self._build_manifest_evaluator)
436437

437-
def _strict_metrics_for_spec(spec_id: int) -> Callable[[DataFile], bool]:
438-
predicate = self._per_spec_predicates.get(spec_id, self._predicate)
439-
return _StrictMetricsEvaluator(schema, predicate, case_sensitive=self._case_sensitive).eval
440-
441-
def _inclusive_metrics_for_spec(spec_id: int) -> Callable[[DataFile], bool]:
442-
predicate = self._per_spec_predicates.get(spec_id, self._predicate)
443-
return _InclusiveMetricsEvaluator(schema, predicate, case_sensitive=self._case_sensitive).eval
438+
strict_metrics_evaluator = _StrictMetricsEvaluator(schema, self._predicate, case_sensitive=self._case_sensitive).eval
439+
inclusive_metrics_evaluator = _InclusiveMetricsEvaluator(
440+
schema, self._predicate, case_sensitive=self._case_sensitive
441+
).eval
444442

445443
existing_manifests = []
446444
total_deleted_entries = []
@@ -460,9 +458,6 @@ def _inclusive_metrics_for_spec(spec_id: int) -> Callable[[DataFile], bool]:
460458
existing_manifests.append(manifest_file)
461459
else:
462460
# It is relevant, let's check out the content
463-
spec_id = manifest_file.partition_spec_id
464-
strict_metrics_evaluator = _strict_metrics_for_spec(spec_id)
465-
inclusive_metrics_evaluator = _inclusive_metrics_for_spec(spec_id)
466461
deleted_entries = []
467462
existing_entries = []
468463
for entry in manifest_file.fetch_manifest_entry(io=self._io, discard_deleted=True):

0 commit comments

Comments
 (0)