Skip to content

[Core] Add compaction.task-threads for multi-thread async compaction per write subtask - #10133

Open
jacklong319 wants to merge 14 commits into
apache:masterfrom
jacklong319:master
Open

jacklong319 wants to merge 14 commits into
apache:masterfrom
jacklong319:master

Conversation

@jacklong319

@jacklong319 jacklong319 commented Sep 23, 2026 •

Copy link
Copy Markdown
Contributor

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-threads to 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:

  • Add CompactionTaskExecutorMode and CoreOptions.COMPACTION_TASK_THREADS.
  • Route executors in AbstractFileStoreWrite.compactExecutor(partition, bucket); release per-bucket executors on writer cleanup where applicable.
  • Default remains 1; no storage format change.

Tests

@JingsongLi

Copy link
Copy Markdown
Contributor

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’ taskFuture guard; the default remains a single shared executor, and the managers appear to keep each bucket’s submitted task serialized.

Merge blockers / production follow-up:

  1. A normal JDK 8 mvn -o -pl paimon-api,paimon-core -DskipTests compile fails Spotless in AbstractFileStoreWrite.java at the perBucketCompactExecutors.remove(...) line. This must be formatted; the PR currently has 14 failing checks.
  2. There are no new tests for FIXED_POOL or PER_BUCKET, cross-bucket concurrency, same-bucket serialization, executor release/recreation, or external executor ownership. The existing targeted compaction/write classes pass locally (29 tests run, 1 skipped), but they exercise mostly the default path. Please add focused behavioral tests before production merge.
  3. Please document the thread/memory trade-off and monitoring behavior described in [Feature] Configurable multi-thread async compaction per Flink write subtask (compaction.task-threads) #10132; compactionThreadBusy can exceed 100 with parallel workers. Also reject unsupported values such as 0 and -2 rather than silently treating them as SINGLE.

The feature has clear value, but the current head is not ready to merge.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@jacklong319

Copy link
Copy Markdown
Contributor Author

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

Copy link
Copy Markdown
Contributor

Reviewed latest head da05812 for production. The linked DV lag issue gives this feature clear end-to-end value; the latest commit addresses the previously reported per-bucket timer retention and adds worker-churn/shared-pool tests. Local JDK 8 verification passed: four focused classes (11 tests), existing compaction/write classes (26 tests), and normal paimon-api,paimon-core compile with Spotless/Checkstyle.

Two merge blockers remain:

  1. P1 — generated configuration docs are missing. Both Core CI jobs fail in ConfigOptionsDocsCompletenessITCase: Option compaction.task-threads in class org.apache.paimon.CoreOptions is not documented. The prose docs were updated, but docs/generated/core_configuration.html was not regenerated. Please regenerate it per paimon-docs/README.md and rerun CI. git diff --check also reports trailing whitespace in docs/docs/primary-key-table/table-mode.md:109.

  2. P1 — timer retirement can race with the next reporter on a shared worker. In ReporterImpl.getCompactTimer(), compactTimers.computeIfAbsent(threadId, ...) happens before acquireCompactTimer(threadId), while unregister() can remove the timer in between. A concrete interleaving: bucket B obtains bucket A's worker timer, bucket A unregisters and removes the last reference, then B acquires the reference and calls start() on the now-detached timer. CompactTask.stopTimer() calls getCompactTimer() again, which creates a fresh timer and calls finish() without a matching start(); CompactTimer.finish() throws IllegalArgumentException. This is reachable with the fixed pool/default shared worker when an idle bucket writer is retired while another bucket starts compaction. Make timer lookup/acquisition/removal atomic for a thread ID, and add a concurrent retirement/start regression test.

The option is worth keeping open, but these issues need fixing before production merge.

@jacklong319

Copy link
Copy Markdown
Contributor Author

Reviewed latest head da05812 for production. The linked DV lag issue gives this feature clear end-to-end value; the latest commit addresses the previously reported per-bucket timer retention and adds worker-churn/shared-pool tests. Local JDK 8 verification passed: four focused classes (11 tests), existing compaction/write classes (26 tests), and normal paimon-api,paimon-core compile with Spotless/Checkstyle.

Two merge blockers remain:

  1. P1 — generated configuration docs are missing. Both Core CI jobs fail in ConfigOptionsDocsCompletenessITCase: Option compaction.task-threads in class org.apache.paimon.CoreOptions is not documented. The prose docs were updated, but docs/generated/core_configuration.html was not regenerated. Please regenerate it per paimon-docs/README.md and rerun CI. git diff --check also reports trailing whitespace in docs/docs/primary-key-table/table-mode.md:109.
  2. P1 — timer retirement can race with the next reporter on a shared worker. In ReporterImpl.getCompactTimer(), compactTimers.computeIfAbsent(threadId, ...) happens before acquireCompactTimer(threadId), while unregister() can remove the timer in between. A concrete interleaving: bucket B obtains bucket A's worker timer, bucket A unregisters and removes the last reference, then B acquires the reference and calls start() on the now-detached timer. CompactTask.stopTimer() calls getCompactTimer() again, which creates a fresh timer and calls finish() without a matching start(); CompactTimer.finish() throws IllegalArgumentException. This is reachable with the fixed pool/default shared worker when an idle bucket writer is retired while another bucket starts compaction. Make timer lookup/acquisition/removal atomic for a thread ID, and add a concurrent retirement/start regression test.

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:

  1. Generated configuration docs — Regenerated docs/generated/core_configuration.html so compaction.task-threads is documented for ConfigOptionsDocsCompletenessITCase. Also removed trailing whitespace in docs/docs/primary-key-table/table-mode.md (git diff --check).

  2. CompactTimer race on shared workers — ReporterImpl.getCompactTimer() and timer retirement now run under the same per-threadId lock, so lookup, ref acquire, and removal are atomic. Reporters still release their compaction-thread timers on unregister() with ref-counting when multiple buckets share one worker thread. Added regression tests in CompactionMetricsTest for PER_BUCKET churn, shared-thread ref retention, and shared-worker handoff after unregister.

