Skip to content

Commit 8924152

Browse files
authored
feat(compaction): support multi-level lookup in LSM tree (alibaba#179)
1 parent d893e0b commit 8924152

72 files changed

Lines changed: 2397 additions & 183 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/data/timestamp.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -100,6 +100,10 @@ class PAIMON_EXPORT Timestamp {
100100
nano_of_millisecond_ == other.nano_of_millisecond_;
101101
}
102102

103+
bool operator!=(const Timestamp& other) const {
104+
return !(*this == other);
105+
}
106+
103107
bool operator<(const Timestamp& other) const {
104108
if (millisecond_ == other.millisecond_) {
105109
return nano_of_millisecond_ < other.nano_of_millisecond_;

include/paimon/disk/io_manager.h

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
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 <cstdint>
20+
#include <memory>
21+
#include <string>
22+
23+
#include "paimon/result.h"
24+
#include "paimon/visibility.h"
25+
26+
namespace paimon {
27+
/// The facade for the provided disk I/O services.
28+
class PAIMON_EXPORT IOManager {
29+
public:
30+
virtual ~IOManager() = default;
31+
static std::unique_ptr<IOManager> Create(const std::string& tmp_dir);
32+
33+
/// @return Temp directory path.
34+
virtual const std::string& GetTempDir() const = 0;
35+
36+
virtual Result<std::string> GenerateTempFilePath(const std::string& prefix) const = 0;
37+
};
38+
} // namespace paimon

include/paimon/reader/file_batch_reader.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@ class PAIMON_EXPORT FileBatchReader : public BatchReader {
4747
using BatchReader::NextBatchWithBitmap;
4848

4949
/// Get the row number of the first row in the previously read batch.
50-
virtual uint64_t GetPreviousBatchFirstRowNumber() const = 0;
50+
virtual Result<uint64_t> GetPreviousBatchFirstRowNumber() const = 0;
5151

5252
/// Get the number of rows in the file.
5353
virtual Result<uint64_t> GetNumberOfRows() const = 0;

src/paimon/CMakeLists.txt

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,7 @@ set(PAIMON_COMMON_SRCS
144144
common/utils/string_utils.cpp)
145145

146146
set(PAIMON_CORE_SRCS
147+
core/disk/io_manager.cpp
147148
core/append/append_only_writer.cpp
148149
core/append/bucketed_append_compact_manager.cpp
149150
core/casting/binary_to_string_cast_executor.cpp
@@ -233,6 +234,8 @@ set(PAIMON_CORE_SRCS
233234
core/mergetree/compact/sort_merge_reader_with_loser_tree.cpp
234235
core/mergetree/compact/sort_merge_reader_with_min_heap.cpp
235236
core/mergetree/merge_tree_writer.cpp
237+
core/mergetree/levels.cpp
238+
core/mergetree/lookup_levels.cpp
236239
core/migrate/file_meta_utils.cpp
237240
core/operation/data_evolution_file_store_scan.cpp
238241
core/operation/data_evolution_split_read.cpp
@@ -537,6 +540,9 @@ if(PAIMON_BUILD_TESTS)
537540
core/manifest/partition_entry_test.cpp
538541
core/manifest/file_entry_test.cpp
539542
core/manifest/index_manifest_entry_serializer_test.cpp
543+
core/mergetree/levels_test.cpp
544+
core/mergetree/lookup_file_test.cpp
545+
core/mergetree/lookup_levels_test.cpp
540546
core/mergetree/compact/aggregate/aggregate_merge_function_test.cpp
541547
core/mergetree/compact/aggregate/field_aggregator_factory_test.cpp
542548
core/mergetree/compact/aggregate/field_bool_agg_test.cpp

src/paimon/common/data/timestamp_test.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@ TEST_F(TimestampTest, EqualityOperator) {
6868
Timestamp ts3(1622547800000, 654321);
6969
ASSERT_EQ(ts1, ts1);
7070
ASSERT_EQ(ts1, ts2);
71-
ASSERT_FALSE(ts1 == ts3);
71+
ASSERT_NE(ts1, ts3);
7272
}
7373

7474
TEST_F(TimestampTest, LessThanOperator) {

src/paimon/common/file_index/bitmap/apply_bitmap_index_batch_reader.h

Lines changed: 23 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,6 @@
2525
#include "arrow/c/helpers.h"
2626
#include "paimon/common/reader/reader_utils.h"
2727
#include "paimon/file_index/bitmap_index_result.h"
28-
#include "paimon/reader/batch_reader.h"
2928
#include "paimon/reader/file_batch_reader.h"
3029
#include "paimon/result.h"
3130
#include "paimon/status.h"
@@ -34,7 +33,7 @@
3433
namespace paimon {
3534
class Metrics;
3635

37-
class ApplyBitmapIndexBatchReader : public BatchReader {
36+
class ApplyBitmapIndexBatchReader : public FileBatchReader {
3837
public:
3938
ApplyBitmapIndexBatchReader(std::unique_ptr<FileBatchReader>&& reader, RoaringBitmap32&& bitmap)
4039
: reader_(std::move(reader)), bitmap_(std::move(bitmap)) {
@@ -72,10 +71,31 @@ class ApplyBitmapIndexBatchReader : public BatchReader {
7271
return reader_->GetReaderMetrics();
7372
}
7473

74+
Result<std::unique_ptr<::ArrowSchema>> GetFileSchema() const override {
75+
return reader_->GetFileSchema();
76+
}
77+
78+
Status SetReadSchema(::ArrowSchema* read_schema, const std::shared_ptr<Predicate>& predicate,
79+
const std::optional<RoaringBitmap32>& selection_bitmap) override {
80+
return Status::Invalid("ApplyBitmapIndexBatchReader does not support SetReadSchema");
81+
}
82+
83+
Result<uint64_t> GetPreviousBatchFirstRowNumber() const override {
84+
return reader_->GetPreviousBatchFirstRowNumber();
85+
}
86+
87+
Result<uint64_t> GetNumberOfRows() const override {
88+
return reader_->GetNumberOfRows();
89+
}
90+
91+
bool SupportPreciseBitmapSelection() const override {
92+
return reader_->SupportPreciseBitmapSelection();
93+
}
94+
7595
private:
7696
Result<RoaringBitmap32> Filter(int32_t batch_size) const {
7797
RoaringBitmap32 is_valid;
78-
int32_t start_pos = reader_->GetPreviousBatchFirstRowNumber();
98+
PAIMON_ASSIGN_OR_RAISE(int32_t start_pos, reader_->GetPreviousBatchFirstRowNumber());
7999
int32_t length = batch_size;
80100
for (auto iter = bitmap_.EqualOrLarger(start_pos);
81101
iter != bitmap_.End() && *iter < start_pos + length; ++iter) {

src/paimon/common/reader/delegating_prefetch_reader.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ class DelegatingPrefetchReader : public FileBatchReader {
5454
return prefetch_reader_->SetReadSchema(read_schema, predicate, selection_bitmap);
5555
}
5656

57-
uint64_t GetPreviousBatchFirstRowNumber() const override {
57+
Result<uint64_t> GetPreviousBatchFirstRowNumber() const override {
5858
return GetReader()->GetPreviousBatchFirstRowNumber();
5959
}
6060

src/paimon/common/reader/prefetch_file_batch_reader_impl.cpp

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -409,7 +409,8 @@ Status PrefetchFileBatchReaderImpl::EnsureReaderPosition(
409409
Status PrefetchFileBatchReaderImpl::HandleReadResult(
410410
size_t reader_idx, const std::pair<uint64_t, uint64_t>& read_range,
411411
ReadBatchWithBitmap&& read_batch_with_bitmap) {
412-
uint64_t first_row_number = readers_[reader_idx]->GetPreviousBatchFirstRowNumber();
412+
PAIMON_ASSIGN_OR_RAISE(uint64_t first_row_number,
413+
readers_[reader_idx]->GetPreviousBatchFirstRowNumber());
413414
auto& prefetch_queue = prefetch_queues_[reader_idx];
414415
if (!BatchReader::IsEofBatch(read_batch_with_bitmap)) {
415416
auto& [read_batch, bitmap] = read_batch_with_bitmap;
@@ -570,7 +571,7 @@ Result<std::unique_ptr<::ArrowSchema>> PrefetchFileBatchReaderImpl::GetFileSchem
570571
return readers_[0]->GetFileSchema();
571572
}
572573

573-
uint64_t PrefetchFileBatchReaderImpl::GetPreviousBatchFirstRowNumber() const {
574+
Result<uint64_t> PrefetchFileBatchReaderImpl::GetPreviousBatchFirstRowNumber() const {
574575
return previous_batch_first_row_num_;
575576
}
576577

src/paimon/common/reader/prefetch_file_batch_reader_impl.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,7 @@ class PrefetchFileBatchReaderImpl : public PrefetchFileBatchReader {
7676
const std::optional<RoaringBitmap32>& selection_bitmap) override;
7777

7878
Status SeekToRow(uint64_t row_number) override;
79-
uint64_t GetPreviousBatchFirstRowNumber() const override;
79+
Result<uint64_t> GetPreviousBatchFirstRowNumber() const override;
8080
Result<uint64_t> GetNumberOfRows() const override;
8181
uint64_t GetNextRowToRead() const override;
8282
void Close() override;

src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp

Lines changed: 17 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -206,11 +206,11 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestSimple) {
206206
/*enable_adaptive_prefetch_strategy=*/false, executor_,
207207
/*initialize_read_ranges=*/true, /*prefetch_cache_mode=*/PrefetchCacheMode::ALWAYS,
208208
CacheConfig(), GetDefaultPool()));
209-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
209+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
210210
ASSERT_OK_AND_ASSIGN(auto result_array,
211211
ReadResultCollector::CollectResult(
212212
reader.get(), /*max simulated data processing time*/ 100));
213-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 101);
213+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 101);
214214
auto expected_array = std::make_shared<arrow::ChunkedArray>(data_array);
215215
ASSERT_TRUE(result_array->Equals(expected_array));
216216
}
@@ -396,11 +396,11 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestReadWithLargeBatchSize) {
396396
prefetch_max_parallel_num * 2, /*enable_adaptive_prefetch_strategy=*/false, executor_,
397397
/*initialize_read_ranges=*/true, /*prefetch_cache_mode=*/PrefetchCacheMode::ALWAYS,
398398
CacheConfig(), GetDefaultPool()));
399-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
399+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
400400
ASSERT_OK_AND_ASSIGN(auto result_array,
401401
ReadResultCollector::CollectResult(
402402
reader.get(), /*max simulated data processing time*/ 100));
403-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 101);
403+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 101);
404404
auto expected_array = std::make_shared<arrow::ChunkedArray>(data_array);
405405
ASSERT_TRUE(result_array->Equals(expected_array));
406406
}
@@ -424,11 +424,11 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestPartialReaderSuccessRead) {
424424
}
425425

426426
arrow::ArrayVector result_array_vector;
427-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
427+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
428428
ASSERT_OK_AND_ASSIGN(auto batch_with_bitmap, reader->NextBatchWithBitmap());
429429
auto& [batch, bitmap] = batch_with_bitmap;
430430
ASSERT_EQ(batch.first->length, bitmap.Cardinality());
431-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 0);
431+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 0);
432432
ASSERT_OK_AND_ASSIGN(auto array, ReadResultCollector::GetArray(std::move(batch)));
433433
result_array_vector.push_back(array);
434434
ASSERT_OK(prefetch_reader->GetReadStatus());
@@ -469,9 +469,9 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestAllReaderFailedWithIOError) {
469469
->SetNextBatchStatus(Status::IOError("mock error"));
470470
}
471471

472-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
472+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
473473
auto batch_result = reader->NextBatchWithBitmap();
474-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
474+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
475475
ASSERT_FALSE(batch_result.ok());
476476
ASSERT_TRUE(batch_result.status().IsIOError());
477477
ASSERT_FALSE(prefetch_reader->is_shutdown_);
@@ -480,7 +480,7 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestAllReaderFailedWithIOError) {
480480

481481
// call NextBatch again, will still return error status
482482
auto batch_result2 = reader->NextBatchWithBitmap();
483-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
483+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
484484
ASSERT_FALSE(batch_result2.ok());
485485
ASSERT_TRUE(batch_result2.status().IsIOError());
486486
}
@@ -497,11 +497,11 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestPrefetchWithEmptyData) {
497497
prefetch_max_parallel_num * 2, /*enable_adaptive_prefetch_strategy=*/false, executor_,
498498
/*initialize_read_ranges=*/true, /*prefetch_cache_mode=*/PrefetchCacheMode::ALWAYS,
499499
CacheConfig(), GetDefaultPool()));
500-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
500+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
501501
ASSERT_OK_AND_ASSIGN(auto result_array,
502502
ReadResultCollector::CollectResult(
503503
reader.get(), /*max simulated data processing time*/ 100));
504-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 0);
504+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 0);
505505
ASSERT_FALSE(result_array);
506506
}
507507

