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

   ### Apache Iceberg version
   
   _No response_
   
   ### Query engine
   
   _No response_
   
   ### Please describe the bug 🐞
   
   We recently set `usePrefixListing(true)` when we run 
`DeleteOrphanFilesSparkAction` which seemed to trigger an OutOfMemoryError for 
some of our tables which have a large(100+ million) number of candidate files 
returned by `listedFileDS()`.
   
   Our initial assessment is that this is likely being triggered by this code 
path in `listedFileDS()` hard coding the partition count in the returned 
Dataset to 1. It appears that this is causing ParallelCollectionRDD to attempt 
to place the file list in a single Java Array which is failing due to hitting 
the 2GiB limit for Arrays.
   ```
       if (usePrefixListing) {
         Preconditions.checkArgument(
             table.io() instanceof SupportsPrefixOperations,
             "Cannot use prefix listing with FileIO {} which does not support 
prefix operations.",
             table.io());
   
         Predicate<org.apache.iceberg.io.FileInfo> predicate =
             fileInfo -> fileInfo.createdAtMillis() < olderThanTimestamp;
         FileSystemWalker.listDirRecursivelyWithFileIO(
             (SupportsPrefixOperations) table.io(),
             location,
             table.specs(),
             predicate,
             matchingFiles::add);
   
         JavaRDD<String> matchingFileRDD = 
sparkContext().parallelize(matchingFiles, 1);
         return spark().createDataset(matchingFileRDD.rdd(), Encoders.STRING());
       } else {
   ```
   
   Example stack trace running on Spark 3.5.x:
   ```
   2026-07-21 11:16:58,534-0500 ERROR START:1784594222 [netty.Inbox] 
[dispatcher-CoarseGrainedScheduler] - {} - An error happened while processing 
message in the inbox for CoarseGrainedScheduler
   java.lang.OutOfMemoryError: Required array length 2147483639 + 489 is too 
large
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOfferSingleTaskSet$2$adapted(TaskSchedulerImpl.scala:409)
       at jdk.internal.util.ArraysSupport.hugeLength(ArraysSupport.java:649) 
~[?:?]
       at jdk.internal.util.ArraysSupport.newLength(ArraysSupport.java:642) 
~[?:?]
       at 
java.io.ByteArrayOutputStream.ensureCapacity(ByteArrayOutputStream.java:100) 
~[?:?]
       at java.io.ByteArrayOutputStream.write(ByteArrayOutputStream.java:130) 
~[?:?]
       at 
org.apache.spark.util.ByteBufferOutputStream.write(ByteBufferOutputStream.scala:41)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
java.io.ObjectOutputStream$BlockDataOutputStream.write(ObjectOutputStream.java:1862)
 ~[?:?]
       at java.io.ObjectOutputStream.write(ObjectOutputStream.java:714) ~[?:?]
       at scala.Option.foreach(Option.scala:407)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOfferSingleTaskSet$1(TaskSchedulerImpl.scala:409)
       at org.apache.spark.util.Utils$$anon$1.write(Utils.scala:143) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at com.esotericsoftware.kryo.io.Output.flush(Output.java:185) 
~[kryo-shaded-4.0.2.jar:?]
       at com.esotericsoftware.kryo.io.Output.close(Output.java:196) 
~[kryo-shaded-4.0.2.jar:?]
       at 
org.apache.spark.serializer.KryoSerializationStream.close(KryoSerializer.scala:292)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.util.Utils$.serializeViaNestedStream(Utils.scala:148) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.rdd.ParallelCollectionPartition.$anonfun$writeObject$1(ParallelCollectionRDD.scala:64)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) 
~[scala-library-2.12.18.jar:?]
       at 
org.apache.spark.util.SparkErrorUtils.tryOrIOException(SparkErrorUtils.scala:35)
 ~[spark-common-utils_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.util.SparkErrorUtils.tryOrIOException$(SparkErrorUtils.scala:33)
 ~[spark-common-utils_2.12-3.5.3.jar:3.5.3]
       at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:94) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.rdd.ParallelCollectionPartition.writeObject(ParallelCollectionRDD.scala:50)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 
~[?:?]
       at 
jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
 ~[?:?]
       at 
jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
 ~[?:?]
       at java.lang.reflect.Method.invoke(Method.java:569) ~[?:?]
       at 
java.io.ObjectStreamClass.invokeWriteObject(ObjectStreamClass.java:1070) ~[?:?]
       at scala.collection.immutable.Range.foreach$mVc$sp(Range.scala:158)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.resourceOfferSingleTaskSet(TaskSchedulerImpl.scala:399)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$20(TaskSchedulerImpl.scala:606)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$20$adapted(TaskSchedulerImpl.scala:601)
       at 
scala.collection.IndexedSeqOptimized.foreach(IndexedSeqOptimized.scala:36)
       at 
scala.collection.IndexedSeqOptimized.foreach$(IndexedSeqOptimized.scala:33)
       at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:198)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$16(TaskSchedulerImpl.scala:601)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$16$adapted(TaskSchedulerImpl.scala:574)
       at 
scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
       at 
scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
       at 
java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1516) ~[?:?]
       at 
java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1438) 
~[?:?]
       at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1181) 
~[?:?]
       at 
java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1572) 
~[?:?]
       at 
java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1529) ~[?:?]
       at 
java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1438) 
~[?:?]
       at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1181) 
~[?:?]
       at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:350) 
~[?:?]
       at 
org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:46)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:115)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.TaskSetManager.prepareLaunchingTask(TaskSetManager.scala:530)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.TaskSetManager.$anonfun$resourceOffer$2(TaskSetManager.scala:494)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at scala.Option.map(Option.scala:230) ~[scala-library-2.12.18.jar:?]
       at 
org.apache.spark.scheduler.TaskSetManager.resourceOffer(TaskSetManager.scala:470)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOfferSingleTaskSet$2(TaskSchedulerImpl.scala:414)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.resourceOffers(TaskSchedulerImpl.scala:574)
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend$DriverEndpoint.$anonfun$makeOffers$4(CoarseGrainedSchedulerBackend.scala:403)
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOfferSingleTaskSet$2$adapted(TaskSchedulerImpl.scala:409)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at scala.Option.foreach(Option.scala:407) ~[scala-library-2.12.18.jar:?]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOfferSingleTaskSet$1(TaskSchedulerImpl.scala:409)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend.org$apache$spark$scheduler$cluster$CoarseGrainedSchedulerBackend$$withLock(CoarseGrainedSchedulerBackend.scala:1058)
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend$DriverEndpoint.org$apache$spark$scheduler$cluster$CoarseGrainedSchedulerBackend$DriverEndpoint$$makeOffers(CoarseGrainedSchedulerBackend.scala:400)
       at scala.collection.immutable.Range.foreach$mVc$sp(Range.scala:158) 
~[scala-library-2.12.18.jar:?]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.resourceOfferSingleTaskSet(TaskSchedulerImpl.scala:399)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$20(TaskSchedulerImpl.scala:606)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$20$adapted(TaskSchedulerImpl.scala:601)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
scala.collection.IndexedSeqOptimized.foreach(IndexedSeqOptimized.scala:36) 
~[scala-library-2.12.18.jar:?]
       at 
scala.collection.IndexedSeqOptimized.foreach$(IndexedSeqOptimized.scala:33) 
~[scala-library-2.12.18.jar:?]
       at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:198) 
~[scala-library-2.12.18.jar:?]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$16(TaskSchedulerImpl.scala:601)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.$anonfun$resourceOffers$16$adapted(TaskSchedulerImpl.scala:574)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) 
~[scala-library-2.12.18.jar:?]
       at 
scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) 
~[scala-library-2.12.18.jar:?]
       at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) 
~[scala-library-2.12.18.jar:?]
       at 
org.apache.spark.scheduler.TaskSchedulerImpl.resourceOffers(TaskSchedulerImpl.scala:574)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend$DriverEndpoint.$anonfun$makeOffers$4(CoarseGrainedSchedulerBackend.scala:403)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend.org$apache$spark$scheduler$cluster$CoarseGrainedSchedulerBackend$$withLock(CoarseGrainedSchedulerBackend.scala:1058)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend$DriverEndpoint.org$apache$spark$scheduler$cluster$CoarseGrainedSchedulerBackend$DriverEndpoint$$makeOffers(CoarseGrainedSchedulerBackend.scala:400)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend$DriverEndpoint$$anonfun$receive$1.applyOrElse(CoarseGrainedSchedulerBackend.scala:232)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:115) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:213) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:100) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.rpc.netty.MessageLoop.org$apache$spark$rpc$netty$MessageLoop$$receiveLoop(MessageLoop.scala:75)
 ~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
org.apache.spark.rpc.netty.MessageLoop$$anon$1.run(MessageLoop.scala:41) 
~[spark-core_2.12-3.5.3.jar:3.5.3]
       at 
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) 
~[?:?]
       at 
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) 
~[?:?]
       at java.lang.Thread.run(Thread.java:840) [?:?]
       at 
org.apache.spark.scheduler.cluster.CoarseGrainedSchedulerBackend$DriverEndpoint$$anonfun$receive$1.applyOrElse(CoarseGrainedSchedulerBackend.scala:232)
       at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:115)
       at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:213)
       at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:100)
       at 
org.apache.spark.rpc.netty.MessageLoop.org$apache$spark$rpc$netty$MessageLoop$$receiveLoop(MessageLoop.scala:75)
       at 
org.apache.spark.rpc.netty.MessageLoop$$anon$1.run(MessageLoop.scala:41)
       at 
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
       at 
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
       at java.base/java.lang.Thread.run(Thread.java:840)
   ```
   
   For now we worked around the issue by building a Dataset in the same manner 
as the usePrefixListing path in `listedFileDS()` except for using multiple 
Spark partitions and passing it to the Action with `compareToFileList()`. That 
feels like evidence that the root issue is as we suspected.
   
   Hopefully fixing that could be a minor change. Please let us know if we can 
provide any additional supporting information.
   
   ### 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
   - [x] 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