QlikFrederic opened a new issue, #4031:
URL: https://github.com/apache/iceberg-python/issues/4031

   ### Apache Iceberg version
   
   0.12.0 (latest release)
   
   ### Please describe the bug 🐞
   
   An overwrite that deletes a live data file is sometimes rejected with:
   
   `ValidationException: Missing required files to delete: <path>`
   The file is live in the parent snapshot, and it is still live afterwards. A 
later overwrite of the same files commits normally.
   
   ### Cause
   
   `_OverwriteFiles._deleted_entries` (pyiceberg/table/update/snapshot.py) 
builds one manifest evaluator per partition spec and calls it from the 
`ExecutorFactory `thread pool:
   
   ```
   manifest_evaluators = KeyDefaultDict(self._build_manifest_evaluator)
   
   def _get_entries(manifest):
       if not manifest_evaluators[manifest.partition_spec_id](manifest):
           return []
       ...
   
   list_of_entries = executor.map(_get_entries, 
previous_snapshot.manifests(self._io))
   ```
   manifest_evaluator() returns the bound eval of a single 
_ManifestEvalVisitor, and eval keeps its input on the instance:
   
   ```
   def eval(self, manifest: ManifestFile) -> bool:
       if partitions := manifest.partitions:
           self.partition_fields = partitions          # shared by every thread
           return visit(self.partition_filter, self)   # each leaf reads 
self.partition_fields
       return ROWS_MIGHT_MATCH
   
   ```
   If thread A stores manifest A's summaries, and thread B then stores manifest 
B's summaries before A reads its first leaf, A's manifest is judged by B's 
partition bounds. When A holds the files being deleted, it is skipped, those 
files are never found, and `_validate_required_deletes` raises.
   
   The same race can also keep a manifest that should be skipped. That 
direction is harmless, because the entry comparison finds nothing in it.
   
   As far as I can see, `_deleted_entries` is the only caller that shares an 
evaluator across threads. `DataScan.plan_files`, `_existing_manifests `and 
`_DeleteFiles` evaluate manifests one at a time.
   
   ### Impact
   
   We hit this on production tables with many manifests. Roughly 1 in 700 
overwrites is rejected, each time for every planned file that sits in one 
manifest (hundreds of live files). No data is lost, but the commit and its work 
are thrown away. It started when we upgraded from 0.11 to 0.12, which added 
`_validate_required_deletes`.
   
   ### Reproduction
   
   The race window is small, so the script below forces the interleaving. It 
