This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new ac92748b945 [refine](exec) Reduce RowDescriptor to an ordered tuple
layout (#68670)
ac92748b945 is described below
commit ac92748b9450e37727187552851859553ac8ffda
Author: Mryange <[email protected]>
AuthorDate: Thu Oct 8 15:38:53 2026 +0800
[refine](exec) Reduce RowDescriptor to an ordered tuple layout (#68670)
This is the first step in a broader BE descriptor refactor.
`RowDescriptor` previously stored a tuple index map, slot counts, and a
variable-length slot flag that could all be derived from its ordered
`TupleDescriptor` list. Those duplicate values made the class heavier
and allowed constructors to report inconsistent slot counts.
`RowDescriptor` now stores only the ordered tuple list, derives the
layout values when needed, and drops unused interfaces. The affected
execution checks use the tuple list directly.
### Release note
None
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [ ] Regression test
- [ ] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [ ] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
be/src/exec/operator/analytic_source_operator.cpp | 6 -
be/src/exec/operator/analytic_source_operator.h | 2 -
be/src/exec/operator/cache_source_operator.cpp | 2 +-
be/src/exec/operator/hashjoin_probe_operator.cpp | 4 +-
.../operator/nested_loop_join_build_operator.cpp | 19 ++--
.../operator/nested_loop_join_probe_operator.cpp | 5 +-
be/src/exec/operator/operator.cpp | 7 +-
be/src/exec/operator/repeat_operator.cpp | 4 +-
.../operator/streaming_aggregation_operator.cpp | 3 +-
be/src/exec/operator/table_function_operator.cpp | 6 +-
be/src/exec/operator/union_source_operator.cpp | 3 +-
be/src/exec/pipeline/pipeline_task.cpp | 3 +-
be/src/exec/scan/scanner.cpp | 8 +-
be/src/runtime/descriptors.cpp | 123 ++-------------------
be/src/runtime/descriptors.h | 79 +------------
.../exec/operator/query_cache_operator_test.cpp | 19 ++--
be/test/runtime/descriptor_test.cpp | 24 ++++
be/test/testutil/mock/mock_descriptors.h | 4 +-
18 files changed, 77 insertions(+), 244 deletions(-)
diff --git a/be/src/exec/operator/analytic_source_operator.cpp
b/be/src/exec/operator/analytic_source_operator.cpp
index 65553e69555..e1dc8b728fe 100644
--- a/be/src/exec/operator/analytic_source_operator.cpp
+++ b/be/src/exec/operator/analytic_source_operator.cpp
@@ -83,10 +83,4 @@ Status AnalyticSourceOperatorX::get_block_impl(RuntimeState*
state, Block* outpu
return Status::OK();
}
-Status AnalyticSourceOperatorX::prepare(RuntimeState* state) {
- RETURN_IF_ERROR(OperatorX<AnalyticLocalState>::prepare(state));
-
DCHECK(_child->operator_row_desc_after_projection().is_prefix_of(_row_descriptor));
- return Status::OK();
-}
-
} // namespace doris
diff --git a/be/src/exec/operator/analytic_source_operator.h
b/be/src/exec/operator/analytic_source_operator.h
index 72273ebd98a..433c2b68763 100644
--- a/be/src/exec/operator/analytic_source_operator.h
+++ b/be/src/exec/operator/analytic_source_operator.h
@@ -50,8 +50,6 @@ public:
bool is_source() const override { return true; }
- Status prepare(RuntimeState* state) override;
-
private:
friend class AnalyticLocalState;
};
diff --git a/be/src/exec/operator/cache_source_operator.cpp
b/be/src/exec/operator/cache_source_operator.cpp
index e4fc8d08df1..a96061c3dd1 100644
--- a/be/src/exec/operator/cache_source_operator.cpp
+++ b/be/src/exec/operator/cache_source_operator.cpp
@@ -186,7 +186,7 @@ Status CacheSourceOperatorX::get_block_impl(RuntimeState*
state, Block* block, b
auto& local_state = get_local_state(state);
SCOPED_TIMER(local_state.exec_time_counter());
- block->clear_column_data(_row_descriptor.num_materialized_slots());
+ block->clear_column_data(_row_descriptor.num_slots());
bool need_clone_empty = block->columns() == 0;
const bool has_cached_block =
diff --git a/be/src/exec/operator/hashjoin_probe_operator.cpp
b/be/src/exec/operator/hashjoin_probe_operator.cpp
index a88bc35a51a..48f896eebf5 100644
--- a/be/src/exec/operator/hashjoin_probe_operator.cpp
+++ b/be/src/exec/operator/hashjoin_probe_operator.cpp
@@ -172,7 +172,7 @@ void HashJoinProbeLocalState::_prepare_probe_block() {
}
_key_columns_holder.clear();
_probe_block.clear_column_data(
-
_parent->get_child()->operator_row_desc_after_projection().num_materialized_slots());
+
_parent->get_child()->operator_row_desc_after_projection().num_slots());
}
HashJoinProbeOperatorX::HashJoinProbeOperatorX(ObjectPool* pool, const
TPlanNode& tnode,
@@ -233,7 +233,7 @@ Status HashJoinProbeOperatorX::pull(doris::RuntimeState*
state, Block* output_bl
RETURN_IF_ERROR(local_state.filter_data_and_build_output(state,
output_block, eos,
&local_state._probe_block, false));
local_state._probe_block.clear_column_data(
-
_child->operator_row_desc_after_projection().num_materialized_slots());
+ _child->operator_row_desc_after_projection().num_slots());
return Status::OK();
}
diff --git a/be/src/exec/operator/nested_loop_join_build_operator.cpp
b/be/src/exec/operator/nested_loop_join_build_operator.cpp
index bc4367fe85a..a472272c9aa 100644
--- a/be/src/exec/operator/nested_loop_join_build_operator.cpp
+++ b/be/src/exec/operator/nested_loop_join_build_operator.cpp
@@ -17,6 +17,7 @@
#include "exec/operator/nested_loop_join_build_operator.h"
+#include <algorithm>
#include <memory>
#include "exec/operator/operator.h"
@@ -83,14 +84,16 @@ Status NestedLoopJoinBuildSinkOperatorX::init(const
TPlanNode& tnode, RuntimeSta
Status NestedLoopJoinBuildSinkOperatorX::prepare(RuntimeState* state) {
RETURN_IF_ERROR(JoinBuildSinkOperatorX<NestedLoopJoinBuildSinkLocalState>::prepare(state));
- size_t num_build_tuples =
-
_child->operator_row_desc_after_projection().tuple_descriptors().size();
-
- for (size_t i = 0; i < num_build_tuples; ++i) {
- TupleDescriptor* build_tuple_desc =
-
_child->operator_row_desc_after_projection().tuple_descriptors()[i];
- auto tuple_idx = _row_descriptor.get_tuple_idx(build_tuple_desc->id());
- RETURN_IF_INVALID_TUPLE_IDX(build_tuple_desc->id(), tuple_idx);
+ const auto& row_tuples = _row_descriptor.tuple_descriptors();
+ for (const auto* build_tuple_desc :
+ _child->operator_row_desc_after_projection().tuple_descriptors()) {
+ if (std::none_of(row_tuples.begin(), row_tuples.end(),
+ [tuple_id = build_tuple_desc->id()](const auto*
tuple_desc) {
+ return tuple_desc->id() == tuple_id;
+ })) {
+ return Status::InternalError("build tuple id {} is not part of the
row descriptor",
+ build_tuple_desc->id());
+ }
}
RETURN_IF_ERROR(VExpr::prepare(_filter_src_expr_ctxs, state,
_child->operator_row_desc_after_projection()));
diff --git a/be/src/exec/operator/nested_loop_join_probe_operator.cpp
b/be/src/exec/operator/nested_loop_join_probe_operator.cpp
index 7b2adb02e0e..2cd8dfd8cfe 100644
--- a/be/src/exec/operator/nested_loop_join_probe_operator.cpp
+++ b/be/src/exec/operator/nested_loop_join_probe_operator.cpp
@@ -1118,9 +1118,8 @@ Status
NestedLoopJoinProbeOperatorX::prepare(RuntimeState* state) {
for (auto& conjunct : _mark_join_conjuncts) {
RETURN_IF_ERROR(conjunct->prepare(state, join_row_desc()));
}
- _num_probe_side_columns =
_child->operator_row_desc_after_projection().num_materialized_slots();
- _num_build_side_columns =
-
_build_side_child->operator_row_desc_after_projection().num_materialized_slots();
+ _num_probe_side_columns =
_child->operator_row_desc_after_projection().num_slots();
+ _num_build_side_columns =
_build_side_child->operator_row_desc_after_projection().num_slots();
for (const auto& conjunct : _join_conjuncts) {
conjunct->root()->collect_slot_column_ids(_lazy_eval_column_ids);
}
diff --git a/be/src/exec/operator/operator.cpp
b/be/src/exec/operator/operator.cpp
index 796cc116969..a885277a745 100644
--- a/be/src/exec/operator/operator.cpp
+++ b/be/src/exec/operator/operator.cpp
@@ -305,8 +305,7 @@ Status OperatorXBase::close(RuntimeState* state) {
}
void PipelineXLocalStateBase::clear_origin_block() {
- _origin_block.clear_column_data(
-
_parent->operator_row_desc_before_projection().num_materialized_slots());
+
_origin_block.clear_column_data(_parent->operator_row_desc_before_projection().num_slots());
}
Status PipelineXLocalStateBase::filter_block(const VExprContextSPtrs&
expr_contexts, Block* block) {
@@ -398,7 +397,7 @@ Status OperatorXBase::do_projections(RuntimeState* state,
Block* origin_block,
}
origin_block->clear_column_data(
-
local_state->_parent->operator_row_desc_before_projection().num_materialized_slots());
+
local_state->_parent->operator_row_desc_before_projection().num_slots());
DCHECK_EQ(output_block->rows(), rows);
return Status::OK();
@@ -746,7 +745,7 @@ Status
StatefulOperatorX<LocalStateType>::get_block_impl(RuntimeState* state, Bl
if (need_more_input_data(state)) {
local_state._child_block->clear_column_data(
OperatorX<LocalStateType>::_child->operator_row_desc_after_projection()
- .num_materialized_slots());
+ .num_slots());
RETURN_IF_ERROR(OperatorX<LocalStateType>::_child->get_block_after_projects(
state, local_state._child_block.get(),
&local_state._child_eos));
*eos = local_state._child_eos;
diff --git a/be/src/exec/operator/repeat_operator.cpp
b/be/src/exec/operator/repeat_operator.cpp
index a68009eff22..acfd4fbd843 100644
--- a/be/src/exec/operator/repeat_operator.cpp
+++ b/be/src/exec/operator/repeat_operator.cpp
@@ -228,7 +228,7 @@ Status RepeatOperatorX::pull(doris::RuntimeState* state,
Block* output_block, bo
if (_repeat_id_idx >= _repeat_id_list_size) {
_intermediate_block->clear();
_child_block.clear_column_data(
-
_child->operator_row_desc_after_projection().num_materialized_slots());
+
_child->operator_row_desc_after_projection().num_slots());
_repeat_id_idx = 0;
}
} else if (local_state._expr_ctxs.empty()) {
@@ -246,7 +246,7 @@ Status RepeatOperatorX::pull(doris::RuntimeState* state,
Block* output_block, bo
if (_repeat_id_idx >= _repeat_id_list_size) {
_intermediate_block->clear();
_child_block.clear_column_data(
-
_child->operator_row_desc_after_projection().num_materialized_slots());
+
_child->operator_row_desc_after_projection().num_slots());
_repeat_id_idx = 0;
}
}
diff --git a/be/src/exec/operator/streaming_aggregation_operator.cpp
b/be/src/exec/operator/streaming_aggregation_operator.cpp
index f7549c15724..de35cec0da7 100644
--- a/be/src/exec/operator/streaming_aggregation_operator.cpp
+++ b/be/src/exec/operator/streaming_aggregation_operator.cpp
@@ -1127,8 +1127,7 @@ Status StreamingAggOperatorX::push(RuntimeState* state,
Block* in_block, bool eo
RETURN_IF_ERROR(
local_state.do_pre_agg(state, in_block,
local_state._pre_aggregated_block.get()));
}
- in_block->clear_column_data(
-
_child->operator_row_desc_after_projection().num_materialized_slots());
+
in_block->clear_column_data(_child->operator_row_desc_after_projection().num_slots());
return Status::OK();
}
diff --git a/be/src/exec/operator/table_function_operator.cpp
b/be/src/exec/operator/table_function_operator.cpp
index 1fa3cd7dbb9..8557474b599 100644
--- a/be/src/exec/operator/table_function_operator.cpp
+++ b/be/src/exec/operator/table_function_operator.cpp
@@ -467,7 +467,7 @@ Status
TableFunctionLocalState::_get_expanded_block_block_fast_path(
}
_child_block->clear_column_data(_parent->cast<TableFunctionOperatorX>()
._child->operator_row_desc_after_projection()
- .num_materialized_slots());
+ .num_slots());
_reset_block_fast_path_state();
}
@@ -763,7 +763,7 @@ Status
TableFunctionLocalState::_get_expanded_block_for_outer_conjuncts(RuntimeS
_child_rows_has_output.clear();
_child_block->clear_column_data(_parent->cast<TableFunctionOperatorX>()
._child->operator_row_desc_after_projection()
- .num_materialized_slots());
+ .num_slots());
}
}
@@ -791,7 +791,7 @@ void TableFunctionLocalState::process_next_child_row() {
if (!_need_to_handle_outer_conjuncts) {
_child_block->clear_column_data(_parent->cast<TableFunctionOperatorX>()
._child->operator_row_desc_after_projection()
- .num_materialized_slots());
+ .num_slots());
}
_cur_child_offset = -1;
_reset_block_fast_path_state();
diff --git a/be/src/exec/operator/union_source_operator.cpp
b/be/src/exec/operator/union_source_operator.cpp
index 9854137d04d..55538d89cff 100644
--- a/be/src/exec/operator/union_source_operator.cpp
+++ b/be/src/exec/operator/union_source_operator.cpp
@@ -117,8 +117,7 @@ Status UnionSourceOperatorX::get_block_impl(RuntimeState*
state, Block* block, b
return Status::OK();
}
block->swap(*queue_block.block);
- queue_block.block->clear_column_data(
-
operator_row_desc_before_projection().num_materialized_slots());
+
queue_block.block->clear_column_data(operator_row_desc_before_projection().num_slots());
local_state._shared_state->data_queue.push_free_block(std::move(queue_block));
}
local_state.reached_limit(block, eos);
diff --git a/be/src/exec/pipeline/pipeline_task.cpp
b/be/src/exec/pipeline/pipeline_task.cpp
index 515f10a6657..1bffe2e764a 100644
--- a/be/src/exec/pipeline/pipeline_task.cpp
+++ b/be/src/exec/pipeline/pipeline_task.cpp
@@ -569,8 +569,7 @@ Status PipelineTask::execute(bool* done) {
Defer defer {[&]() {
// If this run is pended by a spilling request, the block will be
output in next run.
if (!_spilling) {
- _block->clear_column_data(
-
_root->operator_row_desc_after_projection().num_materialized_slots());
+
_block->clear_column_data(_root->operator_row_desc_after_projection().num_slots());
}
}};
// `_wake_up_early` must be after `_is_blocked()`
diff --git a/be/src/exec/scan/scanner.cpp b/be/src/exec/scan/scanner.cpp
index 33deac162a9..f16a46f3a4a 100644
--- a/be/src/exec/scan/scanner.cpp
+++ b/be/src/exec/scan/scanner.cpp
@@ -91,7 +91,7 @@ Status Scanner::get_block_after_projects(RuntimeState* state,
Block* block, bool
SCOPED_CONCURRENCY_COUNT(ConcurrencyStatsManager::instance().vscanner_get_block);
const auto& row_descriptor =
_local_state->_parent->operator_row_desc_before_projection();
if (_has_projection) {
-
_origin_block.clear_column_data(row_descriptor.num_materialized_slots());
+ _origin_block.clear_column_data(row_descriptor.num_slots());
if (!_can_merge_padding_blocks(_padding_block, _origin_block)) {
DORIS_CHECK(_padding_block.empty())
<< "padding policy must remain stable for one scanner";
@@ -111,7 +111,7 @@ Status Scanner::get_block_after_projects(RuntimeState*
state, Block* block, bool
// The merged tail can be larger than the target batch, but
each source block is
// already bounded by the lower scanner.
RETURN_IF_ERROR(_merge_padding_block());
-
_origin_block.clear_column_data(row_descriptor.num_materialized_slots());
+ _origin_block.clear_column_data(row_descriptor.num_slots());
break;
}
if (_origin_block.rows() >= min_batch_size) {
@@ -121,7 +121,7 @@ Status Scanner::get_block_after_projects(RuntimeState*
state, Block* block, bool
if (_origin_block.rows() + _padding_block.rows() <=
state->batch_size() &&
_origin_block.bytes() + _padding_block.bytes() <=
block_max_bytes) {
RETURN_IF_ERROR(_merge_padding_block());
-
_origin_block.clear_column_data(row_descriptor.num_materialized_slots());
+ _origin_block.clear_column_data(row_descriptor.num_slots());
} else {
if (_origin_block.rows() < _padding_block.rows()) {
_padding_block.swap(_origin_block);
@@ -267,7 +267,7 @@ Status Scanner::_do_projections(Block* origin_block, Block*
output_block) {
}
origin_block->clear_column_data(
-
_local_state->_parent->operator_row_desc_before_projection().num_materialized_slots());
+
_local_state->_parent->operator_row_desc_before_projection().num_slots());
DCHECK_EQ(output_block->rows(), rows);
return Status::OK();
diff --git a/be/src/runtime/descriptors.cpp b/be/src/runtime/descriptors.cpp
index 66b408879dd..d3b8862e3a4 100644
--- a/be/src/runtime/descriptors.cpp
+++ b/be/src/runtime/descriptors.cpp
@@ -46,7 +46,6 @@
#include "util/string_util.h"
namespace doris {
-const int RowDescriptor::INVALID_IDX = -1;
SlotDescriptor::SlotDescriptor(const TSlotDescriptor& tdesc)
: _id(tdesc.id),
@@ -484,114 +483,22 @@ int TupleDescriptor::get_column_id(SlotId slot_id) const
{
RowDescriptor::RowDescriptor(const DescriptorTbl& desc_tbl,
const std::vector<TTupleId>& row_tuples) {
DCHECK_GT(row_tuples.size(), 0);
- _num_materialized_slots = 0;
- _num_slots = 0;
for (int row_tuple : row_tuples) {
TupleDescriptor* tupleDesc = desc_tbl.get_tuple_descriptor(row_tuple);
- _num_materialized_slots += tupleDesc->num_materialized_slots();
- _num_slots += tupleDesc->slots().size();
_tuple_desc_map.push_back(tupleDesc);
DCHECK(_tuple_desc_map.back() != nullptr);
}
-
- init_tuple_idx_map();
- init_has_varlen_slots();
-}
-
-RowDescriptor::RowDescriptor(TupleDescriptor* tuple_desc) : _tuple_desc_map(1,
tuple_desc) {
- init_tuple_idx_map();
- init_has_varlen_slots();
- _num_slots = static_cast<int32_t>(tuple_desc->slots().size());
-}
-
-RowDescriptor::RowDescriptor(const RowDescriptor& lhs_row_desc, const
RowDescriptor& rhs_row_desc) {
- _tuple_desc_map.insert(_tuple_desc_map.end(),
lhs_row_desc._tuple_desc_map.begin(),
- lhs_row_desc._tuple_desc_map.end());
- _tuple_desc_map.insert(_tuple_desc_map.end(),
rhs_row_desc._tuple_desc_map.begin(),
- rhs_row_desc._tuple_desc_map.end());
- init_tuple_idx_map();
- init_has_varlen_slots();
-
- _num_slots = lhs_row_desc.num_slots() + rhs_row_desc.num_slots();
-}
-
-void RowDescriptor::init_tuple_idx_map() {
- // find max id
- TupleId max_id = 0;
- for (auto& i : _tuple_desc_map) {
- max_id = std::max(i->id(), max_id);
- }
-
- _tuple_idx_map.resize(max_id + 1, INVALID_IDX);
- for (int i = 0; i < _tuple_desc_map.size(); ++i) {
- _tuple_idx_map[_tuple_desc_map[i]->id()] = i;
- }
}
-void RowDescriptor::init_has_varlen_slots() {
- _has_varlen_slots = false;
- for (auto& i : _tuple_desc_map) {
- if (i->has_varlen_slots()) {
- _has_varlen_slots = true;
- break;
- }
- }
-}
-
-int RowDescriptor::get_tuple_idx(TupleId id) const {
- // comment CHECK temporarily to make fuzzy test run smoothly
- // DCHECK_LT(id, _tuple_idx_map.size()) << "RowDescriptor: " <<
debug_string();
- if (_tuple_idx_map.size() <= id) {
- return RowDescriptor::INVALID_IDX;
- }
- return _tuple_idx_map[id];
-}
-
-void RowDescriptor::to_thrift(std::vector<TTupleId>* row_tuple_ids) {
- row_tuple_ids->clear();
-
- for (auto& i : _tuple_desc_map) {
- row_tuple_ids->push_back(i->id());
- }
-}
+RowDescriptor::RowDescriptor(TupleDescriptor* tuple_desc) : _tuple_desc_map(1,
tuple_desc) {}
-void RowDescriptor::to_protobuf(
- google::protobuf::RepeatedField<google::protobuf::int32>*
row_tuple_ids) const {
- row_tuple_ids->Clear();
- for (auto* desc : _tuple_desc_map) {
- row_tuple_ids->Add(desc->id());
+int RowDescriptor::num_slots() const {
+ int count = 0;
+ for (const auto* tuple_desc : _tuple_desc_map) {
+ count += tuple_desc->slots().size();
}
-}
-
-bool RowDescriptor::is_prefix_of(const RowDescriptor& other_desc) const {
- if (_tuple_desc_map.size() > other_desc._tuple_desc_map.size()) {
- return false;
- }
-
- for (int i = 0; i < _tuple_desc_map.size(); ++i) {
- // pointer comparison okay, descriptors are unique
- if (_tuple_desc_map[i] != other_desc._tuple_desc_map[i]) {
- return false;
- }
- }
-
- return true;
-}
-
-bool RowDescriptor::equals(const RowDescriptor& other_desc) const {
- if (_tuple_desc_map.size() != other_desc._tuple_desc_map.size()) {
- return false;
- }
-
- for (int i = 0; i < _tuple_desc_map.size(); ++i) {
- // pointer comparison okay, descriptors are unique
- if (_tuple_desc_map[i] != other_desc._tuple_desc_map[i]) {
- return false;
- }
- }
-
- return true;
+ return count;
}
std::string RowDescriptor::debug_string() const {
@@ -606,27 +513,17 @@ std::string RowDescriptor::debug_string() const {
}
ss << "] ";
- ss << "tuple_id_map: [";
- for (int i = 0; i < _tuple_idx_map.size(); ++i) {
- ss << _tuple_idx_map[i];
- if (i != _tuple_idx_map.size() - 1) {
- ss << ", ";
- }
- }
- ss << "] ";
-
return ss.str();
}
int RowDescriptor::get_column_id(int slot_id) const {
int column_id_counter = 0;
for (auto* const tuple_desc : _tuple_desc_map) {
- for (auto* const slot : tuple_desc->slots()) {
- if (slot->id() == slot_id) {
- return column_id_counter;
- }
- column_id_counter++;
+ int tuple_column_id = tuple_desc->get_column_id(slot_id);
+ if (tuple_column_id != -1) {
+ return column_id_counter + tuple_column_id;
}
+ column_id_counter += tuple_desc->slots().size();
}
return -1;
}
diff --git a/be/src/runtime/descriptors.h b/be/src/runtime/descriptors.h
index 7e2e048c27c..cee1a3048f6 100644
--- a/be/src/runtime/descriptors.h
+++ b/be/src/runtime/descriptors.h
@@ -24,7 +24,6 @@
#include <gen_cpp/Exprs_types.h>
#include <gen_cpp/Types_types.h>
#include <glog/logging.h>
-#include <google/protobuf/stubs/port.h>
#include <cstdint>
#include <ostream>
@@ -42,11 +41,6 @@
#include "core/data_type/define_primitive_type.h"
#include "storage/utils.h"
-namespace google::protobuf {
-template <typename Element>
-class RepeatedField;
-} // namespace google::protobuf
-
namespace doris {
class ObjectPool;
class PTupleDescriptor;
@@ -441,96 +435,31 @@ private:
#endif
};
-#define RETURN_IF_INVALID_TUPLE_IDX(tuple_id, tuple_idx)
\
- do {
\
- if (UNLIKELY(RowDescriptor::INVALID_IDX == tuple_idx)) {
\
- return Status::InternalError("failed to get tuple idx with tuple
id: {}", tuple_id); \
- }
\
- } while (false)
-
-// Records positions of tuples within row produced by ExecNode.
-// TODO: this needs to differentiate between tuples contained in row
-// and tuples produced by ExecNode (parallel to PlanNode.rowTupleIds and
-// PlanNode.tupleIds); right now, we conflate the two (and distinguish based on
-// context; for instance, HdfsScanNode uses these tids to create row batches,
ie, the
-// first case, whereas TopNNode uses these tids to copy output rows, ie, the
second
-// case)
+// Describes the ordered tuples whose slots form an operator's block layout.
class RowDescriptor {
public:
RowDescriptor(const DescriptorTbl& desc_tbl, const std::vector<TTupleId>&
row_tuples);
- // standard copy c'tor, made explicit here
- RowDescriptor(const RowDescriptor& desc)
- : _tuple_desc_map(desc._tuple_desc_map),
- _tuple_idx_map(desc._tuple_idx_map),
- _has_varlen_slots(desc._has_varlen_slots) {
- auto it = desc._tuple_desc_map.begin();
- for (; it != desc._tuple_desc_map.end(); ++it) {
- _num_materialized_slots += (*it)->num_materialized_slots();
- _num_slots += (*it)->slots().size();
- }
- }
-
- RowDescriptor& operator=(const RowDescriptor&) = default;
-
RowDescriptor(TupleDescriptor* tuple_desc);
- RowDescriptor(const RowDescriptor& lhs_row_desc, const RowDescriptor&
rhs_row_desc);
-
- // dummy descriptor, needed for the JNI EvalPredicate() function
+ // Empty layout for expressions that do not reference slots.
RowDescriptor() = default;
MOCK_DEFINE(virtual ~RowDescriptor() = default;)
- int num_materialized_slots() const { return _num_materialized_slots; }
-
- int num_slots() const { return _num_slots; }
-
- static const int INVALID_IDX;
-
- // Returns INVALID_IDX if id not part of this row.
- int get_tuple_idx(TupleId id) const;
-
- // Return true if any Tuple has variable length slots.
- bool has_varlen_slots() const { return _has_varlen_slots; }
+ int num_slots() const;
// Return descriptors for all tuples in this row, in order of appearance.
MOCK_FUNCTION const std::vector<TupleDescriptor*>& tuple_descriptors()
const {
return _tuple_desc_map;
}
- // Populate row_tuple_ids with our ids.
- void to_thrift(std::vector<TTupleId>* row_tuple_ids);
- void to_protobuf(google::protobuf::RepeatedField<google::protobuf::int32>*
row_tuple_ids) const;
-
- // Return true if the tuple ids of this descriptor are a prefix
- // of the tuple ids of other_desc.
- bool is_prefix_of(const RowDescriptor& other_desc) const;
-
- // Return true if the tuple ids of this descriptor match tuple ids of
other desc.
- bool equals(const RowDescriptor& other_desc) const;
-
std::string debug_string() const;
int get_column_id(int slot_id) const;
private:
- // Initializes tupleIdxMap during c'tor using the _tuple_desc_map.
- void init_tuple_idx_map();
-
- // Initializes _has_varlen_slots during c'tor using the _tuple_desc_map.
- void init_has_varlen_slots();
-
- // map from position of tuple w/in row to its descriptor
+ // Tuples in block column order; descriptors are owned elsewhere.
std::vector<TupleDescriptor*> _tuple_desc_map;
-
- // map from TupleId to position of tuple w/in row
- std::vector<int> _tuple_idx_map;
-
- // Provide quick way to check if there are variable length slots.
- bool _has_varlen_slots = false;
-
- int _num_materialized_slots = 0;
- int _num_slots = 0;
};
} // namespace doris
diff --git a/be/test/exec/operator/query_cache_operator_test.cpp
b/be/test/exec/operator/query_cache_operator_test.cpp
index dc385d89ed5..1e0735eab38 100644
--- a/be/test/exec/operator/query_cache_operator_test.cpp
+++ b/be/test/exec/operator/query_cache_operator_test.cpp
@@ -748,14 +748,9 @@ TEST_F(QueryCacheOperatorTest,
test_hit_cache_multi_block_reordered_slots) {
source->_cache_param = cache_param;
source->_query_cache_runtime =
std::make_shared<QueryCacheRuntime>(cache_param, query_cache);
- // In production this operator is also built without a plan node, so its
- // _row_descriptor reports zero materialized slots and clear_column_data()
- // wipes the block every pull -- the schema-carrying reused-block shape is
- // accidentally unreachable today. Report a real slot count here to pin
- // the permute-before-merge invariant against the natural cleanups (a real
- // row descriptor, or removing the redundant per-operator wipe) that would
- // make that shape live.
- source->_row_descriptor._num_materialized_slots = 2;
+ // Give the source the real output layout so clear_column_data() preserves
+ // the schema-carrying reused block between pulls.
+ source->_row_descriptor = *row_desc;
create_local_state();
EXPECT_EQ(source_local_state->_slot_orders, (std::vector<int> {10, 11}));
@@ -849,7 +844,7 @@ TEST_F(QueryCacheOperatorTest,
test_incremental_reordered_write_back) {
}
// Pin the schema-carrying reused-block shape; see the rationale in
// test_hit_cache_multi_block_reordered_slots.
- source->_row_descriptor._num_materialized_slots = 2;
+ source->_row_descriptor = *row_desc;
create_local_state();
EXPECT_TRUE(source_local_state->_is_incremental);
@@ -941,8 +936,8 @@ TEST_F(QueryCacheOperatorTest,
test_failed_final_delta_merge_publishes_nothing)
// malformed two-column final block makes the merge fail determinately.
sink = std::make_unique<CacheSinkOperatorX>();
source = std::make_unique<CacheSourceOperatorX>();
- child_op->set_mock_row_desc(std::unique_ptr<MockRowDescriptor>(
- new MockRowDescriptor {{std::make_shared<DataTypeInt64>()},
&pool}));
+ auto* row_desc = new MockRowDescriptor
{{std::make_shared<DataTypeInt64>()}, &pool};
+ child_op->set_mock_row_desc(std::unique_ptr<MockRowDescriptor>(row_desc));
EXPECT_TRUE(source->set_child(child_op));
TQueryCacheParam cache_param;
cache_param.node_id = 0;
@@ -982,7 +977,7 @@ TEST_F(QueryCacheOperatorTest,
test_failed_final_delta_merge_publishes_nothing)
// descriptor the reused block is wiped every pull and re-cloned from the
// incoming block, so the malformed final block would merge into a clone
// of itself and SUCCEED, unbinding this test from the bug it guards.
- source->_row_descriptor._num_materialized_slots = 1;
+ source->_row_descriptor = *row_desc;
create_local_state();
EXPECT_TRUE(source_local_state->_need_insert_cache);
diff --git a/be/test/runtime/descriptor_test.cpp
b/be/test/runtime/descriptor_test.cpp
index 69402ed1b35..7905c7875f4 100644
--- a/be/test/runtime/descriptor_test.cpp
+++ b/be/test/runtime/descriptor_test.cpp
@@ -216,6 +216,30 @@ TEST(TupleDescriptorTest, GetColumnId) {
EXPECT_EQ(tuple_desc->get_column_id(100), -1);
}
+TEST(RowDescriptorTest, MultiTupleLayoutAndSingleTupleCount) {
+ ObjectPool pool;
+ DescriptorTblBuilder builder(&pool);
+ builder.declare_tuple() << std::make_shared<DataTypeInt32>();
+ builder.declare_tuple() << std::make_shared<DataTypeInt64>()
+ << std::make_shared<DataTypeInt32>();
+ auto* desc_tbl = builder.build();
+
+ RowDescriptor single(desc_tbl->get_tuple_descriptor(1));
+ EXPECT_EQ(single.num_slots(), 2);
+
+ RowDescriptor row(*desc_tbl, {0, 1});
+ EXPECT_EQ(row.num_slots(), 3);
+ ASSERT_EQ(row.tuple_descriptors().size(), 2);
+ EXPECT_EQ(row.tuple_descriptors()[0], desc_tbl->get_tuple_descriptor(0));
+ EXPECT_EQ(row.tuple_descriptors()[1], desc_tbl->get_tuple_descriptor(1));
+
EXPECT_EQ(row.get_column_id(desc_tbl->get_tuple_descriptor(1)->slots()[0]->id()),
1);
+
EXPECT_EQ(row.get_column_id(desc_tbl->get_tuple_descriptor(1)->slots()[1]->id()),
2);
+
+ RowDescriptor copied(row);
+ EXPECT_EQ(copied.num_slots(), row.num_slots());
+ EXPECT_EQ(copied.tuple_descriptors(), row.tuple_descriptors());
+}
+
TEST_F(SlotDescriptorTest, AccessPathsPreservedThroughProtobuf) {
TColumnAccessPath data_path;
data_path.__set_version(g_Descriptors_constants.TCOLUMN_ACCESS_PATH_VERSION_TYPED);
diff --git a/be/test/testutil/mock/mock_descriptors.h
b/be/test/testutil/mock/mock_descriptors.h
index cb8833cf8d8..70b84c46861 100644
--- a/be/test/testutil/mock/mock_descriptors.h
+++ b/be/test/testutil/mock/mock_descriptors.h
@@ -45,13 +45,11 @@ public:
auto* slot = pool->add(new SlotDescriptor());
slot->_type = type;
slots.push_back(slot);
- _num_slots++;
}
auto* tuple_desc = pool->add(new MockTupleDescriptor());
tuple_desc->Slots = slots;
tuple_desc_map.push_back(tuple_desc);
_tuple_desc_map.push_back(tuple_desc);
- _num_materialized_slots = static_cast<int>(types.size());
}
const std::vector<TupleDescriptor*>& tuple_descriptors() const override {
return tuple_desc_map;
@@ -145,4 +143,4 @@ private:
mutable std::unordered_map<SlotId, std::unique_ptr<SlotDescriptor>>
_slot_descriptors;
};
-} // namespace doris
\ No newline at end of file
+} // namespace doris
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]