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]