holds the thread evaluating manifest A at its first predicate leaf until 
another thread has evaluated manifest B. Nothing else is changed. It uses 
`InMemoryCatalog` and needs `pyiceberg[sql-sqlite]` and `pyarrow`.
   
   repro_shared_manifest_evaluator.py
   ```
   """Reproduce: overwrite's delete check rejects live files when threads share 
one manifest evaluator.
   _OverwriteFiles._deleted_entries evaluates the parent snapshot's manifests 
on the ExecutorFactory thread pool, but all threads share one 
_ManifestEvalVisitor per spec, and _ManifestEvalVisitor.eval stores the 
manifest's partition summaries on self.partition_fields. A thread switch at the 
wrong moment makes one manifest be judged by another manifest's bounds, so the 
manifest holding the file is skipped and the commit fails with
   "ValidationException: Missing required files to delete" although the file is 
live.
   
   The switch is rare in practice, so this script forces it: it holds the 
thread evaluating manifest A at its first
   predicate leaf until another thread has evaluated manifest B. Nothing else 
is changed.
   """
   
   import tempfile
   import threading
   
   import pyarrow as pa
   import pyiceberg
   from pyiceberg.catalog.memory import InMemoryCatalog
   from pyiceberg.exceptions import ValidationException
   from pyiceberg.expressions import visitors
   from pyiceberg.partitioning import PartitionField, PartitionSpec
   from pyiceberg.schema import Schema
   from pyiceberg.transforms import IdentityTransform
   from pyiceberg.types import LongType, NestedField, StringType
   
   print(f"pyiceberg {pyiceberg.__version__}")
   
   warehouse = tempfile.mkdtemp()
   catalog = InMemoryCatalog("default", warehouse=warehouse)
   catalog.create_namespace("default")
   schema = Schema(
       NestedField(1, "tenant", StringType(), required=False),
       NestedField(2, "value", LongType(), required=False),
   )
   spec = PartitionSpec(PartitionField(source_id=1, field_id=1000, 
transform=IdentityTransform(), name="tenant"))
   table = catalog.create_table("default.events", schema=schema, 
partition_spec=spec)
   
   # Two appends -> two manifests: A holds tenant "a", B holds tenant "z".
   table.append(pa.table({"tenant": ["a"], "value": [1]}))
   table.append(pa.table({"tenant": ["z"], "value": [2]}))
   file_a = next(task.file for task in table.scan().plan_files() if 
task.file.partition[0] == "a")
   print(f"manifests: {len(table.current_snapshot().manifests(table.io))}, file 
to delete: {file_a.file_path.rsplit('/', 1)[-1]}")
   
   # --- Force the interleaving 
-------------------------------------------------------------------------------------
   original_eval = visitors._ManifestEvalVisitor.eval
   original_visit_equal = visitors._ManifestEvalVisitor.visit_equal
   a_at_first_leaf = threading.Event()
   b_evaluated = threading.Event()
   local = threading.local()
   
   
   def lower_tenant(manifest) -> bytes:
       return manifest.partitions[0].lower_bound
   
   
   def eval_with_hook(self, manifest):
       local.first_leaf = True
       local.is_a = lower_tenant(manifest) == b"a"
       if lower_tenant(manifest) == b"z":
           a_at_first_leaf.wait(timeout=2)  # start B only once A has stored 
its summaries
           result = original_eval(self, manifest)
           b_evaluated.set()
           return result
       return original_eval(self, manifest)
   
   
   def visit_equal_with_hook(self, term, literal):
       if getattr(local, "first_leaf", False) and local.is_a:
           local.first_leaf = False
           # self.partition_fields now holds A's summaries; let another thread 
evaluate B before reading them.
           a_at_first_leaf.set()
           b_evaluated.wait(timeout=2)
       return original_visit_equal(self, term, literal)
   
   
   visitors._ManifestEvalVisitor.eval = eval_with_hook
   visitors._ManifestEvalVisitor.visit_equal = visit_equal_with_hook
   # 
-----------------------------------------------------------------------------------------------------------------
   
   try:
       with table.transaction() as transaction:
           with transaction.update_snapshot().overwrite() as overwrite:
               overwrite.delete_data_file(file_a)
       print("commit succeeded (no race)")
   except ValidationException as error:
       print(f"ValidationException: {error}")
   finally:
       visitors._ManifestEvalVisitor.eval = original_eval
       visitors._ManifestEvalVisitor.visit_equal = original_visit_equal
   
   table.refresh()
   live = {task.file.file_path for task in table.scan().plan_files()}
   print(f"file is live in the current snapshot: {file_a.file_path in live}")
   
   # Control: the same delete without the forced interleaving commits.
   with table.transaction() as transaction:
       with transaction.update_snapshot().overwrite() as overwrite:
           overwrite.delete_data_file(file_a)
   print("same delete without the hook: committed")
   ```
   
   Output on 0.12.0:
   
   ```
   pyiceberg 0.12.0
   manifests: 2, file to delete: 
00000-0-1474ca0a-3540-41f6-94a5-121577be40ea.parquet
   ValidationException: Missing required files to delete: 
/tmp/tmp6fb_dff5/default/events/data/tenant=a/00000-0-1474ca0a-3540-41f6-94a5-121577be40ea.parquet
   file is live in the current snapshot: True
   same delete without the hook: committed
   ```
   ### Possible fix
   
   Keep the per-manifest state out of the shared instance. For example, 
evaluate on a shallow copy, so the bound filter stays shared and read-only:
   
   ```
   def eval(self, manifest: ManifestFile) -> bool:
       if partitions := manifest.partitions:
           visitor = copy.copy(self)
           visitor.partition_fields = partitions
           return visit(self.partition_filter, visitor)
       return ROWS_MIGHT_MATCH
   ```
   With this change, the script above commits the delete. Alternatives: create 
one evaluator per task inside `_get_entries`, or pass the summaries through the 
visit instead of storing them on `self`.
   
   ### Willingness to contribute
   
   - [ ] I can contribute a fix for this bug independently
   - [x] I would be willing to contribute a fix for this bug with guidance from 
the Iceberg community
   - [ ] I cannot contribute a fix for this bug at this time


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to