Skip to content

Commit e85d7d5

Browse files
zjw1111claude
andcommitted
refactor: rename LeveledMerger to SpillFileMerger and improve code clarity
- Rename class LeveledMerger → SpillFileMerger, files leveled_merger → spill_file_merger - Rename Compact/Compaction methods to Merge (RunMergeIfNeeded, RunFinalMergeIfNeeded, etc.) - Move kMinFanIn validation from CoreOptions to ExternalSortBuffer::Create - Improve variable names (l→level_idx, a/b→lhs/rhs, n→files_to_merge, f→file) - Convert 6 spill tests from TEST_P to TEST_F (no parameterization needed) - Fix 4 misused TEST_P in btree_global_index and file_system tests - Make test assertions exact instead of range-based Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 18f3fe8 commit e85d7d5

11 files changed

Lines changed: 157 additions & 164 deletions

src/paimon/CMakeLists.txt

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -253,7 +253,7 @@ set(PAIMON_CORE_SRCS
253253
core/mergetree/merge_tree_writer.cpp
254254
core/mergetree/in_memory_sort_buffer.cpp
255255
core/mergetree/external_sort_buffer.cpp
256-
core/mergetree/leveled_merger.cpp
256+
core/mergetree/spill_file_merger.cpp
257257
core/mergetree/write_buffer.cpp
258258
core/mergetree/levels.cpp
259259
core/mergetree/lookup_file.cpp
@@ -649,7 +649,7 @@ if(PAIMON_BUILD_TESTS)
649649
core/mergetree/merge_tree_writer_test.cpp
650650
core/mergetree/write_buffer_test.cpp
651651
core/mergetree/sort_buffer_test.cpp
652-
core/mergetree/leveled_merger_test.cpp
652+
core/mergetree/spill_file_merger_test.cpp
653653
core/mergetree/sorted_run_test.cpp
654654
core/mergetree/spill_channel_manager_test.cpp
655655
core/mergetree/spill_reader_writer_test.cpp

src/paimon/common/fs/file_system_test.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,7 @@ class FileSystemTest : public ::testing::Test, public ::testing::WithParamInterf
173173
std::string test_root_;
174174
};
175175

