apache / apache/paimon

[Feature] Support maximum compaction concurrency control in Compact Database Action

Open
#3,466 0 comments 0 reactions 0 assignees View on GitHub
enhancement
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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.