gimgit opened a new issue, #17489:
URL: https://github.com/apache/iceberg/issues/17489

   ### Feature Request / Improvement
   
   **TL;DR: `rewrite_data_files` candidate selection only looks at file size, 
so a table whose files are all within the healthy size band is *permanently* 
excluded from sort/z-order compaction. The job reports success on every run 
while sort-key overlap between files never improves — and there is currently no 
way to even detect this state.**
   
   ### Mechanism
   
   File selection for all rewrite strategies lives in 
`SizeBasedFileRewritePlanner#outsideDesiredFileSizeRange`:
   
   ```java
   protected boolean outsideDesiredFileSizeRange(T task) {
     return task.length() < minFileSize || task.length() > maxFileSize;
   }
   ```
   
   `SparkShufflingDataRewritePlanner` (sort/z-order) extends 
`BinPackRewriteFilePlanner` and inherits this selection unchanged — sort key 
overlap between files is never considered. So:
   
   - files inside `[min-file-size-bytes, max-file-size-bytes]` → never 
candidates, no matter how badly they overlap on the sort key
   - the rewrite completes with `rewritten_data_files_count = 0` and reports 
**success**
   
   The perverse consequence: a *fragmented* table (undersized files) gets 
selected and accidentally comes out perfectly clustered, while a *size-healthy* 
table accumulates overlap indefinitely. On one of our production tables, 
scheduled daily sort compaction "succeeded" every day while max overlap depth 
(number of file ranges covering a single point of the sort key) kept growing by 
~1 per day.
   
   ### Deterministic reproduction (Iceberg 1.10.0, Spark 3.5, local)
   
   Two tables with identical, fully-overlapping data — only file size differs:
   
   | case | file sizes | rewrite result | overlap depth |
   |---|---|---|---|
   | A: size-healthy | 8 × ~2.3 MB (inside band) | `rewritten=0`, success | 8 → 
**8** (unchanged) |
   | B: fragmented (control) | 8 × ~0.2 MB (below min) | `rewritten=8` | 8 → 