Could you please take another look when CI is green?

@JingsongLi

Copy link
Copy Markdown
Contributor

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 (AbstractFileStoreWrite.java:800–802, also the PER_BUCKET path). Multiple workers now call CompactionMetrics.ReporterImpl.increaseCompactionsCompletedCount() concurrently. In Flink, FlinkMetricGroup.counter() delegates to the default Flink SimpleCounter, whose long increment/decrement is not atomic. The core TestMetricRegistry uses an AtomicLong counter, so the new concurrency tests hide this production difference.

I reproduced this through the real FlinkMetricRegistry/FlinkMetricGroup adapters and actual CompactTask.call() updates on Flink 1.20.4. One worker completed 10,000 tasks and reported exactly 10,000 / queued 0. Four workers completed 40,000 tasks, but three runs reported completion counts 39,952, 39,986, and 39,913, with queued counts 5,607, 1,972, and 1,461 after all tasks finished. The new multi-worker completion-count loss is the introduced regression; the queued counter already had a main-thread/worker race, which this change also exposes under greater concurrency. Flink 2.2's default counter has the same implementation.

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.

@jacklong319

Copy link
Copy Markdown
Contributor Author

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 (AbstractFileStoreWrite.java:800–802, also the PER_BUCKET path). Multiple workers now call CompactionMetrics.ReporterImpl.increaseCompactionsCompletedCount() concurrently. In Flink, FlinkMetricGroup.counter() delegates to the default Flink SimpleCounter, whose long increment/decrement is not atomic. The core TestMetricRegistry uses an AtomicLong counter, so the new concurrency tests hide this production difference.

I reproduced this through the real FlinkMetricRegistry/FlinkMetricGroup adapters and actual CompactTask.call() updates on Flink 1.20.4. One worker completed 10,000 tasks and reported exactly 10,000 / queued 0. Four workers completed 40,000 tasks, but three runs reported completion counts 39,952, 39,986, and 39,913, with queued counts 5,607, 1,972, and 1,461 after all tasks finished. The new multi-worker completion-count loss is the introduced regression; the queued counter already had a main-thread/worker race, which this change also exposes under greater concurrency. Flink 2.2's default counter has the same implementation.

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 4344ed7.

For the P2 shared compaction counters: multiple fixed-pool workers could update the same Flink SimpleCounter concurrently and lose counts. I synchronized shared counter updates in CompactionMetrics and added CompactionMetricsTest#testSharedCompactionCountersWithConcurrentNonAtomicBackend using a non-thread-safe counter backend (similar to Flink's SimpleCounter).

Please take another look when CI is green.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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");

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

@jacklong319

Copy link
Copy Markdown
Contributor Author

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.

@JingsongLi Thanks for the detailed review on the parallel compaction findings. I pushed a follow-up that addresses the P1/P2 items while keeping compaction.task-threads=-1 (PER_BUCKET).

[P1] Append compaction reader / cast state

  • Root cause: BaseAppendFileStoreWrite shared a single RawFileSplitRead whose cached FormatReaderMapping (including nested cast state) was reused across parallel compaction workers.
  • Fix: for FIXED_POOL and PER_BUCKET, each compaction worker thread uses its own RawFileSplitRead copy with a fresh mapping cache (copyWithFreshReaderMappings). SINGLE keeps the previous shared reader (no extra overhead).
  • Regression: ParallelAppendCompactionReaderIsolationTest — two buckets run compactRewrite in parallel after nested schema evolution (payload.val INT → BIGINT), with task-threads=2 and -1; each bucket’s marker values stay isolated.

[P2] compactionThreadBusy / CompactTimer lifetime

  • Fix: retire CompactTimer on reporter unregister() only when the write uses PER_BUCKET (dedicated per-bucket workers). For SINGLE and FIXED_POOL, timers are not removed on bucket reporter unregister, so the 60s busy window is preserved for long-lived shared workers.
  • Removed ref-count bookkeeping on the default path; updated CompactionMetricsTest accordingly.

PER_BUCKET executor routing and earlier counter synchronization (sharedCounterLock) are unchanged.

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Feature] Configurable multi-thread async compaction per Flink write subtask (compaction.task-threads)

2 participants