andygrove commented on PR #6773:
URL: 
https://github.com/apache/datafusion-comet/pull/6773#issuecomment-6083434253

   > Could there be too many small files when the rows' partitions are 
interleaved (e.g. p0, p1, ... pN, p0, p1, ... pN, p0, p1, ...)?
   
   Yes, that's the case where it happens, once the pool can't hold the write: 
every partition keeps getting rows, so a partition closed early opens a new 
file on its next rows. When the pool holds the write, nothing closes early and 
each partition gets one file, as before.
   
   How many files per partition depends on how far short the pool falls, not on 
how many partitions there are. I measured it with rows taking turns across 
partitions as in your example (40 batches of 4096 rows with 100-byte payloads, 
about 18 MB), and the pool set to a fraction of the most the same write 
reserves with room to spare. Files per partition:
   
   | pool | 16 partitions | 64 partitions | 256 partitions |
   | ---- | ------------- | ------------- | -------------- |
   | 1/2  | 1.69          | 1.66          | 1.66           |
   | 1/4  | 2.62          | 2.55          | 2.56           |
   | 1/8  | 4.38          | 4.31          | 4.27           |
   | 1/16 | 7.44          | 7.36          | 7.29           |
   
   So a pool 1/k of what the write needs gives roughly (k + 1) / 2 files per 
partition. Closing the partition holding the most is what keeps it that low, 
since that frees as much as one close can. At 1/8, closing every partition at 
once gives 7 files per partition, and closing them in first-seen order, which 
is smallest first here, gives 17 to 25. Nothing tested that order, so 
bb1788fc35 adds 
`a_fanout_write_short_of_memory_closes_the_partitions_holding_the_most`, which 
fails under all three.
   
   Before this PR the same write failed its task, and the retry failed the same 
way. iceberg-java's fanout writer keeps every file open until the task ends, so 
it writes one file per partition as long as the JVM heap holds all of them. 
Where the extra files matter, `write.spark.fanout.enabled=false` makes Spark 
sort each task's rows by partition, spilling as it needs to, and the task keeps 
one file open at a time. More off-heap memory, or `rewrite_data_files` 
afterwards, also work. The user guide now says roughly how many files to expect.
   
   Avoiding the small files without more memory would take what Spark's own 
file writer does past `spark.sql.maxConcurrentOutputFileWriters`: sort the rest 
of the task's rows by partition and write them one partition at a time. That's 
a much bigger change, so I'd leave it for a follow-up if it turns out to matter.
   


-- 
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