Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 50 additions & 0 deletions be/src/exec/operator/olap_scan_operator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

#include <fmt/format.h>

#include <algorithm>
#include <memory>
#include <numeric>
#include <optional>
Expand Down Expand Up @@ -140,6 +141,13 @@ Status OlapScanLocalState::_init_profile() {
// 2. init timer and counters
_reader_init_timer = ADD_TIMER(_scanner_profile, "ReaderInitTime");
_scanner_init_timer = ADD_TIMER(_scanner_profile, "ScannerInitTime");
_rowset_tso_prune_timer = ADD_TIMER(custom_profile(), "RowsetTsoPruneTime");
_rowsets_pruned_by_tso_counter =
ADD_COUNTER(custom_profile(), "RowsetsPrunedByTso", TUnit::UNIT);
_segments_pruned_by_tso_counter =
ADD_COUNTER(custom_profile(), "SegmentsPrunedByTso", TUnit::UNIT);
_tablets_pruned_by_tso_counter =
ADD_COUNTER(custom_profile(), "TabletsPrunedByTso", TUnit::UNIT);
_process_conjunct_timer = ADD_TIMER(custom_profile(), "ProcessConjunctTime");
_read_compressed_counter = ADD_COUNTER(_segment_profile, "CompressedBytesRead", TUnit::BYTES);
_read_uncompressed_counter =
Expand Down Expand Up @@ -635,6 +643,32 @@ bool OlapScanLocalState::_is_binlog_merge_scan() const {
return scan_type == TBinlogScanType::MIN_DELTA || scan_type == TBinlogScanType::DETAIL;
}

void OlapScanLocalState::_prune_rowsets_by_tso(const TPaloScanRange& scan_range,
TabletReadSource& read_source) {
SCOPED_TIMER(_rowset_tso_prune_timer);
int64_t pruned_segments = 0;
const auto pruned_rowsets = std::erase_if(read_source.rs_splits, [&](const auto& split) {
const auto& rowset = split.rs_reader->rowset();
const auto tso = rowset->commit_tso();
// Old rowsets can lack commit TSO metadata, including compaction inputs with an
// unknown endpoint. Keep them for the existing segment/row-level predicates.
if (tso.start_tso() < 0 || tso.end_tso() < 0) {
return false;
}
DCHECK_LE(tso.start_tso(), tso.end_tso());
// Rowset metadata is inclusive [min, max]; the query is half-open [start, end).
const bool outside_window =
(scan_range.__isset.start_tso && tso.end_tso() < scan_range.start_tso) ||
(scan_range.__isset.end_tso && tso.start_tso() >= scan_range.end_tso);
if (outside_window) {
pruned_segments += rowset->num_segments();
}
return outside_window;
});
COUNTER_UPDATE(_rowsets_pruned_by_tso_counter, pruned_rowsets);
COUNTER_UPDATE(_segments_pruned_by_tso_counter, pruned_segments);
}

Status OlapScanLocalState::_init_scanners(std::list<ScannerSPtr>* scanners) {
if (_scan_ranges.empty()) {
_eos = true;
Expand Down Expand Up @@ -795,6 +829,22 @@ Status OlapScanLocalState::_init_scanners(std::list<ScannerSPtr>* scanners) {
int scanners_per_tablet = std::max(1, 64 / (int)_scan_ranges.size());
for (size_t scan_range_idx = 0; scan_range_idx < _scan_ranges.size(); scan_range_idx++) {
const auto& palo_scan_range = *_scan_ranges[scan_range_idx];
if (read_row_binlog &&
(palo_scan_range.__isset.start_tso || palo_scan_range.__isset.end_tso)) {
auto& read_source = _read_sources[scan_range_idx];
// The version-consistent read source and delete predicates have already been
// captured. Prune before cloning readers or opening any segment footers.
_prune_rowsets_by_tso(palo_scan_range, read_source);
if (std::all_of(read_source.rs_splits.begin(), read_source.rs_splits.end(),
[](const auto& split) {
return split.rs_reader->rowset()->num_rows() == 0;
})) {
// Empty bootstrap rowsets may have no TSO. Skip the tablet even if those
// remain, and do not let OlapScanner recapture an empty read source.
COUNTER_UPDATE(_tablets_pruned_by_tso_counter, 1);
continue;
}
}
int64_t version = 0;
std::from_chars(palo_scan_range.version.data(),
palo_scan_range.version.data() + palo_scan_range.version.size(), version);
Expand Down
6 changes: 6 additions & 0 deletions be/src/exec/operator/olap_scan_operator.h
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,8 @@ class OlapScanLocalState final : public ScanLocalState<OlapScanLocalState> {

Status _init_scanners(std::list<ScannerSPtr>* scanners) override;

void _prune_rowsets_by_tso(const TPaloScanRange& scan_range, TabletReadSource& read_source);

Status _build_key_ranges_and_filters();

bool _is_tablet_pruned_by_runtime_filter(int64_t partition_id, int32_t bucket_seq,
Expand Down Expand Up @@ -165,6 +167,10 @@ class OlapScanLocalState final : public ScanLocalState<OlapScanLocalState> {
RuntimeProfile::Counter* _key_range_counter = nullptr;
RuntimeProfile::Counter* _reader_init_timer = nullptr;
RuntimeProfile::Counter* _scanner_init_timer = nullptr;
RuntimeProfile::Counter* _rowset_tso_prune_timer = nullptr;
RuntimeProfile::Counter* _rowsets_pruned_by_tso_counter = nullptr;
RuntimeProfile::Counter* _segments_pruned_by_tso_counter = nullptr;
RuntimeProfile::Counter* _tablets_pruned_by_tso_counter = nullptr;
RuntimeProfile::Counter* _process_conjunct_timer = nullptr;

RuntimeProfile::Counter* _io_timer = nullptr;
Expand Down
168 changes: 168 additions & 0 deletions be/test/exec/operator/olap_scan_operator_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -20,11 +20,13 @@
#include <gtest/gtest.h>

#include <memory>
#include <optional>

#include "common/object_pool.h"
#include "core/data_type/data_type_number.h"
#include "gen_cpp/PlanNodes_types.h"
#include "gen_cpp/QueryCache_types.h"
#include "storage/rowset/beta_rowset.h"
#include "testutil/desc_tbl_builder.h"
#include "testutil/mock/mock_runtime_state.h"

Expand Down Expand Up @@ -93,4 +95,170 @@ TEST_F(OlapScanOperatorBinlogPushDownTest, AppendOnlyKeepsValuePredicatePushDown
EXPECT_TRUE(_local_state->can_push_down_column_predicate(_value_slot));
}

class OlapScanOperatorTsoPruningTest : public OlapScanOperatorBinlogPushDownTest {
protected:
void SetUp() override {
OlapScanOperatorBinlogPushDownTest::SetUp();
_parent->_olap_scan_node.__set_read_row_binlog(true);
_profile = std::make_unique<RuntimeProfile>("TsoPruningTest");
_local_state->_scanner_init_timer = ADD_TIMER(_profile, "ScannerInitTime");
_local_state->_rowset_tso_prune_timer = ADD_TIMER(_profile, "RowsetTsoPruneTime");
_local_state->_rowsets_pruned_by_tso_counter =
ADD_COUNTER(_profile, "RowsetsPrunedByTso", TUnit::UNIT);
_local_state->_segments_pruned_by_tso_counter =
ADD_COUNTER(_profile, "SegmentsPrunedByTso", TUnit::UNIT);
_local_state->_tablets_pruned_by_tso_counter =
ADD_COUNTER(_profile, "TabletsPrunedByTso", TUnit::UNIT);
_local_state->_scan_dependency = Dependency::create_shared(0, 0, "Scan", false);
}

RowsetReaderSharedPtr make_reader(std::optional<TsoRange> tso, int64_t segments = 1,
int64_t rows = 1) {
auto meta = std::make_shared<RowsetMeta>();
meta->set_version(tso && tso->start_tso() == tso->end_tso() ? Version(2, 2)
: Version(2, 4));
meta->set_num_segments(segments);
meta->set_num_rows(rows);
if (tso) {
meta->set_commit_tso(*tso);
}
auto rowset = std::make_shared<BetaRowset>(std::make_shared<TabletSchema>(), meta,
"/nonexistent/tso_pruning_test");
RowsetReaderSharedPtr reader;
EXPECT_TRUE(rowset->create_reader(&reader).ok());
return reader;
}

std::unique_ptr<RuntimeProfile> _profile;
};

TEST_F(OlapScanOperatorTsoPruningTest, HalfOpenBoundsPreserveOverlappingAndUnknownRowsets) {
TPaloScanRange range;
range.__set_start_tso(100);
range.__set_end_tso(200);

auto before = make_reader(TsoRange {99, 99}, 2);
auto lower = make_reader(TsoRange {100, 100});
auto inside = make_reader(TsoRange {199, 199}, 3);
auto upper = make_reader(TsoRange {200, 200}, 4);
auto compacted_before = make_reader(TsoRange {10, 99}, 2);
auto touches_lower = make_reader(TsoRange {90, 100}, 3);
auto compacted_inside = make_reader(TsoRange {120, 170}, 2);
auto spans_window = make_reader(TsoRange {0, 250}, 6);
auto compacted_after = make_reader(TsoRange {200, 250}, 5);
auto missing = make_reader(std::nullopt, 2);
auto unknown_lower = make_reader(TsoRange {-1, 170});
auto unknown_upper = make_reader(TsoRange {100, -1});

TabletReadSource source;
for (const auto& reader :
{before, lower, inside, upper, compacted_before, touches_lower, compacted_inside,
spans_window, compacted_after, missing, unknown_lower, unknown_upper}) {
source.rs_splits.emplace_back(reader);
}
// Pruning the data source must not discard separately captured delete predicates.
source.delete_predicates.push_back(before->rowset()->rowset_meta());
_local_state->_prune_rowsets_by_tso(range, source);

const std::vector<RowsetReaderSharedPtr> expected {
lower, inside, touches_lower, compacted_inside,
spans_window, missing, unknown_lower, unknown_upper};
ASSERT_EQ(source.rs_splits.size(), expected.size());
for (size_t i = 0; i < expected.size(); ++i) {
EXPECT_EQ(source.rs_splits[i].rs_reader, expected[i]);
}
ASSERT_EQ(source.delete_predicates.size(), 1);
EXPECT_EQ(source.delete_predicates.front(), before->rowset()->rowset_meta());
EXPECT_EQ(_local_state->_rowsets_pruned_by_tso_counter->value(), 4);
EXPECT_EQ(_local_state->_segments_pruned_by_tso_counter->value(), 13);
}

TEST_F(OlapScanOperatorTsoPruningTest, LowerBoundOnly) {
TPaloScanRange range;
range.__set_start_tso(100);
TabletReadSource source;
auto at_lower = make_reader(TsoRange {100, 100});
auto later = make_reader(TsoRange {200, 200});
source.rs_splits.emplace_back(make_reader(TsoRange {99, 99}));
source.rs_splits.emplace_back(at_lower);
source.rs_splits.emplace_back(later);

_local_state->_prune_rowsets_by_tso(range, source);

ASSERT_EQ(source.rs_splits.size(), 2);
EXPECT_EQ(source.rs_splits[0].rs_reader, at_lower);
EXPECT_EQ(source.rs_splits[1].rs_reader, later);
EXPECT_EQ(_local_state->_rowsets_pruned_by_tso_counter->value(), 1);
}

TEST_F(OlapScanOperatorTsoPruningTest, UpperBoundOnly) {
TPaloScanRange range;
range.__set_end_tso(100);
TabletReadSource source;
auto earlier = make_reader(TsoRange {99, 99});
source.rs_splits.emplace_back(earlier);
source.rs_splits.emplace_back(make_reader(TsoRange {100, 100}));
source.rs_splits.emplace_back(make_reader(TsoRange {200, 200}));

_local_state->_prune_rowsets_by_tso(range, source);

ASSERT_EQ(source.rs_splits.size(), 1);
EXPECT_EQ(source.rs_splits.front().rs_reader, earlier);
EXPECT_EQ(_local_state->_rowsets_pruned_by_tso_counter->value(), 2);
}

TEST_F(OlapScanOperatorTsoPruningTest, NoBoundsPreservesAllRowsets) {
TPaloScanRange range;
TabletReadSource source;
source.rs_splits.emplace_back(make_reader(TsoRange {100, 100}));
source.rs_splits.emplace_back(make_reader(std::nullopt));

_local_state->_prune_rowsets_by_tso(range, source);

EXPECT_EQ(source.rs_splits.size(), 2);
EXPECT_EQ(_local_state->_rowsets_pruned_by_tso_counter->value(), 0);
EXPECT_EQ(_local_state->_segments_pruned_by_tso_counter->value(), 0);
}

TEST_F(OlapScanOperatorTsoPruningTest, EmptyWindowSkipsScannersAndSignalsEos) {
auto& range = *_local_state->_scan_ranges.front();
range.__set_binlog_scan_type(TBinlogScanType::DETAIL);
range.__set_start_tso(100);
range.__set_end_tso(200);
_local_state->_read_sources.resize(1);
auto& source = _local_state->_read_sources.front();
source.rs_splits.emplace_back(make_reader(TsoRange {99, 99}, 2));
source.rs_splits.emplace_back(make_reader(TsoRange {200, 200}, 3));

// There are no tablet or segment files. Successful preparation proves that the
// fully pruned source is not recaptured or passed to a scanner for initialization.
EXPECT_TRUE(_local_state->_prepare_scanners().ok());

EXPECT_TRUE(_local_state->_scanners.empty());
EXPECT_TRUE(_local_state->_eos);
EXPECT_EQ(_local_state->_scan_dependency->is_blocked_by(nullptr), nullptr);
EXPECT_EQ(_local_state->_rowsets_pruned_by_tso_counter->value(), 2);
EXPECT_EQ(_local_state->_segments_pruned_by_tso_counter->value(), 5);
EXPECT_EQ(_local_state->_tablets_pruned_by_tso_counter->value(), 1);
}

TEST_F(OlapScanOperatorTsoPruningTest, EmptyBootstrapDoesNotStartMinDeltaScanner) {
auto& range = *_local_state->_scan_ranges.front();
range.__set_start_tso(100);
range.__set_end_tso(200);
_local_state->_read_sources.resize(1);
auto& source = _local_state->_read_sources.front();
source.rs_splits.emplace_back(make_reader(std::nullopt, 0, 0));
source.rs_splits.emplace_back(make_reader(TsoRange {99, 99}, 2));

EXPECT_TRUE(_local_state->_prepare_scanners().ok());

EXPECT_TRUE(_local_state->_scanners.empty());
EXPECT_TRUE(_local_state->_eos);
EXPECT_EQ(_local_state->_scan_dependency->is_blocked_by(nullptr), nullptr);
EXPECT_EQ(_local_state->_rowsets_pruned_by_tso_counter->value(), 1);
EXPECT_EQ(_local_state->_segments_pruned_by_tso_counter->value(), 2);
EXPECT_EQ(_local_state->_tablets_pruned_by_tso_counter->value(), 1);
}

} // namespace doris
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
-- This file is automatically generated. You should know what you did if you want to edit this
-- !detail --
1 10 2
1 11 2
1 11 3
1 12 3
2 20 1
4 40 0

-- !min_delta --
1 10 2
1 12 3
2 20 1
4 40 0

-- !start_only --
1 10 2
1 11 2
1 11 3
1 12 3
2 20 1
4 40 0
5 50 0

-- !end_only --
1 10 0
2 20 0
3 30 0

-- !empty_detail --

-- !empty_min_delta --

-- !snapshot --
1 12
3 30
4 40
5 50

-- !compacted_detail --
1 10 2
1 11 2
1 11 3
1 12 3
2 20 1
4 40 0

-- !compacted_min_delta --
1 10 2
1 12 3
2 20 1
4 40 0

-- !compacted_snapshot --
1 12
3 30
4 40
5 50

Loading
Loading