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

   ### Apache Iceberg version
   
   1.11.0 (latest release)
   
   ### Query engine
   
   Spark
   
   ### Please describe the bug 🐞
   
   v3 deletion-vector MERGE (merge-on-read) grows executor memory without bound 
under
   high-churn streaming upserts until executors OOM; reproduced on vanilla 
Spark 4.0.2 +
   Iceberg 1.11.0 + Glue/S3 with NO compaction (follow-up to #17209 with the 
OSS repro
   requested there)
   
   Spark 4.0.2 (Scala 2.13), official `apache/spark:4.0.2` image, Spark on 
Kubernetes
   (Amazon EKS, 6 x m5.2xlarge), GlueCatalog + S3FileIO (plain S3 bucket, 
SSE-KMS),
   `s3a://` streaming checkpoints. **No compaction of any kind** (no 
maintenance job, no
   managed service, no `rewrite_data_files` / `rewrite_manifests`, ever).
   
   ### Please describe the bug
   
   **Summary.** A continuous Structured Streaming CDC "mirror" `MERGE INTO` 
(upsert keyed
   on the bucket column) against a format-version **3** merge-on-read table 
exhibits
   executor memory growth batch over batch until executors are 
container-OOM-killed
   (exit 137) and the job dies. The identical job against a format-version 
**2** table is
   stable for hours (shown at length in #17209 on EMR; being re-confirmed on 
this stack,
   see "control run" below).
   
   **Workload** (fully scripted, see *Reproduction*):
   
   - Table: `PARTITIONED BY (bucket(64, account_id))`, MoR
     (`write.{delete,update,merge}.mode=merge-on-read`, 
`write.merge.distribution-mode=hash`,
     `write.spark.fanout.enabled=true`), `format-version=3` -> deletion vectors.
   - Feed: DMS-like CDC over Kafka at ~87k msg/s; 80% of changes hit 100k hot 
keys (of 2M),
     ~85% updates / ~15% deletes after first insert -> nearly every batch is 
`WHEN MATCHED`,
     i.e. sustained row-level deletes against already-written data files across 
all 64 buckets.
   - Per 60s micro-batch: dedup to latest change per key (`row_number() over 
(partition by
     account_id order by seq desc)`), then
     `MERGE INTO mirror ON account_id` with `WHEN MATCHED DELETE / UPDATE`, 
`WHEN NOT MATCHED INSERT`.
   - Sizing: 6 executors x 5 cores x 16g heap (+10% overhead -> 17.6 GiB pod 
limit), 4g driver.
   
   **Observed (v3, fresh table, single job on an otherwise idle cluster)**
   
   | batch | rows | merge time |
   |---|---|---|
   | 1-8   | 3.7-5.2M | 26-32s |
   | 9     | 4.4M | 59s |
   | 10    | 5.2M | 73s |
   | ...   | growing | 150-220s stalls |
   | ~batch 25+ (t+27min, ~130M rows merged) | - | executors OOM in a loop, job 
dies |
   
   Driver log at the end (executor IDs reached 19 - i.e. 13 replacements had 
already been
   OOM-killed and respawned before the context gave up):
   
   ```
   ERROR TaskSchedulerImpl: Lost executor 13 on 172.31.87.124:
   The executor with id 13 exited with exit code 137(SIGKILL, possible 
container OOM).
   ...
   ERROR MicroBatchExecution: Query streaming-cdc-mirror-accounts_mirror_dv 
terminated with error
   org.apache.spark.SparkException: Job 51 cancelled because SparkContext was 
shut down
   ```
   
   Spark UI symptoms while degrading (same as on EMR in #17209): tasks with 
very high
   *task deserialization time* on the MERGE write stage, executors freezing 
>120s (GC) and
   being evicted for missed heartbeats, then 137s.
   
   <img width="2432" height="708" alt="Image" 
src="https://github.com/user-attachments/assets/3924260a-21ff-4b9e-87f9-5b23feffa621";
 />
   
   **Control run (v2).** Same binary, same feed, same sizing, same cluster, only
   `fv=2`: stable indefinitely, batch time flat (~40-55s), zero executor 
losses, while
   carrying ~10x more delete files than the v3 table had when it started dying
   (detailed v2/v3 table in #17209; re-run on this OSS stack ongoing, will 
attach the
   event logs of both runs).
   
   ### Proposed mechanism ( this cames from Fable reading the whole source 
code, I need to study what is proposed below).
   
   `SparkPositionDeltaWrite` makes the v3 DV write a *read-modify-write* whose 
