Skip to content
Open
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
2 changes: 2 additions & 0 deletions be/src/exec/operator/olap_scan_operator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -217,6 +217,8 @@ Status OlapScanLocalState::_init_profile() {
_lazy_read_seek_timer = ADD_TIMER(_segment_profile, "LazyReadSeekTime");
_lazy_read_seek_counter = ADD_COUNTER(_segment_profile, "LazyReadSeekCount", TUnit::UNIT);

_lazy_read_pruned_timer = ADD_TIMER(_segment_profile, "LazyReadPrunedTime");

_output_col_timer = ADD_TIMER(_segment_profile, "OutputColumnTime");

_stats_filtered_counter = ADD_COUNTER(_segment_profile, "RowsStatsFiltered", TUnit::UNIT);
Expand Down
1 change: 1 addition & 0 deletions be/src/exec/operator/olap_scan_operator.h
Original file line number Diff line number Diff line change
Expand Up @@ -210,6 +210,7 @@ class OlapScanLocalState final : public ScanLocalState<OlapScanLocalState> {
RuntimeProfile::Counter* _lazy_read_timer = nullptr;
RuntimeProfile::Counter* _lazy_read_seek_timer = nullptr;
RuntimeProfile::Counter* _lazy_read_seek_counter = nullptr;
RuntimeProfile::Counter* _lazy_read_pruned_timer = nullptr;

// total pages read
// used by segment v2
Expand Down
1 change: 1 addition & 0 deletions be/src/exec/scan/olap_scanner.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -762,6 +762,7 @@ void OlapScanner::_collect_profile_before_close() {
COUNTER_UPDATE(local_state->_predicate_column_read_seek_counter,
stats.predicate_column_read_seek_num);
COUNTER_UPDATE(local_state->_lazy_read_timer, stats.lazy_read_ns);
COUNTER_UPDATE(local_state->_lazy_read_pruned_timer, stats.lazy_read_pruned_ns);
COUNTER_UPDATE(local_state->_lazy_read_seek_timer, stats.block_lazy_read_seek_ns);
COUNTER_UPDATE(local_state->_lazy_read_seek_counter, stats.block_lazy_read_seek_num);
COUNTER_UPDATE(local_state->_output_col_timer, stats.output_col_ns);
Expand Down
6 changes: 6 additions & 0 deletions be/src/runtime/descriptors.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,9 @@ SlotDescriptor::SlotDescriptor(const PSlotDescriptor& pdesc)
auto convert_to_thrift_column_access_path = [](const PColumnAccessPath& pb_path) {
TColumnAccessPath thrift_path;
thrift_path.type = (TAccessPathType::type)pb_path.type();
if (pb_path.has_version()) {
thrift_path.__set_version(pb_path.version());
}
if (pb_path.has_data_access_path()) {
thrift_path.__isset.data_access_path = true;
for (int i = 0; i < pb_path.data_access_path().path_size(); ++i) {
Expand Down Expand Up @@ -161,6 +164,9 @@ void SlotDescriptor::to_protobuf(PSlotDescriptor* pslot) const {
doris::PColumnAccessPath* pb_path) {
pb_path->Clear();
pb_path->set_type((PAccessPathType)thrift_path.type); // 使用 reinterpret_cast 进行类型转换
if (thrift_path.__isset.version) {
pb_path->set_version(thrift_path.version);
}
if (thrift_path.__isset.data_access_path) {
auto* pb_data = pb_path->mutable_data_access_path();
pb_data->Clear();
Expand Down
5 changes: 5 additions & 0 deletions be/src/runtime/runtime_state.h
Original file line number Diff line number Diff line change
Expand Up @@ -642,6 +642,11 @@ class RuntimeState {
_query_options.enable_aggregate_function_null_v2;
}

bool enable_prune_nested_column() const {
return _query_options.__isset.enable_prune_nested_column &&
_query_options.enable_prune_nested_column;
}

bool is_read_csv_empty_line_as_null() const {
return _query_options.__isset.read_csv_empty_line_as_null &&
_query_options.read_csv_empty_line_as_null;
Expand Down
1 change: 1 addition & 0 deletions be/src/storage/olap_common.h
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,7 @@ struct OlapReaderStatistics {
int64_t lazy_read_ns = 0;
int64_t block_lazy_read_seek_num = 0;
int64_t block_lazy_read_seek_ns = 0;
int64_t lazy_read_pruned_ns = 0;

int64_t raw_rows_read = 0;

Expand Down
1,229 changes: 872 additions & 357 deletions be/src/storage/segment/column_reader.cpp

Large diffs are not rendered by default.

261 changes: 219 additions & 42 deletions be/src/storage/segment/column_reader.h

Large diffs are not rendered by default.

128 changes: 112 additions & 16 deletions be/src/storage/segment/segment_iterator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
#include <gen_cpp/Opcodes_types.h>
#include <gen_cpp/Types_types.h>
#include <gen_cpp/olap_file.pb.h>
#include <glog/logging.h>

#include <algorithm>
#include <boost/iterator/iterator_facade.hpp>
Expand Down Expand Up @@ -119,6 +120,30 @@ namespace segment_v2 {

#include "common/compile_check_begin.h"

class ScopedColumnIteratorReadPhase {
public:
ScopedColumnIteratorReadPhase(ColumnIterator* column_iter, ColumnIterator::ReadPhase mode)
: _column_iter(column_iter) {
DORIS_CHECK(_column_iter != nullptr);
_column_iter->set_read_phase(mode);
}

ScopedColumnIteratorReadPhase(const ScopedColumnIteratorReadPhase&) = delete;
ScopedColumnIteratorReadPhase& operator=(const ScopedColumnIteratorReadPhase&) = delete;

~ScopedColumnIteratorReadPhase() {
// ReadPhase is a per-read phase knob. SegmentIterator only needs a
// temporary PREDICATE/LAZY mode while reading one column in one phase; it
// must be restored before the next column or later normal reads reuse the
// same ColumnIterator. Keep the restoration in one scoped helper instead
// of open-coding the same Defer block at every call site.
_column_iter->set_read_phase(ColumnIterator::ReadPhase::NORMAL);
}

private:
ColumnIterator* _column_iter = nullptr;
};

SegmentIterator::~SegmentIterator() = default;

void SegmentIterator::_init_row_bitmap_by_condition_cache() {
Expand Down Expand Up @@ -419,6 +444,10 @@ Status SegmentIterator::_init_impl(const StorageReadOptions& opts) {
_score_runtime = _opts.score_runtime;
_ann_topn_runtime = _opts.ann_topn_runtime;

_enable_prune_nested_column = _opts.io_ctx.reader_type == ReaderType::READER_QUERY &&
_opts.runtime_state &&
_opts.runtime_state->enable_prune_nested_column();

if (opts.output_columns != nullptr) {
_output_columns = *(opts.output_columns);
}
Expand Down Expand Up @@ -657,8 +686,14 @@ void SegmentIterator::_init_segment_prefetchers() {
? PrefetcherInitMethod::FROM_ROWIDS
: PrefetcherInitMethod::ALL_DATA_BLOCKS;
std::map<PrefetcherInitMethod, std::vector<SegmentPrefetcher*>> prefetchers;
for (const auto& column_iter : _column_iterators) {
for (size_t idx = 0; idx < _column_iterators.size(); ++idx) {
auto cid = cast_set<ColumnId>(idx);
auto* column_iter = _column_iterators[cid].get();
if (column_iter != nullptr) {
ScopedColumnIteratorReadPhase scoped_read_phase {
column_iter, _support_lazy_read_pruned_columns.contains(cid)
? ColumnIterator::ReadPhase::PREDICATE
: ColumnIterator::ReadPhase::NORMAL};
column_iter->collect_prefetchers(prefetchers, init_method);
}
}
Expand Down Expand Up @@ -1984,6 +2019,25 @@ Status SegmentIterator::_vec_init_lazy_materialization() {
if (_is_common_expr_column[cid] || _is_pred_column[cid]) {
auto loc = _schema_block_id_map[cid];
_columns_to_filter.push_back(loc);

const auto field_type = _schema->column(cid)->type();
if (_is_common_expr_column[cid] && _enable_prune_nested_column &&
(field_type == FieldType::OLAP_FIELD_TYPE_STRUCT ||
field_type == FieldType::OLAP_FIELD_TYPE_ARRAY ||
field_type == FieldType::OLAP_FIELD_TYPE_MAP)) {
DCHECK(_column_iterators[cid]);
if (_column_iterators[cid]->read_requirement() ==
ColumnIterator::ReadRequirement::PREDICATE &&
_column_iterators[cid]->has_lazy_read_target()) {
// Only split lazy recovery for complex common expr columns that have
// both predicate-only and non-predicate nested targets. The two requirement
// checks already imply that nested-column pruning happened: without an
// explicit predicate sub-path the parent would not be
// PREDICATE, and without a pruned non-predicate child there
// would be no lazy target to recover after filtering.
_support_lazy_read_pruned_columns.emplace(cid);
}
}
}
}

Expand Down Expand Up @@ -2301,16 +2355,22 @@ Status SegmentIterator::_read_columns_by_index(uint32_t nrows_read_limit, uint16
})
}

auto* column_iter = _column_iterators[cid].get();
ScopedColumnIteratorReadPhase scoped_read_phase {
column_iter, _support_lazy_read_pruned_columns.contains(cid)
? ColumnIterator::ReadPhase::PREDICATE
: ColumnIterator::ReadPhase::NORMAL};

if (is_continuous) {
size_t rows_read = nrows_read;
_opts.stats->predicate_column_read_seek_num += 1;
if (_opts.runtime_state && _opts.runtime_state->enable_profile()) {
SCOPED_RAW_TIMER(&_opts.stats->predicate_column_read_seek_ns);
RETURN_IF_ERROR(_column_iterators[cid]->seek_to_ordinal(_block_rowids[0]));
RETURN_IF_ERROR(column_iter->seek_to_ordinal(_block_rowids[0]));
} else {
RETURN_IF_ERROR(_column_iterators[cid]->seek_to_ordinal(_block_rowids[0]));
RETURN_IF_ERROR(column_iter->seek_to_ordinal(_block_rowids[0]));
}
RETURN_IF_ERROR(_column_iterators[cid]->next_batch(&rows_read, column));
RETURN_IF_ERROR(column_iter->next_batch(&rows_read, column));
if (rows_read != nrows_read) {
return Status::Error<ErrorCode::INTERNAL_ERROR>("nrows({}) != rows_read({})",
nrows_read, rows_read);
Expand All @@ -2330,20 +2390,18 @@ Status SegmentIterator::_read_columns_by_index(uint32_t nrows_read_limit, uint16
_opts.stats->predicate_column_read_seek_num += 1;
if (_opts.runtime_state && _opts.runtime_state->enable_profile()) {
SCOPED_RAW_TIMER(&_opts.stats->predicate_column_read_seek_ns);
RETURN_IF_ERROR(
_column_iterators[cid]->seek_to_ordinal(_block_rowids[processed]));
RETURN_IF_ERROR(column_iter->seek_to_ordinal(_block_rowids[processed]));
} else {
RETURN_IF_ERROR(
_column_iterators[cid]->seek_to_ordinal(_block_rowids[processed]));
RETURN_IF_ERROR(column_iter->seek_to_ordinal(_block_rowids[processed]));
}
RETURN_IF_ERROR(_column_iterators[cid]->next_batch(&rows_read, column));
RETURN_IF_ERROR(column_iter->next_batch(&rows_read, column));
if (rows_read != current_batch_size) {
return Status::Error<ErrorCode::INTERNAL_ERROR>(
"batch nrows({}) != rows_read({})", current_batch_size, rows_read);
}
} else {
RETURN_IF_ERROR(_column_iterators[cid]->read_by_rowids(
&_block_rowids[processed], current_batch_size, column));
RETURN_IF_ERROR(column_iter->read_by_rowids(&_block_rowids[processed],
current_batch_size, column));
}
processed += current_batch_size;
}
Expand Down Expand Up @@ -2483,7 +2541,8 @@ Status SegmentIterator::_read_columns_by_rowids(std::vector<ColumnId>& read_colu
std::vector<rowid_t>& rowid_vector,
uint16_t* sel_rowid_idx, size_t select_size,
MutableColumns* mutable_columns,
bool init_condition_cache) {
bool init_condition_cache,
bool read_for_predicate) {
SCOPED_RAW_TIMER(&_opts.stats->lazy_read_ns);
std::vector<rowid_t> rowids(select_size);

Expand Down Expand Up @@ -2527,10 +2586,45 @@ Status SegmentIterator::_read_columns_by_rowids(std::vector<ColumnId>& read_colu
"SegmentIterator meet invalid column, return columns size {}, cid {}",
_current_return_columns.size(), cid);
}
RETURN_IF_ERROR(_column_iterators[cid]->read_by_rowids(rowids.data(), select_size,
_current_return_columns[cid]));

auto* column_iter = _column_iterators[cid].get();
ScopedColumnIteratorReadPhase scoped_read_phase {
column_iter, read_for_predicate && _support_lazy_read_pruned_columns.contains(cid)
? ColumnIterator::ReadPhase::PREDICATE
: ColumnIterator::ReadPhase::NORMAL};

RETURN_IF_ERROR(column_iter->read_by_rowids(rowids.data(), select_size,
_current_return_columns[cid]));
}

return Status::OK();
}

Status SegmentIterator::_read_lazy_pruned_columns(Block* block) {
if (_support_lazy_read_pruned_columns.empty()) {
return Status::OK();
}

SCOPED_RAW_TIMER(&_opts.stats->lazy_read_pruned_ns);
DorisVector<rowid_t> rowids(_selected_size);
for (size_t i = 0; i < _selected_size; ++i) {
rowids[i] = _block_rowids[_sel_rowid_idx[i]];
}

for (auto cid : _support_lazy_read_pruned_columns) {
// branch-4.2 keeps the column-id -> block-position map inside SegmentIterator;
// master's Schema::column_index() helper does not exist on this branch.
auto loc = cast_set<size_t>(_schema_block_id_map[cid]);
auto column = IColumn::mutate(std::move(block->get_by_position(loc).column));
auto* column_iter = _column_iterators[cid].get();
ScopedColumnIteratorReadPhase scoped_read_phase {column_iter,
ColumnIterator::ReadPhase::LAZY};
if (_selected_size > 0) {
RETURN_IF_ERROR(column_iter->read_by_rowids(rowids.data(), _selected_size, column));
}
column_iter->finalize_lazy_phase(column);
block->get_by_position(loc).column = std::move(column);
}
return Status::OK();
}

Expand Down Expand Up @@ -2735,7 +2829,7 @@ Status SegmentIterator::_next_batch_internal(Block* block) {
SCOPED_RAW_TIMER(&_opts.stats->non_predicate_read_ns);
RETURN_IF_ERROR(_read_columns_by_rowids(
_non_predicate_column_ids, _block_rowids, _sel_rowid_idx.data(),
_selected_size, &_current_return_columns));
_selected_size, &_current_return_columns, false, true));
_replace_version_col_if_needed(_non_predicate_column_ids, _selected_size);
RETURN_IF_ERROR(_process_columns(_non_predicate_column_ids, block));
}
Expand Down Expand Up @@ -2766,7 +2860,7 @@ Status SegmentIterator::_next_batch_internal(Block* block) {
RETURN_IF_ERROR(_read_columns_by_rowids(
_non_predicate_columns, _block_rowids, _sel_rowid_idx.data(),
_selected_size, &_current_return_columns,
_opts.condition_cache_digest && !_find_condition_cache));
_opts.condition_cache_digest && !_find_condition_cache, false));
_replace_version_col_if_needed(_non_predicate_columns, _selected_size);
} else {
if (_opts.condition_cache_digest && !_find_condition_cache) {
Expand All @@ -2778,6 +2872,8 @@ Status SegmentIterator::_next_batch_internal(Block* block) {
}
}
}

RETURN_IF_ERROR(_read_lazy_pruned_columns(block));
}

// step5: output columns
Expand Down
7 changes: 6 additions & 1 deletion be/src/storage/segment/segment_iterator.h
Original file line number Diff line number Diff line change
Expand Up @@ -227,7 +227,9 @@ class SegmentIterator : public RowwiseIterator {
std::vector<rowid_t>& rowid_vector,
uint16_t* sel_rowid_idx, size_t select_size,
MutableColumns* mutable_columns,
bool init_condition_cache = false);
bool init_condition_cache = false,
bool read_for_predicate = false);
[[nodiscard]] Status _read_lazy_pruned_columns(Block* block);

Status copy_column_data_by_selector(IColumn* input_col_ptr, MutableColumnPtr& output_col,
uint16_t* sel_rowid_idx, uint16_t select_size,
Expand Down Expand Up @@ -372,6 +374,9 @@ class SegmentIterator : public RowwiseIterator {
bool _is_need_short_eval = false;
bool _is_need_expr_eval = false;

std::set<ColumnId> _support_lazy_read_pruned_columns;
bool _enable_prune_nested_column = false;

// fields for vectorization execution
std::vector<ColumnId>
_vec_pred_column_ids; // keep columnId of columns for vectorized predicate evaluation
Expand Down
31 changes: 26 additions & 5 deletions be/src/storage/segment/variant/variant_column_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ class ReaderOwnedColumnIterator final : public ColumnIterator {
: _inner(std::move(inner)), _owner(std::move(owner)) {
DCHECK(_inner != nullptr);
set_column_name(_inner->column_name());
set_reading_flag(_inner->reading_flag());
set_read_requirement(_inner->read_requirement());
}

Status init(const ColumnIteratorOptions& opts) override { return _inner->init(opts); }
Expand Down Expand Up @@ -130,15 +130,36 @@ class ReaderOwnedColumnIterator final : public ColumnIterator {
Status set_access_paths(const TColumnAccessPaths& all_access_paths,
const TColumnAccessPaths& predicate_access_paths) override {
RETURN_IF_ERROR(_inner->set_access_paths(all_access_paths, predicate_access_paths));
set_reading_flag(_inner->reading_flag());
ColumnIterator::set_read_requirement_self(_inner->read_requirement());
return Status::OK();
}

void set_need_to_read() override {
_inner->set_need_to_read();
set_reading_flag(_inner->reading_flag());
void set_read_requirement(ReadRequirement requirement) override {
_inner->set_read_requirement(requirement);
ColumnIterator::set_read_requirement_self(_inner->read_requirement());
}

void set_read_requirement_self(ReadRequirement requirement) override {
_inner->set_read_requirement_self(requirement);
ColumnIterator::set_read_requirement_self(_inner->read_requirement());
}

void set_lazy_output_requirement() override {
_inner->set_lazy_output_requirement();
ColumnIterator::set_read_requirement_self(_inner->read_requirement());
}

void set_read_phase(ReadPhase mode) override {
ColumnIterator::set_read_phase(mode);
_inner->set_read_phase(mode);
}

void finalize_lazy_phase(MutableColumnPtr& dst) override { _inner->finalize_lazy_phase(dst); }

bool has_lazy_read_target() const override { return _inner->has_lazy_read_target(); }

bool need_to_read() const override { return _inner->need_to_read(); }

void remove_pruned_sub_iterators() override { _inner->remove_pruned_sub_iterators(); }

Status init_prefetcher(const SegmentPrefetchParams& params) override {
Expand Down
Loading
Loading