@@ -517,11 +517,11 @@ TEST_F(PrefetchFileBatchReaderImplTest, TestCallNextBatchAfterReadingEof) {
517517
prefetch_max_parallel_num * 2, /*enable_adaptive_prefetch_strategy=*/false, executor_,
518518
/*initialize_read_ranges=*/true, /*prefetch_cache_mode=*/PrefetchCacheMode::ALWAYS,
519519
CacheConfig(), GetDefaultPool()));
520-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
520+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
521521
ASSERT_OK_AND_ASSIGN(auto result_array,
522522
ReadResultCollector::CollectResult(
523523
reader.get(), /*max simulated data processing time*/ 100));
524-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 10);
524+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 10);
525525
auto expected_array = std::make_shared<arrow::ChunkedArray>(data_array);
526526
ASSERT_TRUE(result_array->Equals(expected_array));
527527

@@ -623,11 +623,11 @@ TEST_P(PrefetchFileBatchReaderImplTest, TestPrefetchWithPredicatePushdownWithCom
623623
PreparePrefetchReader(file_format, schema.get(), predicate,
624624
/*selection_bitmap=*/std::nullopt,
625625
/*batch_size=*/10, /*prefetch_max_parallel_num=*/3, cache_mode);
626-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
626+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
627627
ASSERT_OK_AND_ASSIGN(auto result_array,
628628
ReadResultCollector::CollectResult(
629629
reader.get(), /*max simulated data processing time*/ 100));
630-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 90);
630+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 90);
631631

632632
arrow::ArrayVector expected_array_vector;
633633
expected_array_vector.push_back(data_array->Slice(0, 30));
@@ -659,11 +659,11 @@ TEST_P(PrefetchFileBatchReaderImplTest,
659659
/*selection_bitmap=*/std::nullopt,
660660
/*batch_size=*/10, /*prefetch_max_parallel_num=*/3, cache_mode);
661661
ASSERT_OK(reader->RefreshReadRanges());
662-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), -1);
662+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), -1);
663663
ASSERT_OK_AND_ASSIGN(auto result_array,
664664
ReadResultCollector::CollectResult(
665665
reader.get(), /*max simulated data processing time*/ 100));
666-
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber(), 90);
666+
ASSERT_EQ(reader->GetPreviousBatchFirstRowNumber().value(), 90);
667667

668668
arrow::ArrayVector expected_array_vector;
669669
expected_array_vector.push_back(data_array->Slice(0, 20));

0 commit comments

Comments
 (0)