Skip to content

Commit b542e57

Browse files
authored
feat(compaction): support compaction for key table in framework (alibaba#195)
1 parent 8f4102a commit b542e57

63 files changed

Lines changed: 2802 additions & 389 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

include/paimon/defs.h

Lines changed: 37 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -203,6 +203,43 @@ struct PAIMON_EXPORT Options {
203203
/// "commit.max-retries" - Maximum number of retries when commit failed. Default value is 10.
204204
static const char COMMIT_MAX_RETRIES[];
205205

206+
/// "compaction.max-size-amplification-percent" - The size amplification is defined as the
207+
/// amount (in percentage) of additional storage needed to store a single byte of data in the
208+
/// merge tree for changelog mode table. Default value is 200.
209+
static const char COMPACTION_MAX_SIZE_AMPLIFICATION_PERCENT[];
210+
211+
/// "compaction.size-ratio" - Percentage flexibility while comparing sorted run size for
212+
/// changelog mode table. If the candidate sorted run(s) size is 1% smaller than the next
213+
/// sorted run's size, then include next sorted run into this candidate set. Default value is 1.
214+
static const char COMPACTION_SIZE_RATIO[];
215+
216+
/// "num-sorted-run.compaction-trigger" - The sorted run number to trigger compaction. Includes
217+
/// level0 files (one file one sorted run) and high-level runs (one level one sorted run).
218+
/// Default value is 5.
219+
static const char NUM_SORTED_RUNS_COMPACTION_TRIGGER[];
220+
221+
/// "num-sorted-run.stop-trigger" - The number of sorted runs that trigger the stopping of
222+
/// writes, the default value is 'num-sorted-run.compaction-trigger' + 3.
223+
static const char NUM_SORTED_RUNS_STOP_TRIGGER[];
224+
225+
/// "num-levels" - Total level number, for example, there are 3 levels, including 0,1,2 levels.
226+
/// No default value.
227+
static const char NUM_LEVELS[];
228+
229+
/// "lookup-compact" - Lookup compact mode used for lookup compaction. Default value is
230+
/// LookupCompactMode::RADICAL.
231+
static const char LOOKUP_COMPACT[];
232+
233+
/// "compaction.force-up-level-0" - If set to true, compaction strategy will always include all
234+
/// level 0 files in candidates. Default value is false.
235+
static const char COMPACTION_FORCE_UP_LEVEL_0[];
236+
237+
/// "lookup-compact.max-interval" - The max interval for a gentle mode lookup compaction to be
238+
/// triggered. For every interval, a forced lookup compaction will be performed to flush L0
239+
/// files to higher level. This option is only valid when lookup-compact mode is gentle. No
240+
/// default value.
241+
static const char LOOKUP_COMPACT_MAX_INTERVAL[];
242+
206243
/// "sequence.field" - The field that generates the sequence number for primary key table, the
207244
/// sequence number determines which data is the most recent. Value use "," as delimiter.
208245
static const char SEQUENCE_FIELD[];
@@ -304,7 +341,6 @@ struct PAIMON_EXPORT Options {
304341
static const char SCAN_TAG_NAME[];
305342
/// "write-only" - If set to "true", compactions and snapshot expiration will be skipped. This
306343
/// option is used along with dedicated compact jobs. Default value is "false".
307-
/// @note: This option will be ignore until compaction is supported.
308344
static const char WRITE_ONLY[];
309345
/// "compaction.min.file-num" - For file set [f_0,...,f_N], the minimum file number to trigger a
310346
/// compaction for append-only table. Default value is 5.

include/paimon/write_context.h

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929

3030
namespace paimon {
3131
class Executor;
32+
class IOManager;
3233
class MemoryPool;
3334

3435
/// `WriteContext` is some configuration for write operations.
@@ -44,6 +45,7 @@ class PAIMON_EXPORT WriteContext {
4445
const std::vector<std::string>& write_schema,
4546
const std::shared_ptr<MemoryPool>& memory_pool,
4647
const std::shared_ptr<Executor>& executor,
48+
const std::shared_ptr<IOManager>& io_manager,
4749
const std::shared_ptr<FileSystem>& specific_file_system,
4850
const std::map<std::string, std::string>& fs_scheme_to_identifier_map,
4951
const std::map<std::string, std::string>& options);
@@ -98,6 +100,10 @@ class PAIMON_EXPORT WriteContext {
98100
return executor_;
99101
}
100102

103+
std::shared_ptr<IOManager> GetIOManager() const {
104+
return io_manager_;
105+
}
106+
101107
std::shared_ptr<FileSystem> GetSpecificFileSystem() const {
102108
return specific_file_system_;
103109
}
@@ -113,6 +119,7 @@ class PAIMON_EXPORT WriteContext {
113119
std::vector<std::string> write_schema_;
114120
std::shared_ptr<MemoryPool> memory_pool_;
115121
std::shared_ptr<Executor> executor_;
122+
std::shared_ptr<IOManager> io_manager_;
116123
std::shared_ptr<FileSystem> specific_file_system_;
117124
std::map<std::string, std::string> fs_scheme_to_identifier_map_;
118125
std::map<std::string, std::string> options_;
@@ -163,6 +170,11 @@ class PAIMON_EXPORT WriteContextBuilder {
163170
/// @return Reference to this builder for method chaining.
164171
WriteContextBuilder& WithExecutor(const std::shared_ptr<Executor>& executor);
165172

173+
/// Set custom IO manager for lookup and external disk spill operations.
174+
/// @param io_manager The IO manager to use.
175+
/// @return Reference to this builder for method chaining.
176+
WriteContextBuilder& WithIOManager(const std::shared_ptr<IOManager>& io_manager);
177+
166178
/// For postpone bucket mode in pk table, `WithWriteId()` supposed to be used.
167179
///
168180
/// Each worker must have its own unique `write_id` within a task, which is used as the prefix

src/paimon/CMakeLists.txt

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -229,7 +229,9 @@ set(PAIMON_CORE_SRCS
229229
core/mergetree/compact/aggregate/field_sum_agg.cpp
230230
core/mergetree/compact/interval_partition.cpp
231231
core/mergetree/compact/loser_tree.cpp
232+
core/mergetree/compact/merge_tree_compact_manager.cpp
232233
core/mergetree/compact/merge_tree_compact_rewriter.cpp
234+
core/mergetree/compact/merge_tree_compact_task.cpp
233235
core/mergetree/compact/partial_update_merge_function.cpp
234236
core/mergetree/compact/sort_merge_reader_with_loser_tree.cpp
235237
core/mergetree/compact/sort_merge_reader_with_min_heap.cpp
@@ -510,6 +512,7 @@ if(PAIMON_BUILD_TESTS)
510512
core/catalog/identifier_test.cpp
511513
core/compact/compact_deletion_file_test.cpp
512514
core/core_options_test.cpp
515+
core/options/lookup_strategy_test.cpp
513516
core/deletionvectors/apply_deletion_vector_batch_reader_test.cpp
514517
core/deletionvectors/bitmap_deletion_vector_test.cpp
515518
core/deletionvectors/bucketed_dv_maintainer_test.cpp
@@ -577,6 +580,7 @@ if(PAIMON_BUILD_TESTS)
577580
core/mergetree/compact/universal_compaction_test.cpp
578581
core/mergetree/compact/force_up_level0_compaction_test.cpp
579582
core/mergetree/compact/compact_strategy_test.cpp
583+
core/mergetree/compact/merge_tree_compact_manager_test.cpp
580584
core/mergetree/compact/merge_tree_compact_rewriter_test.cpp
581585
core/mergetree/lookup/persist_processor_test.cpp
582586
core/mergetree/drop_delete_reader_test.cpp

src/paimon/common/defs.cpp

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,4 +102,14 @@ const char Options::LOOKUP_CACHE_SPILL_COMPRESSION[] = "lookup.cache-spill-compr
102102
const char Options::SPILL_COMPRESSION_ZSTD_LEVEL[] = "spill-compression.zstd-level";
103103
const char Options::CACHE_PAGE_SIZE[] = "cache-page-size";
104104
const char Options::FILE_FORMAT_PER_LEVEL[] = "file.format.per.level";
105+
const char Options::COMPACTION_MAX_SIZE_AMPLIFICATION_PERCENT[] =
106+
"compaction.max-size-amplification-percent";
107+
const char Options::COMPACTION_SIZE_RATIO[] = "compaction.size-ratio";
108+
const char Options::NUM_SORTED_RUNS_COMPACTION_TRIGGER[] = "num-sorted-run.compaction-trigger";
109+
const char Options::NUM_SORTED_RUNS_STOP_TRIGGER[] = "num-sorted-run.stop-trigger";
110+
const char Options::NUM_LEVELS[] = "num-levels";
111+
const char Options::COMPACTION_FORCE_UP_LEVEL_0[] = "compaction.force-up-level-0";
112+
const char Options::LOOKUP_COMPACT[] = "lookup-compact";
113+
const char Options::LOOKUP_COMPACT_MAX_INTERVAL[] = "lookup-compact.max-interval";
114+
105115
} // namespace paimon

src/paimon/core/append/append_only_writer.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ Status AppendOnlyWriter::Flush(bool wait_for_latest_compaction, bool forced_full
140140
}
141141
// add new generated files
142142
for (const auto& flushed_file : flushed_files) {
143-
compact_manager_->AddNewFile(flushed_file);
143+
PAIMON_RETURN_NOT_OK(compact_manager_->AddNewFile(flushed_file));
144144
}
145145
PAIMON_RETURN_NOT_OK(TrySyncLatestCompaction(wait_for_latest_compaction));
146146
PAIMON_RETURN_NOT_OK(compact_manager_->TriggerCompaction(forced_full_compaction));

src/paimon/core/append/bucketed_append_compact_manager.cpp

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@
1616

1717
#include "paimon/core/append/bucketed_append_compact_manager.h"
1818

19+
#include <cassert>
20+
1921
#include "paimon/common/executor/future.h"
2022

2123
namespace paimon {
@@ -26,7 +28,7 @@ BucketedAppendCompactManager::BucketedAppendCompactManager(
2628
const std::shared_ptr<BucketedDvMaintainer>& dv_maintainer, int32_t min_file_num,
2729
int64_t target_file_size, int64_t compaction_file_size, bool force_rewrite_all_files,
2830
CompactRewriter rewriter, const std::shared_ptr<CompactionMetrics::Reporter>& reporter,
29-
const std::shared_ptr<std::atomic_bool>& cancel_flag)
31+
const std::shared_ptr<CancellationController>& cancellation_controller)
3032
: executor_(executor),
3133
dv_maintainer_(dv_maintainer),
3234
min_file_num_(min_file_num),
@@ -39,15 +41,16 @@ BucketedAppendCompactManager::BucketedAppendCompactManager(
3941
[](const std::shared_ptr<DataFileMeta>& lhs, const std::shared_ptr<DataFileMeta>& rhs) {
4042
return lhs->min_sequence_number > rhs->min_sequence_number;
4143
}),
42-
cancel_flag_(cancel_flag ? cancel_flag : std::make_shared<std::atomic_bool>(false)),
44+
cancellation_controller_(cancellation_controller),
4345
logger_(Logger::GetLogger("BucketedAppendCompactManager")) {
46+
assert(cancellation_controller_ != nullptr);
4447
for (const auto& file : restored) {
4548
to_compact_.push(file);
4649
}
4750
}
4851

4952
void BucketedAppendCompactManager::CancelCompaction() {
50-
cancel_flag_->store(true, std::memory_order_relaxed);
53+
cancellation_controller_->Cancel();
5154
CompactFutureManager::CancelCompaction();
5255
}
5356

@@ -78,7 +81,7 @@ Status BucketedAppendCompactManager::TriggerFullCompaction() {
7881
compacting.push_back(to_compact_.top());
7982
to_compact_.pop();
8083
}
81-
cancel_flag_->store(false, std::memory_order_relaxed);
84+
cancellation_controller_->Reset();
8285
auto compact_task = std::make_shared<FullCompactTask>(reporter_, dv_maintainer_, compacting,
8386
compaction_file_size_,
8487
force_rewrite_all_files_, rewriter_);
@@ -95,7 +98,7 @@ void BucketedAppendCompactManager::TriggerCompactionWithBestEffort() {
9598
}
9699
std::optional<std::vector<std::shared_ptr<DataFileMeta>>> picked = PickCompactBefore();
97100
if (picked) {
98-
cancel_flag_->store(false, std::memory_order_relaxed);
101+
cancellation_controller_->Reset();
99102
compacting_ = picked.value();
100103
auto compact_task = std::make_shared<AutoCompactTask>(reporter_, dv_maintainer_,
101104
compacting_.value(), rewriter_);

src/paimon/core/append/bucketed_append_compact_manager.h

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@
1616

1717
#pragma once
1818

19-
#include <atomic>
2019
#include <deque>
2120
#include <functional>
2221
#include <memory>
@@ -25,6 +24,7 @@
2524
#include <vector>
2625

2726
#include "paimon/common/executor/future.h"
27+
#include "paimon/core/compact/cancellation_controller.h"
2828
#include "paimon/core/compact/compact_deletion_file.h"
2929
#include "paimon/core/compact/compact_future_manager.h"
3030
#include "paimon/core/compact/compact_task.h"
@@ -67,14 +67,13 @@ class BucketedAppendCompactManager : public CompactFutureManager {
6767
};
6868
}
6969

70-
BucketedAppendCompactManager(const std::shared_ptr<Executor>& executor,
71-
const std::vector<std::shared_ptr<DataFileMeta>>& restored,
72-
const std::shared_ptr<BucketedDvMaintainer>& dv_maintainer,
73-
int32_t min_file_num, int64_t target_file_size,
74-
int64_t compaction_file_size, bool force_rewrite_all_files,
75-
CompactRewriter rewriter,
76-
const std::shared_ptr<CompactionMetrics::Reporter>& reporter,
77-
const std::shared_ptr<std::atomic_bool>& cancel_flag);
70+
BucketedAppendCompactManager(
71+
const std::shared_ptr<Executor>& executor,
72+
const std::vector<std::shared_ptr<DataFileMeta>>& restored,
73+
const std::shared_ptr<BucketedDvMaintainer>& dv_maintainer, int32_t min_file_num,
74+
int64_t target_file_size, int64_t compaction_file_size, bool force_rewrite_all_files,
75+
CompactRewriter rewriter, const std::shared_ptr<CompactionMetrics::Reporter>& reporter,
76+
const std::shared_ptr<CancellationController>& cancellation_controller);
7877
~BucketedAppendCompactManager() override = default;
7978

8079
void CancelCompaction() override;
@@ -88,8 +87,9 @@ class BucketedAppendCompactManager : public CompactFutureManager {
8887
return false;
8988
}
9089

91-
void AddNewFile(const std::shared_ptr<DataFileMeta>& file) override {
90+
Status AddNewFile(const std::shared_ptr<DataFileMeta>& file) override {
9291
to_compact_.push(file);
92+
return Status::OK();
9393
}
9494

9595
std::vector<std::shared_ptr<DataFileMeta>> AllFiles() const override;
@@ -199,7 +199,7 @@ class BucketedAppendCompactManager : public CompactFutureManager {
199199
std::shared_ptr<CompactionMetrics::Reporter> reporter_;
200200
std::optional<std::vector<std::shared_ptr<DataFileMeta>>> compacting_;
201201
DataFileMetaPriorityQueue to_compact_;
202-
std::shared_ptr<std::atomic_bool> cancel_flag_;
202+
std::shared_ptr<CancellationController> cancellation_controller_;
203203
std::unique_ptr<Logger> logger_;
204204
};
205205

src/paimon/core/append/bucketed_append_compact_manager_test.cpp

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@
1616

1717
#include "paimon/core/append/bucketed_append_compact_manager.h"
1818

19-
#include <atomic>
2019
#include <chrono>
2120
#include <future>
2221
#include <optional>
@@ -76,7 +75,7 @@ class BucketedAppendCompactManagerTest : public testing::Test {
7675
executor_, to_compact_before_pick,
7776
/*dv_maintainer=*/nullptr, min_file_num, target_file_size, threshold,
7877
/*force_rewrite_all_files=*/false, /*rewriter=*/nullptr, /*reporter=*/nullptr,
79-
/*cancel_flag=*/std::make_shared<std::atomic_bool>(false));
78+
/*cancellation_controller=*/std::make_shared<CancellationController>());
8079
auto actual = manager.PickCompactBefore();
8180
if (expected_present) {
8281
ASSERT_TRUE(actual.has_value());
@@ -267,14 +266,14 @@ TEST_F(BucketedAppendCompactManagerTest, TestPick) {
267266
}
268267

269268
TEST_F(BucketedAppendCompactManagerTest, TestCancelCompactionPropagatesToRewriteLoop) {
270-
auto cancel_flag = std::make_shared<std::atomic_bool>(false);
269+
auto cancellation_controller = std::make_shared<CancellationController>();
271270
auto exit_signal = std::make_shared<std::promise<void>>();
272271
auto exit_future = exit_signal->get_future();
273272

274-
auto rewriter = [cancel_flag,
273+
auto rewriter = [cancellation_controller,
275274
exit_signal](const std::vector<std::shared_ptr<DataFileMeta>>& to_compact)
276275
-> Result<std::vector<std::shared_ptr<DataFileMeta>>> {
277-
while (!cancel_flag->load(std::memory_order_relaxed)) {
276+
while (!cancellation_controller->IsCancelled()) {
278277
std::this_thread::sleep_for(std::chrono::milliseconds(1));
279278
}
280279
exit_signal->set_value();
@@ -288,7 +287,7 @@ TEST_F(BucketedAppendCompactManagerTest, TestCancelCompactionPropagatesToRewrite
288287
/*target_file_size=*/1024,
289288
/*compaction_file_size=*/700,
290289
/*force_rewrite_all_files=*/false, rewriter,
291-
/*reporter=*/nullptr, cancel_flag);
290+
/*reporter=*/nullptr, cancellation_controller);
292291

293292
ASSERT_OK(manager.TriggerCompaction(/*full_compaction=*/true));
294293
manager.CancelCompaction();
@@ -297,7 +296,8 @@ TEST_F(BucketedAppendCompactManagerTest, TestCancelCompactionPropagatesToRewrite
297296
}
298297

299298
TEST_F(BucketedAppendCompactManagerTest, TestTriggerCompactionResetsCancelFlag) {
300-
auto cancel_flag = std::make_shared<std::atomic_bool>(true);
299+
auto cancellation_controller = std::make_shared<CancellationController>();
300+
cancellation_controller->Cancel();
301301
auto rewriter = [](const std::vector<std::shared_ptr<DataFileMeta>>& to_compact)
302302
-> Result<std::vector<std::shared_ptr<DataFileMeta>>> { return to_compact; };
303303

@@ -308,10 +308,10 @@ TEST_F(BucketedAppendCompactManagerTest, TestTriggerCompactionResetsCancelFlag)
308308
/*target_file_size=*/1024,
309309
/*compaction_file_size=*/700,
310310
/*force_rewrite_all_files=*/false, rewriter,
311-
/*reporter=*/nullptr, cancel_flag);
311+
/*reporter=*/nullptr, cancellation_controller);
312312

313313
ASSERT_OK(manager.TriggerCompaction(/*full_compaction=*/true));
314-
EXPECT_FALSE(cancel_flag->load(std::memory_order_relaxed));
314+
EXPECT_FALSE(cancellation_controller->IsCancelled());
315315
}
316316

317317
} // namespace paimon::test
Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,41 @@
1+
/*
2+
* Copyright 2026-present Alibaba Inc.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* http://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
17+
#pragma once
18+
19+
#include <atomic>
20+
21+
namespace paimon {
22+
23+
class CancellationController {
24+
public:
25+
void Cancel() {
26+
cancelled_.store(true, std::memory_order_relaxed);
27+
}
28+
29+
void Reset() {
30+
cancelled_.store(false, std::memory_order_relaxed);
31+
}
32+
33+
bool IsCancelled() const {
34+
return cancelled_.load(std::memory_order_relaxed);
35+
}
36+
37+
private:
38+
std::atomic_bool cancelled_{false};
39+
};
40+
41+
} // namespace paimon

0 commit comments

Comments
 (0)