Skip to content

Commit 0567a6e

Browse files
Gayathri Srividya RajavarapuGayathri Srividya Rajavarapu
authored andcommitted
fix spec-evolution overwrite predicate matching
1 parent 40661d0 commit 0567a6e

2 files changed

Lines changed: 12 additions & 5 deletions

File tree

‎pyiceberg/table/__init__.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -632,6 +632,7 @@ def dynamic_partition_overwrite(
632632
for spec_id, hist_spec in all_specs.items():
633633
hist_source_ids = {field.source_id for field in hist_spec.fields}
634634
missing_source_ids = current_source_ids - hist_source_ids
635+
has_overlap_with_current = bool(hist_source_ids & current_source_ids)
635636

636637
per_record_exprs: list[BooleanExpression] = []
637638
for partition_record in partitions_to_overwrite:
@@ -640,7 +641,7 @@ def dynamic_partition_overwrite(
640641
value = partition_record[source_id_to_pos[source_id]]
641642
if value is not None:
642643
field_pred: BooleanExpression = EqualTo(Reference(col_name), value)
643-
if source_id in missing_source_ids:
644+
if source_id in missing_source_ids and has_overlap_with_current:
644645
field_pred = Or(field_pred, IsNull(Reference(col_name)))
645646
else:
646647
field_pred = IsNull(Reference(col_name))

‎pyiceberg/table/update/snapshot.py‎

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -435,10 +435,13 @@ def _copy_with_new_status(entry: ManifestEntry, status: ManifestEntryStatus) ->
435435

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

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
438+
def _strict_metrics_for_spec(spec_id: int) -> Callable[[DataFile], bool]:
439+
predicate = self._per_spec_predicates.get(spec_id, self._predicate)
440+
return _StrictMetricsEvaluator(schema, predicate, case_sensitive=self._case_sensitive).eval
441+
442+
def _inclusive_metrics_for_spec(spec_id: int) -> Callable[[DataFile], bool]:
443+
predicate = self._per_spec_predicates.get(spec_id, self._predicate)
444+
return _InclusiveMetricsEvaluator(schema, predicate, case_sensitive=self._case_sensitive).eval
442445

443446
existing_manifests = []
444447
total_deleted_entries = []
@@ -458,6 +461,9 @@ def _copy_with_new_status(entry: ManifestEntry, status: ManifestEntryStatus) ->
458461
existing_manifests.append(manifest_file)
459462
else:
460463
# It is relevant, let's check out the content
464+
spec_id = manifest_file.partition_spec_id
465+
strict_metrics_evaluator = _strict_metrics_for_spec(spec_id)
466+
inclusive_metrics_evaluator = _inclusive_metrics_for_spec(spec_id)
461467
deleted_entries = []
462468
existing_entries = []
463469
for entry in manifest_file.fetch_manifest_entry(io=self._io, discard_deleted=True):

0 commit comments

Comments
 (0)