working set
   grows with the table, while the v2 path is a bounded blind append:
   
   1. **Driver side, every batch:** `broadcastRewritableDeletes()` calls
      `scan.rewritableDeletes(useDVs=true)` and broadcasts a
      `Map<String, DeleteFileSet>` of *every scanned data file that has 
existing deletes ->
      its delete files*. With DVs every previously-touched data file carries 
one, and a
      hot-key upsert keyed on the bucket column touches all 64 buckets every 
batch, so this
      map converges to "all data files in the table" and grows monotonically as 
appends add
      files (no compaction to consolidate). A monotonically growing broadcast 
is consistent
      with the growing task-deserialization times we observe.
   2. **Executor side:** for DVs the writer is `PartitioningDVWriter` with a
      `PreviousDeleteLoader`: for each data file a task deletes from, it loads 
that file's
      existing DV bitmap, unions the new positions, and keeps the merged
      `PositionDeleteIndex` in memory until the task completes (one live bitmap 
per touched
      data file, fanout-style, no spill). Files-touched-per-batch grows with 
the table ->
      executor heap grows batch over batch -> GC stalls -> heartbeat evictions 
-> container OOM.
   3. **v2 path for comparison:** with `rewritableDeletes == null` the 
position-delete
      writers (`ClusteredPositionDeleteWriter` / 
`FanoutPositionOnlyDeleteWriter`) just
      append delete files bounded by the *current batch's* deletes. No prior 
state is
      loaded, which matches v2 being stable while accumulating far more delete 
files.
   
   In other words: DVs move work from readers to writers by design 
(supersede-on-write),
   but the Spark implementation holds the entire supersede working set in 
memory, per
   executor, with no bound and no spill, and additionally re-broadcasts the full
   file->deletes map every micro-batch. Under a streaming upsert workload whose 
scan
   touches most files, that working set is O(table data-file count), not 
O(batch).
   
   ### Expected behavior
   
   v3 DV writes should have memory consumption comparable to (or at least 
bounded like) v2
   positional-delete writes for the same MERGE workload - or expose a knob to 
bound/spill
   the DV merge working set. Doubling executor memory only delays onset 
(measured on EMR:
   56 GB -> 120 GB moved failure from ~20 to ~45 min).
   
   ### Reproduction (fully scripted, two commands)
   
   Everything is public in 
https://github.com/aws-samples/iceberg-streaming-examples
    a CDC feed generator (`KafkaCDCSimulator`: unkeyed, monotonic `seq`, 80/20 
hot keys, ~85% U / 15% D) and the streaming MERGE job 
(`SparkStreamingCDCMirror`, shared SQL in `CdcSql`).
   
   **Option A - Spark on EKS (what produced the numbers above):** 
`scripts/eks/` stands up
   everything (or adopts an existing cluster/bucket); then
   
   ```bash
   PRODUCER_CLASS=com.aws.emr.kafka.KafkaCDCSimulator ./scripts/eks/06-run.sh 
producer rate=100000
   JOB_CLASS=com.aws.emr.spark.cdc.SparkStreamingCDCMirror \
   EXECUTORS=6 EXECUTOR_CORES=5 EXECUTOR_MEMORY=16g \
     ./scripts/eks/06-run.sh ingest table=mirror_v3 fv=3 
startingoffsets=earliest
   # control: identical, with table=mirror_v2 fv=2
   ```
   
   **Option B - single machine:** `docker compose up -d` (Kafka), then
   `scripts/run-local.sh` with the same `fv=` knob; the memory growth pattern 
is visible at
   smaller scale in the executor metrics.
   
   Iceberg/Kafka/hadoop-aws jars and versions are pinned in 
`scripts/eks/Dockerfile`
   (iceberg 1.11.0, spark-sql-kafka 4.0.2, kafka-clients 3.9.1, hadoop-aws 
3.4.1,
   awssdk bundle 2.44.4). Spark event logs of the failing v3 run and the stable 
v2 control
   are preserved and can be attached on request.
   
   **Related**: #11122 (v3 position-delete design), #15924 (in-memory DV blob 
handling in
   path rewrite - same "whole DV in memory" pattern in a different code path).
   
   
   ### 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