Skip to content

Commit e5a7da0

Browse files
authored
feat(mergetree): add compact manager, task abstractions and file rewrite task (#105)
1 parent 5ba61ad commit e5a7da0

10 files changed

Lines changed: 2122 additions & 0 deletions
Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
19+
#pragma once
20+
21+
#include <cstdint>
22+
#include <memory>
23+
#include <vector>
24+
25+
#include "paimon/core/compact/compact_task.h"
26+
#include "paimon/core/compact/compact_unit.h"
27+
#include "paimon/core/mergetree/compact/compact_rewriter.h"
28+
#include "paimon/core/mergetree/sorted_run.h"
29+
30+
namespace paimon {
31+
32+
/// Compact task for file rewrite compaction.
33+
class FileRewriteCompactTask : public CompactTask {
34+
public:
35+
FileRewriteCompactTask(const std::shared_ptr<CompactRewriter>& rewriter,
36+
const CompactUnit& unit, bool drop_delete,
37+
const std::shared_ptr<CompactionMetrics::Reporter>& metrics_reporter)
38+
: CompactTask(metrics_reporter),
39+
rewriter_(rewriter),
40+
output_level_(unit.output_level),
41+
files_(unit.files),
42+
drop_delete_(drop_delete) {}
43+
44+
protected:
45+
Result<std::shared_ptr<CompactResult>> DoCompact() override {
46+
auto result = std::make_shared<CompactResult>();
47+
for (const auto& file : files_) {
48+
PAIMON_RETURN_NOT_OK(RewriteFile(file, result.get()));
49+
}
50+
return result;
51+
}
52+
53+
private:
54+
Status RewriteFile(const std::shared_ptr<DataFileMeta>& file, CompactResult* to_update) {
55+
std::vector<std::vector<SortedRun>> candidate = {{SortedRun::FromSingle(file)}};
56+
PAIMON_ASSIGN_OR_RAISE(CompactResult rewritten,
57+
rewriter_->Rewrite(output_level_, drop_delete_, candidate));
58+
return to_update->Merge(rewritten);
59+
}
60+
61+
std::shared_ptr<CompactRewriter> rewriter_;
62+
int32_t output_level_;
63+
std::vector<std::shared_ptr<DataFileMeta>> files_;
64+
bool drop_delete_;
65+
};
66+
67+
} // namespace paimon
Lines changed: 167 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
19+
#include "paimon/core/mergetree/compact/file_rewrite_compact_task.h"
20+
21+
#include <memory>
22+
#include <string>
23+
#include <utility>
24+
#include <vector>
25+
26+
#include "gtest/gtest.h"
27+
#include "paimon/core/io/data_file_meta.h"
28+
#include "paimon/core/manifest/file_source.h"
29+
#include "paimon/core/stats/simple_stats.h"
30+
#include "paimon/testing/utils/testharness.h"
31+
32+
namespace paimon::test {
33+
namespace {
34+
35+
class SpyCompactRewriter final : public CompactRewriter {
36+
public:
37+
Result<CompactResult> Rewrite(int32_t output_level, bool drop_delete,
38+
const std::vector<std::vector<SortedRun>>& sections) override {
39+
output_levels.push_back(output_level);
40+
drop_deletes.push_back(drop_delete);
41+
42+
EXPECT_EQ(sections.size(), 1);
43+
EXPECT_EQ(sections[0].size(), 1);
44+
EXPECT_EQ(sections[0][0].Files().size(), 1);
45+
46+
auto input = sections[0][0].Files()[0];
47+
rewritten_files.push_back(input);
48+
49+
PAIMON_ASSIGN_OR_RAISE(auto upgraded, input->Upgrade(output_level));
50+
return CompactResult({input}, {upgraded});
51+
}
52+
53+
Result<CompactResult> Upgrade(int32_t, const std::shared_ptr<DataFileMeta>&) override {
54+
return Status::Invalid("Upgrade should not be called by FileRewriteCompactTask");
55+
}
56+
57+
Status Close() override {
58+
return Status::OK();
59+
}
60+
61+
std::vector<int32_t> output_levels;
62+
std::vector<bool> drop_deletes;
63+
std::vector<std::shared_ptr<DataFileMeta>> rewritten_files;
64+
};
65+
66+
class FailingCompactRewriter final : public CompactRewriter {
67+
public:
68+
explicit FailingCompactRewriter(size_t fail_on_call) : fail_on_call_(fail_on_call) {}
69+
70+
Result<CompactResult> Rewrite(int32_t, bool,
71+
const std::vector<std::vector<SortedRun>>& sections) override {
72+
EXPECT_EQ(sections.size(), 1);
73+
EXPECT_EQ(sections[0].size(), 1);
74+
EXPECT_EQ(sections[0][0].Files().size(), 1);
75+
76+
rewritten_files.push_back(sections[0][0].Files()[0]);
77+
++call_count_;
78+
if (call_count_ == fail_on_call_) {
79+
return Status::Invalid("injected rewrite failure");
80+
}
81+
82+
auto input = sections[0][0].Files()[0];
83+
PAIMON_ASSIGN_OR_RAISE(auto upgraded, input->Upgrade(/*new_level=*/1));
84+
return CompactResult({input}, {upgraded});
85+
}
86+
87+
Result<CompactResult> Upgrade(int32_t, const std::shared_ptr<DataFileMeta>&) override {
88+
return Status::Invalid("Upgrade should not be called by FileRewriteCompactTask");
89+
}
90+
91+
Status Close() override {
92+
return Status::OK();
93+
}
94+
95+
size_t call_count() const {
96+
return call_count_;
97+
}
98+
99+
std::vector<std::shared_ptr<DataFileMeta>> rewritten_files;
100+
101+
private:
102+
size_t fail_on_call_;
103+
size_t call_count_ = 0;
104+
};
105+
106+
std::shared_ptr<DataFileMeta> NewAppendFile(const std::string& file_name, int64_t row_count,
107+
int64_t min_sequence_number,
108+
int64_t max_sequence_number) {
109+
return DataFileMeta::ForAppend(file_name, /*file_size=*/row_count, row_count,
110+
SimpleStats::EmptyStats(), min_sequence_number,
111+
max_sequence_number, /*schema_id=*/0, FileSource::Append(),
112+
std::nullopt, std::nullopt, std::nullopt, std::nullopt)
113+
.value();
114+
}
115+
116+
} // namespace
117+
118+
TEST(FileRewriteCompactTaskTest, TestRewriteEachFileAndMergeResults) {
119+
auto file1 = NewAppendFile("file-1", /*row_count=*/10, /*min_sequence_number=*/0,
120+
/*max_sequence_number=*/9);
121+
auto file2 = NewAppendFile("file-2", /*row_count=*/20, /*min_sequence_number=*/10,
122+
/*max_sequence_number=*/29);
123+
CompactUnit unit(/*output_level=*/2, {file1, file2}, /*file_rewrite=*/true);
124+
125+
auto rewriter = std::make_shared<SpyCompactRewriter>();
126+
FileRewriteCompactTask task(rewriter, unit, /*drop_delete=*/true, /*metrics_reporter=*/nullptr);
127+
128+
ASSERT_OK_AND_ASSIGN(std::shared_ptr<CompactResult> result, task.Execute());
129+
130+
ASSERT_EQ(rewriter->output_levels, std::vector<int32_t>({2, 2}));
131+
ASSERT_EQ(rewriter->drop_deletes, std::vector<bool>({true, true}));
132+
ASSERT_EQ(rewriter->rewritten_files.size(), 2);
133+
EXPECT_EQ(rewriter->rewritten_files[0], file1);
134+
EXPECT_EQ(rewriter->rewritten_files[1], file2);
135+
136+
ASSERT_EQ(result->Before().size(), 2);
137+
EXPECT_EQ(result->Before()[0], file1);
138+
EXPECT_EQ(result->Before()[1], file2);
139+
140+
ASSERT_EQ(result->After().size(), 2);
141+
EXPECT_EQ(result->After()[0]->file_name, file1->file_name);
142+
EXPECT_EQ(result->After()[1]->file_name, file2->file_name);
143+
EXPECT_EQ(result->After()[0]->level, 2);
144+
EXPECT_EQ(result->After()[1]->level, 2);
145+
}
146+
147+
TEST(FileRewriteCompactTaskTest, TestStopOnRewriteFailure) {
148+
auto file1 = NewAppendFile("file-1", /*row_count=*/10, /*min_sequence_number=*/0,
149+
/*max_sequence_number=*/9);
150+
auto file2 = NewAppendFile("file-2", /*row_count=*/20, /*min_sequence_number=*/10,
151+
/*max_sequence_number=*/29);
152+
auto file3 = NewAppendFile("file-3", /*row_count=*/30, /*min_sequence_number=*/30,
153+
/*max_sequence_number=*/59);
154+
CompactUnit unit(/*output_level=*/2, {file1, file2, file3}, /*file_rewrite=*/true);
155+
156+
auto rewriter = std::make_shared<FailingCompactRewriter>(/*fail_on_call=*/2);
157+
FileRewriteCompactTask task(rewriter, unit, /*drop_delete=*/false,
158+
/*metrics_reporter=*/nullptr);
159+
160+
ASSERT_NOK_WITH_MSG(task.Execute(), "injected rewrite failure");
161+
ASSERT_EQ(rewriter->call_count(), 2);
162+
ASSERT_EQ(rewriter->rewritten_files.size(), 2);
163+
EXPECT_EQ(rewriter->rewritten_files[0], file1);
164+
EXPECT_EQ(rewriter->rewritten_files[1], file2);
165+
}
166+
167+
} // namespace paimon::test

0 commit comments

Comments
 (0)