**1** |
   
   <details>
   <summary>repro.py (self-contained, ~2 min on a laptop)</summary>
   
   ```python
   """
   Reproduction: sort compaction permanently no-ops on size-healthy tables.
   
   Candidate selection in RewriteDataFiles considers only file size
   (min-file-size-bytes / max-file-size-bytes / delete counts). A table whose
   files are all within the healthy size band is never selected for rewrite,
   even when every file overlaps every other file on the sort key. The sort
   strategy therefore reports success daily while clustering never improves.
   
   Control group: the same data written as undersized files IS selected and
   gets perfectly clustered - fragmentation accidentally heals, health
   permanently rots.
   
   Run: .venv/bin/python repro.py
   """
   
   import shutil
   import tempfile
   
   from pyspark.sql import SparkSession
   
   ICEBERG = "org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.10.0"
   WAREHOUSE = tempfile.mkdtemp(prefix="iceberg-overlap-repro-")
   
   TARGET_MB = 2  # scaled-down target-file-size; selection logic only uses 
ratios
   NUM_FILES = 8
   ROWS_PER_FILE_HEALTHY = 90_000  # ~2 MB parquet each -> inside [0.75t, 1.8t]
   ROWS_PER_FILE_SMALL = 8_000     # ~0.3 MB each -> below min threshold
   
   spark = (
       SparkSession.builder.appName("overlap-repro")
       .config("spark.jars.packages", ICEBERG)
       .config("spark.sql.extensions",
               
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
       .config("spark.sql.catalog.local", 
"org.apache.iceberg.spark.SparkCatalog")
       .config("spark.sql.catalog.local.type", "hadoop")
       .config("spark.sql.catalog.local.warehouse", WAREHOUSE)
       .config("spark.driver.memory", "2g")
       .master("local[4]")
       .getOrCreate()
   )
   spark.sparkContext.setLogLevel("ERROR")
   
   
   def create_table(name: str) -> None:
       spark.sql(f"DROP TABLE IF EXISTS local.db.{name}")
       spark.sql(
           f"""CREATE TABLE local.db.{name} (id BIGINT, payload STRING)
               USING iceberg
               TBLPROPERTIES ('format-version'='2')"""
       )
   
   
   def write_overlapping_files(name: str, rows_per_file: int) -> None:
       """Each commit writes one file spanning the FULL id range -> total 
overlap.
   
       The sort order is declared AFTER the appends (see main flow): declaring 
it
       before would make appends locally sorted, which is not the scenario -- 
the
       scenario is an existing unsorted table that sort compaction should fix.
       """
       for i in range(NUM_FILES):
           df = (
               spark.range(rows_per_file)
               .selectExpr(
                   # deterministic pseudo-random spread over the full key range
                   f"(id * 2654435761 + {i}) % 10000000 AS id",
                   "repeat(uuid(), 2) AS payload",
               )
               .coalesce(1)
           )
           df.writeTo(f"local.db.{name}").append()
   
   
   def file_stats(name: str):
       rows = spark.sql(
           f"""SELECT file_size_in_bytes,
                      readable_metrics.id.lower_bound AS lo,
                      readable_metrics.id.upper_bound AS hi
               FROM local.db.{name}.files"""
       ).collect()
       return [(r.file_size_in_bytes, r.lo, r.hi) for r in rows]
   
   
   def overlap_depth(intervals) -> int:
       """Max number of file [lo, hi] ranges covering a single point (sweep 
line)."""
       events = []
       for _, lo, hi in intervals:
           events.append((lo, 1))
           events.append((hi, -1))  # closed intervals: end after start at same 
key
       events.sort(key=lambda e: (e[0], -e[1]))
       depth = best = 0
       for _, delta in events:
           depth += delta
           best = max(best, depth)
       return best
   
   
   def rewrite_sort(name: str):
       return spark.sql(
           f"""CALL local.system.rewrite_data_files(
                 table => 'db.{name}',
                 strategy => 'sort',
                 options => map('target-file-size-bytes', '{TARGET_MB * 1024 * 
1024}')
               )"""
       ).collect()[0]
   
   
   def report(name: str, label: str) -> None:
       stats = file_stats(name)
       sizes = sorted(s // 1024 for s, _, _ in stats)
       print(f"  [{label}] files={len(stats)} sizes_kb={sizes} "
             f"overlap_depth={overlap_depth(stats)}")
   
   
   print(f"warehouse: {WAREHOUSE}\n")
   
   # --- Case A: size-healthy files, fully overlapping ------------------------
   print("Case A: size-healthy files (inside [min,max] band), fully 
overlapping")
   create_table("healthy")
   write_overlapping_files("healthy", ROWS_PER_FILE_HEALTHY)
   spark.sql("ALTER TABLE local.db.healthy WRITE ORDERED BY id")
   report("healthy", "before")
   res = rewrite_sort("healthy")
   print(f"  rewrite result: rewritten={res.rewritten_data_files_count} "
         f"added={res.added_data_files_count}")
   report("healthy", "after ")
   
   # --- Case B: control - same data as undersized files ----------------------
   print("\nCase B (control): undersized files, fully overlapping")
   create_table("fragmented")
   write_overlapping_files("fragmented", ROWS_PER_FILE_SMALL)
   spark.sql("ALTER TABLE local.db.fragmented WRITE ORDERED BY id")
   report("fragmented", "before")
   res = rewrite_sort("fragmented")
   print(f"  rewrite result: rewritten={res.rewritten_data_files_count} "
         f"added={res.added_data_files_count}")
   report("fragmented", "after ")
   
   print("""
   Expected outcome:
     Case A: rewritten=0, depth unchanged (=8). Job reports success, sort order
             never materializes, and this repeats on every scheduled run.
     Case B: rewritten=8, depth -> 1. Being undersized accidentally heals it.
   """)
   
   spark.stop()
   shutil.rmtree(WAREHOUSE, ignore_errors=True)
   ```
   
   </details>
   
   [repro.py](https://github.com/user-attachments/files/30642900/repro.py)
   
   ### Prior discussion
   
   - In #12761, @pvary asked exactly the missing piece — *"Do you have a way to 
identify the non-sorted files somehow?"* — and the thread went stale without an 
answer. This proposal is intended as that answer.
   - #2609's description already sketched overlap-aware selection ("find files 
A, B, C with significant overlap and rewrite them together"), but it never 
landed in selection logic.
   - Recent work (#16827 Hilbert clustering, #16305 k-way merge) improves how 
data is *rewritten*, but not *which* files get selected — the selection gap 
applies to those strategies as well.
   
   ### Proposal
   
   1. **Measurement first**: expose per-partition sort-key overlap depth 
computed purely from manifest `lower_bounds`/`upper_bounds` (no data file 
reads) — as a metadata table column or a procedure. Today every size-based 
health metric reports a badly-clustered table as perfectly healthy.
   2. **Selection (opt-in)**: a planner option (e.g. `min-overlap-depth`) so 
sort/z-order strategies can pick up well-sized but badly-clustered files. The 
`filterFiles`/`filterFileGroups` extension points seem to accommodate this 
without changing default behavior.
   
   Happy to work on a PR for (1) if maintainers think either direction is worth 
pursuing.
   
   ### Query engine
   
   Spark
   
   ### Willingness to contribute
   
   - [ ] I can contribute this improvement/feature independently
   - [x] I would be willing to contribute this improvement/feature with 
guidance from the Iceberg community
   - [ ] I cannot contribute this improvement/feature 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