[Core] Add compaction.task-threads for multi-thread async compaction per write subtask - #10133
jacklong319 wants to merge 14 commits into
Conversation
|
The linked production issue provides a strong end-to-end case for cross-bucket compaction parallelism (reported DV read lag improved from ~40 min at 1 thread to ~8 min at 3). I traced the executor routing and the per-bucket compact managers’ Merge blockers / production follow-up:
The feature has clear value, but the current head is not ready to merge. |
(Same body bullets as above.) Related to apache#10132
JingsongLi
left a comment
There was a problem hiding this comment.
Reviewed head 072f895 for production use. The end-to-end value is clear: allowing buckets within one long-lived Flink write subtask to compact concurrently can reduce the MOW/DV visibility lag described in #10132. The update addresses the previous feedback on invalid values, documentation, focused executor tests, and Spotless. Local JDK 8 verification passed: the three new test classes 10/10; existing MergeTreeCompactManagerTest, KeyValueFileStoreWriteTest, and BucketedAppendFileStoreWriteTest 26/26; and ordinary mvn -pl paimon-api,paimon-core -DskipTests compile including Spotless. Current full CI is still running.
P1 — retire compaction timers when PER_BUCKET workers are retired. AbstractFileStoreWrite.releaseCompactionExecutor shuts down and removes each bucket executor, but CompactionMetrics.ReporterImpl.getCompactTimer() caches a CompactTimer by thread ID in compactTimers and never removes that entry; ReporterImpl.unregister() removes only the reporter. In a long-lived Flink subtask with changing active partitions/buckets, each replacement worker has a new thread ID. The compactionThreadBusy gauge scans every historical timer, so both retained state and scrape cost grow with total buckets ever compacted, even after their writers are gone. I reproduced this with 32 sequential single-thread workers: after shutting each down and unregistering its reporter, the timer map still held all 32 entries (expected 0). The temporary probe was removed from the review checkout. Please bound or retire these timers on worker shutdown and add a churn regression test for PER_BUCKET; the current routing tests only check executor removal.
This is a feature-specific blocker for the new -1 mode. The default single-thread path and fixed-size pool do not create the same unbounded worker churn.
Thanks @JingsongLi — agreed on the P1 for PER_BUCKET timer retention. I’ll fix timer retirement on worker shutdown, add the churn regression test, and update the PR soon. |
Release CompactTimer entries on reporter unregister with ref-counting for shared compaction threads. Add churn regression tests mimicking PER_BUCKET worker rotation. Related to apache#10132
|
Reviewed latest head Two merge blockers remain:
The option is worth keeping open, but these issues need fixing before production merge. |
Thank you, @JingsongLi, for affirming the value of this PR—it is encouraging to see that recognition from the community. I pushed another update to address the two P1 items on the latest head:
Could you please take another look when CI is green? |
|
Reviewed head 4344ed7. The feature has clear production value, and the previous generated-doc and timer-retirement feedback is addressed. Normal JDK 8 verification passed 42 focused/adjacent tests, including worker churn, shared-worker handoff, executor ownership, and bucket serialization. I also ran actual write/update/delete -> compaction -> commit -> read workflows for all three modes with and without deletion vectors: observed concurrency was 1 / 2 / 4 respectively, each bucket stayed serialized, all six workflows returned the correct 16 rows, and rewriters were closed. Generated option documentation, whitespace checks, and the exact-head CI are green. [P2] Make shared compaction counters safe for the new worker concurrency ( I reproduced this through the real Please synchronize updates to these shared counters or register a thread-safe counter through the Flink adapter, and add a regression using the actual production metric backend. This corrupts the completion/queue signals operators are explicitly told to use when sizing the new pool; it did not cause data loss or compaction failure in my tests. |
Thank you, @JingsongLi, for reviewing For the P2 shared compaction counters: multiple fixed-pool workers could update the same Flink Please take another look when CI is green. |
JingsongLi
left a comment
There was a problem hiding this comment.
Reviewed head 0107e685 for production use. Requirement fit: SUPPORTED; implementation: FINDINGS.
The cross-bucket capacity knob addresses a credible production bottleneck. The default still uses one shared compaction thread, and per-bucket task serialization remains in place. However, the new parallel modes expose a shared-reader data corruption path for bucketed append tables, and timer retirement changes the busy metric even with the default configuration; see the inline comments.
Verification: normal JDK 8 build and 47 focused/adjacent tests passed, including formatting checks. Additional deterministic probes reproduced both findings. The append probe used actual Parquet input and output files and the configured executors; a latch controlled the interleaving around the existing cached caster without changing the returned values. One thread preserved both rows; two threads and per-bucket executors persisted an incorrect nested value.
A smaller initial design would support only positive thread counts: keep the existing executor for 1 and use one bounded shared pool for N > 1. This removes the per-bucket executor map, key, release/recreation lifecycle, mode enum, and the need to retire timers for churned workers. Keep the shared-counter synchronization. This simplification still needs isolated append reader/cast state before enabling append concurrency; alternatively, initially restrict the feature to the validated primary-key compaction path. I would resolve the P1 before deploying either parallel mode.
| + "-compaction-bucket-" | ||
| + key.bucket))); | ||
| case FIXED_POOL: | ||
| return sharedCompactionExecutor(compactionTaskThreads, "-compaction-pool"); |
There was a problem hiding this comment.
[P1] Isolate append compaction reader state before enabling cross-bucket concurrency
This routing also applies to BucketedAppendFileStoreWrite. Its bucket rewriters share BaseAppendFileStoreWrite.readForCompact, whose RawFileSplitRead caches one FormatReaderMapping per (schemaId, format). The cached CastFieldGetter[] contains nested ROW/ARRAY/MAP casts that capture mutable CastedRow/CastedArray/CastedMap instances. Each DataFileRecordReader gets a separate outer row wrapper but shares these nested cast instances. Another bucket can therefore replace a nested wrapper's backing value while the first worker serializes it, silently writing the other bucket's data into the compacted file.
I reproduced this on JDK 8 with an actual two-bucket append Parquet table: (id=10, tags=[10]) and (id=20, tags=[20]), followed by the supported schema change tags.element: INT -> BIGINT. A latch-controlled interleaving around the real cached caster produced output files containing (10, [20]) and (20, [20]) with both compaction.task-threads=2 and -1; the single-thread control preserved (10, [10]) and (20, [20]).
Please give each bucket/worker its own compaction reader and mutable cast state, or construct independent nested casts for each actual file reader. Making the cache a ConcurrentHashMap alone does not fix the shared mutable wrappers. Add a regression that checks rewritten nested values across buckets after schema evolution; the executor-routing tests do not exercise this path.
| threadId, | ||
| (id, count) -> { | ||
| if (count == null || count <= 1) { | ||
| compactTimers.remove(id); |
There was a problem hiding this comment.
[P2] Preserve the busy window when a shared worker remains alive
Removing the timer when its last bucket reporter unregisters also changes the default compaction.task-threads=1 behavior. An idle writer can be cleaned by prepareCommit shortly after its compaction finishes, while the shared executor remains alive. This deletion discards that worker's recent activity immediately instead of letting compactionThreadBusy decay over its existing 60-second window.
A deterministic timer probe recording 30 seconds of recent activity reports busy=50 before unregister. The PR head reports 0 immediately afterward; the base still reports 50. Rolling/sparse partitions can therefore under-report compaction utilization even without opting into parallelism.
Please retain timers for the lifetime of SINGLE/FIXED_POOL workers, and limit retirement to workers that actually leave, such as dedicated PER_BUCKET executors. This also avoids introducing reporter-based timer reference bookkeeping into the unchanged default path.
@JingsongLi Thanks for the detailed review on the parallel compaction findings. I pushed a follow-up that addresses the P1/P2 items while keeping [P1] Append compaction reader / cast state
[P2]
PER_BUCKET executor routing and earlier counter synchronization ( Could you take another look when CI is green? Happy to follow up separately if you see gaps on PK merge-tree paths or external compact executor combinations. |
Purpose
Linked issue: close #10132
Today each Flink write subtask uses one shared async compaction thread for all
(partition, bucket)writers. Under peak load, compaction falls behind; in Deletion Vectors (MOW) mode, delayed L0 compaction hurts read freshness (see #10132).This PR adds table option
compaction.task-threadsto run parallel async compaction across buckets within the same write subtask, while compaction for the same bucket remains serialized.Semantics:
1(default):SINGLE— unchanged behavior.N > 1:FIXED_POOL— N threads shared by buckets in the subtask.-1:PER_BUCKET— one thread per active(partition, bucket)writer (higher memory use).Changes:
CompactionTaskExecutorModeandCoreOptions.COMPACTION_TASK_THREADS.AbstractFileStoreWrite.compactExecutor(partition, bucket); release per-bucket executors on writer cleanup where applicable.1; no storage format change.Tests
mvn -pl paimon-api,paimon-core -am -Pfast-build test