[Feature] Support maximum compaction concurrency control in Compact Database Action
- Dominant language
- Java
- Stars
- 3.4k
- Forks
- 1.4k
- Avg merge
- 1d 11h
- Merged PRs (30d)
- 396
Description
### Search before asking
- [X] I searched in the [issues](https://github.com/apache/paimon/issues) and found nothing similar.
### Motivation
When multiple paimon tables are compacting by a single compact database flink job, there might have a sudden surge of the memory usage or disk usage in Flink job when these tables are producing new files simultaneously.
This may lead to TM OOM or 'No space left on device'
`Caused by: java.util.concurrent.ExecutionException: java.io.IOException: No space left on device
at java.util.concurrent.FutureTask.report(FutureTask.java:122)
at java.util.concurrent.FutureTask.get(FutureTask.java:192)
at org.apache.paimon.compact.CompactFutureManager.obtainCompactResult(CompactFutureManager.java:67)
at org.apache.paimon.compact.CompactFutureManager.innerGetCompactionResult(CompactFutureManager.java:53)
at org.apache.paimon.mergetree.compact.MergeTreeCompactManager.getCompactionResult(MergeTreeCompactManager.java:191)
at org.apache.paimon.mergetree.MergeTreeWriter.trySyncLatestCompaction(MergeTreeWriter.java:304)
at org.apache.paimon.mergetree.MergeTreeWriter.prepareCommit(MergeTreeWriter.java:244)
at org.apache.paimon.operation.AbstractFileStoreWrite.prepareCommit(AbstractFileStoreWrite.java:196)
at org.apache.paimon.table.sink.TableWriteImpl.prepareCommit(TableWriteImpl.java:176)
at org.apache.paimon.flink.sink.StoreSinkWriteImpl.prepareCommit(StoreSinkWriteImpl.java:208)
... 31 more
Caused by: java.io.IOException: No space left on device
at sun.nio.ch.FileDispatcherImpl.write0(Native Method)
at sun.nio.ch.FileDispatcherImpl.write(FileDispatcherImpl.java:60)
at sun.nio.ch.IOUtil.writeFromNativeBuffer(IOUtil.java:93)
at sun.nio.ch.IOUtil.write(IOUtil.java:51)
at sun.nio.ch.FileChannelImpl.write(FileChannelImpl.java:211)
at org.apache.paimon.utils.FileIOUtils.writeCompletely(FileIOUtils.java:59)
at org.apache.paimon.disk.BufferFileWriterImpl.writeBlock(BufferFileWriterImpl.java:41)
at org.apache.paimon.disk.ChannelWriterOutputView.writeCompressed(ChannelWriterOutputView.java:82)
at org.apache.paimon.disk.ChannelWriterOutputView.close(ChannelWriterOutputView.java:65)
at org.apache.paimon.sort.AbstractBinaryExternalMerger.mergeChannels(AbstractBinaryExternalMerger.java:190)
at org.apache.paimon.sort.AbstractBinaryExternalMerger.mergeChannelList(AbstractBinaryExternalMerger.java:153)
at org.apache.paimon.sort.BinaryExternalSortBuffer.write(BinaryExternalSortBuffer.java:181)
at org.apache.paimon.mergetree.MergeSorter$ExternalSorterWithLevel.put(MergeSorter.java:213)
at org.apache.paimon.reader.RecordReader.forIOEachRemaining(RecordReader.java:156)
at org.apache.paimon.mergetree.MergeSorter.spillMergeSort(MergeSorter.java:131)
at org.apache.paimon.mergetree.MergeSorter.mergeSort(MergeSorter.java:107)
at org.apache.paimon.mergetree.MergeTreeReaders.readerForSection(MergeTreeReaders.java:79)
at org.apache.paimon.mergetree.MergeTreeReaders.lambda$readerForMergeTree$0(MergeTreeReaders.java:54)
at org.apache.paimon.mergetree.compact.ConcatRecordReader.create(ConcatRecordReader.java:50)
at org.apache.paimon.mergetree.MergeTreeReaders.readerForMergeTree(MergeTreeReaders.java:61)
at org.apache.paimon.mergetree.compact.MergeTreeCompactRewriter.rewriteCompaction(MergeTreeCompactRewriter.java:70)
at org.apache.paimon.mergetree.compact.MergeTreeCompactRewriter.rewrite(MergeTreeCompactRewriter.java:62)
at org.apache.paimon.mergetree.compact.MergeTreeCompactTask.rewrite(MergeTreeCompactTask.java:133)
at org.apache.paimon.mergetree.compact.MergeTreeCompactTask.doCompact(MergeTreeCompactTask.java:97)
at org.apache.paimon.compact.CompactTask.call(CompactTask.java:49)
at org.apache.paimon.compact.CompactTask.call(CompactTask.java:35)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)`
### Solution
_No response_
### Anything else?
_No response_
### Are you willing to submit a PR?
- [ ] I'm willing to submit a PR!
Contributor guide
No contributing guide indexed for this repository
Research direction
Start by tracing compaction scheduling from CompactFutureManager and CompactTask through the Flink compact database action. Determine where concurrent compactions across tables are launched and how a maximum can be configured. Done means the action enforces the configured concurrency limit and coverage demonstrates that simultaneous compactions are bounded.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- data-engineering, performance
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100