176-
TEST_P(FileSystemTest, TestNoneFileSystemFactory) {
176+
TEST(FileSystemStaticTest, TestNoneFileSystemFactory) {
177177
std::map<std::string, std::string> fs_options;
178178
Result<std::unique_ptr<FileSystem>> fs =
179179
FileSystemFactory::Get("not exist", "/tmp", fs_options);
@@ -516,7 +516,7 @@ TEST_P(FileSystemTest, TestReadAndWriteAndAtomicStoreFile) {
516516
ASSERT_EQ(read_content, new_content);
517517
}
518518

519-
TEST_P(FileSystemTest, TestIsObjectStore) {
519+
TEST(FileSystemStaticTest, TestIsObjectStore) {
520520
ASSERT_OK_AND_ASSIGN(bool is_object_store, FileSystem::IsObjectStore("file:///tmp/test.txt"));
521521
ASSERT_FALSE(is_object_store);
522522
ASSERT_OK_AND_ASSIGN(is_object_store, FileSystem::IsObjectStore("/tmp/test.txt"));

src/paimon/common/global_index/btree/btree_global_index_integration_test.cpp

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1491,7 +1491,7 @@ TEST_P(BTreeGlobalIndexIntegrationTest, WriteAndReadAllNull) {
14911491
}
14921492
}
14931493

1494-
TEST_P(BTreeGlobalIndexIntegrationTest, WriteAndReadLargeDataWithSmallBlocks) {
1494+
TEST_F(BTreeGlobalIndexIntegrationTest, WriteAndReadLargeDataWithSmallBlocks) {
14951495
// Use very small block size and cache size to force multiple block evictions
14961496
auto file_writer = std::make_shared<FakeGlobalIndexFileWriter>(fs_, base_path_);
14971497
auto field = arrow::field("int_field", arrow::int32());
@@ -1710,7 +1710,7 @@ TEST_P(BTreeGlobalIndexIntegrationTest, CreateReaderWithMultiFieldSchema) {
17101710
"supposed to have single field");
17111711
}
17121712

1713-
TEST_P(BTreeGlobalIndexIntegrationTest, CreateWriterWithMissingField) {
1713+
TEST_F(BTreeGlobalIndexIntegrationTest, CreateWriterWithMissingField) {
17141714
auto file_writer = std::make_shared<FakeGlobalIndexFileWriter>(fs_, base_path_);
17151715
auto type = arrow::struct_({arrow::field("existing_field", arrow::int32())});
17161716
auto struct_type = std::dynamic_pointer_cast<arrow::StructType>(type);

src/paimon/core/core_options.cpp

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -498,11 +498,6 @@ struct CoreOptions::Impl {
498498
// Parse local-sort.max-num-file-handles - spill file handle cap for local merge
499499
PAIMON_RETURN_NOT_OK(parser.Parse(Options::LOCAL_SORT_MAX_NUM_FILE_HANDLES,
500500
&local_sort_max_num_file_handles));
501-
if (local_sort_max_num_file_handles < kLocalSortFileHandlesMinimalLimit) {
502-
return Status::Invalid(fmt::format(
503-
"invalid '{}': {}, must be at least {}", Options::LOCAL_SORT_MAX_NUM_FILE_HANDLES,
504-
local_sort_max_num_file_handles, kLocalSortFileHandlesMinimalLimit));
505-
}
506501
// Parse spill-compression - compression codec for spill files
507502
PAIMON_RETURN_NOT_OK(
508503
parser.Parse(Options::SPILL_COMPRESSION, &spill_compress_options.compress));

src/paimon/core/core_options.h

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -181,8 +181,6 @@ class PAIMON_EXPORT CoreOptions {
181181

182182
const std::map<std::string, std::string>& ToMap() const;
183183

184-
static constexpr int32_t kLocalSortFileHandlesMinimalLimit = 2;
185-
186184
private:
187185
std::optional<std::string> GetDataFileExternalPaths() const;
188186
std::optional<std::string> GetGlobalIndexExternalPath() const;

src/paimon/core/mergetree/external_sort_buffer.cpp

Lines changed: 16 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,11 @@ Result<std::unique_ptr<ExternalSortBuffer>> ExternalSortBuffer::Create(
4848
const std::shared_ptr<FieldsComparator>& user_defined_seq_comparator,
4949
const CoreOptions& options, const std::shared_ptr<IOManager>& io_manager,
5050
const std::shared_ptr<MemoryPool>& pool) {
51+
if (options.GetLocalSortMaxNumFileHandles() < kSpillMinFanIn) {
52+
return Status::Invalid(fmt::format(
53+
"invalid '{}': {}, must be at least {}", Options::LOCAL_SORT_MAX_NUM_FILE_HANDLES,
54+
options.GetLocalSortMaxNumFileHandles(), kSpillMinFanIn));
55+
}
5156
arrow::FieldVector key_fields;
5257
key_fields.reserve(trimmed_primary_keys.size());
5358
for (const auto& primary_key : trimmed_primary_keys) {
@@ -84,7 +89,7 @@ ExternalSortBuffer::ExternalSortBuffer(
8489
max_fan_in_(options.GetLocalSortMaxNumFileHandles()),
8590
spill_channel_manager_(
8691
std::make_shared<SpillChannelManager>(options_.GetFileSystem(), max_fan_in_)),
87-
leveled_merger_(std::make_unique<LeveledMerger>(max_fan_in_)),
92+
spill_merger_(std::make_unique<SpillFileMerger>(max_fan_in_)),
8893
spill_channel_enumerator_(spill_channel_enumerator),
8994
actual_max_fan_in_(max_fan_in_),
9095
spill_batch_size_(options_.GetWriteBatchSize()) {}
@@ -102,7 +107,7 @@ void ExternalSortBuffer::DoClear() {
102107

103108
spill_channel_manager_->Reset();
104109
total_spill_disk_bytes_ = 0;
105-
leveled_merger_->Clear();
110+
spill_merger_->Clear();
106111
}
107112

108113
void ExternalSortBuffer::Clear() {
@@ -113,7 +118,7 @@ uint64_t ExternalSortBuffer::GetMemorySize() const {
113118
return in_memory_buffer_->GetMemorySize();
114119
}
115120

116-
void ExternalSortBuffer::EstimateSpillParameters() {
121+
void ExternalSortBuffer::UpdateSpillParameters() {
117122
int64_t estimated_row_size = in_memory_buffer_->GetEstimateMemoryUseForEachRow();
118123
if (estimated_row_size <= 0) {
119124
return;
@@ -128,22 +133,21 @@ void ExternalSortBuffer::EstimateSpillParameters() {
128133
spill_batch_size_ = std::clamp(spill_batch_size_, min_batch_size, max_batch_size);
129134

130135
actual_max_fan_in_ = merge_budget / (spill_batch_size_ * estimated_row_size);
131-
actual_max_fan_in_ =
132-
std::clamp(actual_max_fan_in_, CoreOptions::kLocalSortFileHandlesMinimalLimit, max_fan_in_);
136+
actual_max_fan_in_ = std::clamp(actual_max_fan_in_, kSpillMinFanIn, max_fan_in_);
133137

134138
// Re-derive spill_batch_size_ from the clamped actual_max_fan_in_ to stay within merge_budget.
135139
spill_batch_size_ = merge_budget / (actual_max_fan_in_ * estimated_row_size);
136140
spill_batch_size_ = std::clamp(spill_batch_size_, 1, max_batch_size);
137141

138-
leveled_merger_->SetMaxFanIn(actual_max_fan_in_);
142+
spill_merger_->SetMaxFanIn(actual_max_fan_in_);
139143
}
140144

141145
Result<bool> ExternalSortBuffer::FlushMemory() {
142146
if (!in_memory_buffer_->HasData()) {
143147
return true;
144148
}
145149

146-
EstimateSpillParameters();
150+
UpdateSpillParameters();
147151
PAIMON_ASSIGN_OR_RAISE(std::vector<std::unique_ptr<KeyValueRecordReader>> memory_buffer_readers,
148152
in_memory_buffer_->CreateReaders());
149153
PAIMON_RETURN_NOT_OK(SpillMemoryBuffer(std::move(memory_buffer_readers)));
@@ -168,9 +172,9 @@ Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> ExternalSortBuffer::C
168172

169173
int32_t max_spill_files = actual_max_fan_in_ - 1;
170174
PAIMON_RETURN_NOT_OK(
171-
leveled_merger_->RunFinalCleanupIfNeeded(max_spill_files, CreateMergeFn()));
175+
spill_merger_->RunFinalMergeIfNeeded(max_spill_files, CreateSpillFileMergeFn()));
172176
PAIMON_ASSIGN_OR_RAISE(std::vector<std::unique_ptr<KeyValueRecordReader>> readers,
173-
CreateSpillReaders(leveled_merger_->GetAllFiles()));
177+
CreateSpillReaders(spill_merger_->GetAllFiles()));
174178
readers.insert(readers.end(), std::make_move_iterator(memory_readers.begin()),
175179
std::make_move_iterator(memory_readers.end()));
176180
return readers;
@@ -242,11 +246,11 @@ Status ExternalSortBuffer::SpillMemoryBuffer(
242246
PAIMON_ASSIGN_OR_RAISE(FileChannelInfo file_info,
243247
SpillToDisk(std::move(readers), spill_batch_size_));
244248
total_spill_disk_bytes_ += file_info.file_size;
245-
leveled_merger_->AddFile(file_info);
246-
return leveled_merger_->RunCompactionIfNeeded(CreateMergeFn());
249+
spill_merger_->AddFile(file_info);
250+
return spill_merger_->RunMergeIfNeeded(CreateSpillFileMergeFn());
247251
}
248252

249-
LeveledMerger::MergeFn ExternalSortBuffer::CreateMergeFn() {
253+
SpillFileMerger::MergeFn ExternalSortBuffer::CreateSpillFileMergeFn() {
250254
return [this](const std::vector<FileChannelInfo>& files) -> Result<FileChannelInfo> {
251255
return MergeAndReplaceFiles(files);
252256
};

src/paimon/core/mergetree/external_sort_buffer.h

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,8 +25,8 @@
2525
#include "paimon/core/core_options.h"
2626
#include "paimon/core/disk/file_io_channel.h"
2727
#include "paimon/core/mergetree/in_memory_sort_buffer.h"
28-
#include "paimon/core/mergetree/leveled_merger.h"
2928
#include "paimon/core/mergetree/sort_buffer.h"
29+
#include "paimon/core/mergetree/spill_file_merger.h"
3030
#include "paimon/record_batch.h"
3131
#include "paimon/result.h"
3232
#include "paimon/status.h"
@@ -64,16 +64,17 @@ class ExternalSortBuffer : public SortBuffer {
6464
bool HasData() const override;
6565

6666
private:
67+
static constexpr int32_t kSpillMinFanIn = 2;
6768
static constexpr int32_t kSpillMinBatchSize = 256;
6869

6970
void DoClear();
70-
void EstimateSpillParameters();
71+
void UpdateSpillParameters();
7172
bool HasSpilledData() const;
7273
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> CreateSpillReaders(
7374
const std::vector<FileChannelInfo>& files) const;
7475
Result<FileChannelInfo> SpillToDisk(
7576
std::vector<std::unique_ptr<KeyValueRecordReader>>&& readers, int32_t write_batch_size);
76-
LeveledMerger::MergeFn CreateMergeFn();
77+
SpillFileMerger::MergeFn CreateSpillFileMergeFn();
7778
Result<FileChannelInfo> MergeAndReplaceFiles(const std::vector<FileChannelInfo>& files);
7879
Status SpillMemoryBuffer(std::vector<std::unique_ptr<KeyValueRecordReader>>&& readers);
7980

@@ -98,7 +99,7 @@ class ExternalSortBuffer : public SortBuffer {
9899
const int32_t max_fan_in_;
99100
const std::shared_ptr<SpillChannelManager> spill_channel_manager_;
100101

101-
std::unique_ptr<LeveledMerger> leveled_merger_;
102+
std::unique_ptr<SpillFileMerger> spill_merger_;
102103
std::shared_ptr<FileIOChannel::Enumerator> spill_channel_enumerator_;
103104
int64_t total_spill_disk_bytes_ = 0;
104105
int32_t actual_max_fan_in_;

src/paimon/core/mergetree/merge_tree_writer_test.cpp

Lines changed: 10 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1076,7 +1076,7 @@ TEST_P(MergeTreeWriterTest, TestCloseSkipsDeleteForUpgradedFilesInCompactAfter)
10761076
<< "Intermediate file should be deleted because it's not in compact_before_";
10771077
}
10781078

1079-
TEST_P(MergeTreeWriterTest, TestSpillWithSameKeyDeduplicate) {
1079+
TEST_F(MergeTreeWriterTest, TestSpillWithSameKeyDeduplicate) {
10801080
ASSERT_OK_AND_ASSIGN(CoreOptions options,
10811081
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
10821082
{Options::WRITE_BUFFER_SIZE, "1"},
@@ -1111,9 +1111,8 @@ TEST_P(MergeTreeWriterTest, TestSpillWithSameKeyDeduplicate) {
11111111

11121112
WriteBatch(batch1, /*row_kinds=*/{}, merge_writer.get());
11131113
WriteBatch(batch2, /*row_kinds=*/{}, merge_writer.get());
1114-
// Leveled compaction may merge 2 files into 1 when actual_max_fan_in_ is low.
1115-
ASSERT_GE(2u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
1116-
ASSERT_GE(TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"), 1u);
1114+
// actual_max_fan_in_=2 with 1-byte budget, so 2 spill files compact to 1.
1115+
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
11171116

11181117
std::shared_ptr<arrow::Array> batch3 =
11191118
arrow::ipc::internal::json::ArrayFromJSON(value_type_, R"([
@@ -1142,7 +1141,7 @@ TEST_P(MergeTreeWriterTest, TestSpillWithSameKeyDeduplicate) {
11421141
CheckFileContent(expected_data_file_path, expected_array);
11431142
}
11441143

1145-
TEST_P(MergeTreeWriterTest, TestIntermediateMergeSpillFileBound) {
1144+
TEST_F(MergeTreeWriterTest, TestIntermediateMergeSpillFileBound) {
11461145
ASSERT_OK_AND_ASSIGN(CoreOptions options,
11471146
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
11481147
{Options::WRITE_BUFFER_SIZE, "1"},
@@ -1186,9 +1185,8 @@ TEST_P(MergeTreeWriterTest, TestIntermediateMergeSpillFileBound) {
11861185
ASSERT_EQ(1u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
11871186

11881187
WriteBatch(batch3, /*row_kinds=*/{}, merge_writer.get());
1189-
// Leveled compaction keeps files across levels, so file count may be > 1.
1190-
ASSERT_LE(TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"), 2u);
1191-
ASSERT_GE(TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"), 1u);
1188+
// After 3rd spill: 1 file at level 0 (new) + 1 file at level 1 (from prior compact) = 2.
1189+
ASSERT_EQ(2u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
11921190

11931191
ASSERT_OK_AND_ASSIGN(CommitIncrement commit_increment,
11941192
merge_writer->PrepareCommit(/*wait_compaction=*/false));
@@ -1207,7 +1205,7 @@ TEST_P(MergeTreeWriterTest, TestIntermediateMergeSpillFileBound) {
12071205
CheckFileContent(expected_data_file_path, expected_array);
12081206
}
12091207

1210-
TEST_P(MergeTreeWriterTest, TestDiskQuotaExhaustedFallsBackToFlushWriteBuffer) {
1208+
TEST_F(MergeTreeWriterTest, TestDiskQuotaExhaustedFallsBackToFlushWriteBuffer) {
12111209
ASSERT_OK_AND_ASSIGN(CoreOptions options,
12121210
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
12131211
{Options::WRITE_BUFFER_SIZE, "1"},
@@ -1281,7 +1279,7 @@ TEST_P(MergeTreeWriterTest, TestDiskQuotaExhaustedFallsBackToFlushWriteBuffer) {
12811279
}
12821280
}
12831281

1284-
TEST_P(MergeTreeWriterTest, TestFlushMemoryQuotaExhaustedFallsBackToFlushWriteBuffer) {
1282+
TEST_F(MergeTreeWriterTest, TestFlushMemoryQuotaExhaustedFallsBackToFlushWriteBuffer) {
12851283
// WRITE_BUFFER_SIZE is large enough so WriteBatch does NOT auto-spill.
12861284
// SPILL_MAX_DISK_SIZE is tiny so the first FlushMemory() exhausts the quota,
12871285
// triggering the fallback path: FlushMemory() -> quota exhausted -> FlushWriteBuffer.
@@ -1328,7 +1326,7 @@ TEST_P(MergeTreeWriterTest, TestFlushMemoryQuotaExhaustedFallsBackToFlushWriteBu
13281326
ASSERT_OK(merge_writer->Close());
13291327
}
13301328

1331-
TEST_P(MergeTreeWriterTest, TestCloseDeletesSpillTempFiles) {
1329+
TEST_F(MergeTreeWriterTest, TestCloseDeletesSpillTempFiles) {
13321330
ASSERT_OK_AND_ASSIGN(CoreOptions options,
13331331
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"},
13341332
{Options::WRITE_BUFFER_SIZE, "1"},
@@ -1360,7 +1358,7 @@ TEST_P(MergeTreeWriterTest, TestCloseDeletesSpillTempFiles) {
13601358
ASSERT_EQ(0u, TestHelper::CountChannelFiles(file_system_, dir->Str() + "/tmp"));
13611359
}
13621360

1363-
TEST_P(MergeTreeWriterTest, TestMultiplePrepareCommitWithSpill) {
1361+
TEST_F(MergeTreeWriterTest, TestMultiplePrepareCommitWithSpill) {
13641362
ASSERT_OK_AND_ASSIGN(
13651363
CoreOptions options,
13661364
CoreOptions::FromMap({{Options::FILE_FORMAT, "orc"}, {Options::WRITE_ONLY, "true"}}));

0 commit comments

Comments
 (0)