github-actions[bot] commented on code in PR #68569:
URL: https://github.com/apache/doris/pull/68569#discussion_r4121071539
##########
be/src/exec/operator/materialization_opertor.cpp:
##########
@@ -189,24 +189,40 @@ std::string format_counter_array(size_t size, TUnit::type
unit, GetValue get_val
} // namespace
void MaterializationSharedState::get_block(Block* block) {
- for (int i = 0, j = 0, rowid_to_block_loc = rowid_locs[j]; i <
origin_block.columns(); i++) {
- if (i != rowid_to_block_loc) {
+ DCHECK(output_ready);
+ size_t j = 0;
+ for (size_t i = 0; i < origin_block.columns(); i++) {
+ if (j >= rowid_locs.size() || i != static_cast<size_t>(rowid_locs[j]))
{
block->insert(origin_block.get_by_position(i));
} else {
- auto response_block = response_blocks[j].to_block();
- for (int k = 0; k < response_block.columns(); k++) {
- auto& data = response_block.get_by_position(k);
+ auto source_block = response_blocks[j].to_block();
+ for (int k = 0; k < source_block.columns(); k++) {
+ auto& data = source_block.get_by_position(k);
response_blocks[j].mutable_columns()[k] =
data.column->clone_empty();
block->insert(data);
}
- if (++j < rowid_locs.size()) {
- rowid_to_block_loc = rowid_locs[j];
- }
+ ++j;
}
}
origin_block.clear();
}
+void MaterializationSharedState::clear_current_batch() {
+ origin_block.clear();
+ for (auto& response_block : response_blocks) {
+ for (auto& column : response_block.mutable_columns()) {
+ column->clear();
+ }
+ }
+ for (auto& block_order : block_order_results) {
+ block_order.clear();
+ }
+ _backend_rows_count.clear();
+ _max_rows_per_backend = 0;
Review Comment:
[P3] Preserve the peak in `MaxRowsPerBackend`. Resetting
`_max_rows_per_backend` here makes the subsequent `push()` overwrite the
profile counter with only the current fetch batch's maximum. If one backend
handles 20,000 rows in the first batch and 1 in the final batch, the completed
profile reports 1. Keep the profile counter at the maximum seen across batches
while resetting batch-local request accounting.
##########
be/src/exec/operator/materialization_opertor.cpp:
##########
@@ -189,24 +189,40 @@ std::string format_counter_array(size_t size, TUnit::type
unit, GetValue get_val
} // namespace
void MaterializationSharedState::get_block(Block* block) {
- for (int i = 0, j = 0, rowid_to_block_loc = rowid_locs[j]; i <
origin_block.columns(); i++) {
- if (i != rowid_to_block_loc) {
+ DCHECK(output_ready);
Review Comment:
[P2] Update the existing materialization tests for the new state contract.
Three `MaterializationSharedStateTest` cases call `get_block()` directly after
`merge_multi_response()` without setting `output_ready`, so this DCHECK fails
in debug builds. `TestCreateMultiGetResult` also still expects `eos` to become
true during `create_muiltget_result(..., true, ...)`, although only `input_eos`
is set now. These failures follow directly from the changed code; please update
those tests and add a multi-batch case for the new behavior.
##########
be/src/exec/operator/materialization_opertor.cpp:
##########
@@ -545,13 +566,13 @@ Status
MaterializationSharedState::create_muiltget_result(const Columns& columns
}
}
- eos = child_eos;
+ input_eos = input_eos || child_eos;
if (eos && gc_id_map) {
Review Comment:
[P2] Set GC on the final input batch. `input_eos` becomes true immediately
above, but this condition checks output `eos`, which is only set in `pull`
after the RPCs are sent. Thus the top materialization node never sets
`gc_id_map` on its final request. `RowIdStorageReader` only removes the query's
ID map when that flag is true, leaving file mappings and pinned temporary
rowsets until timeout-based GC.
```suggestion
if (input_eos && gc_id_map) {
```
##########
be/src/exec/operator/materialization_opertor.cpp:
##########
@@ -693,11 +720,39 @@ Status MaterializationOperator::push(RuntimeState* state,
Block* in_block, bool
in_block->get_by_position(local_state._materialization_state.rowid_locs[i])
.column);
}
- local_state._materialization_state.origin_block.swap(*in_block);
+ auto& origin_block =
local_state._materialization_state.origin_block;
+ if (origin_block.columns() == 0) {
+ origin_block.swap(*in_block);
+ } else {
+ DCHECK_EQ(origin_block.columns(), in_block->columns());
+ DCHECK_EQ(origin_block.get_names(), in_block->get_names());
+ for (size_t i = 0; i < origin_block.columns(); ++i) {
+ origin_block.replace_by_position_if_const(i);
+ }
+ auto destination_columns =
origin_block.mutate_columns_scoped();
+ for (size_t i = 0; i < origin_block.columns(); ++i) {
+ const auto& source_column =
+
in_block->get_by_position(i).column->convert_to_full_column_if_const();
+
destination_columns.mutable_columns()[i]->insert_range_from(*source_column, 0,
+
in_block->rows());
+ }
+ in_block->clear();
+ }
}
RETURN_IF_ERROR(local_state._materialization_state.create_muiltget_result(columns,
eos,
_gc_id_map));
+ const auto should_fetch =
Review Comment:
[P1] Bound accumulated batches by bytes as well as rows. With the default
20,000-row target, 8 KiB incompressible lazy values produce a 156 MiB fetch;
the deserialized and copied columns coexist at over 312 MiB, so a 256 MiB query
limit that fits the old 8,160-row fetches can now fail. The same whole-block
append can combine two individually valid 2.07 GiB eager string blocks past
`ColumnString`'s 4 GiB limit and fail even with enough query memory. Flush or
split before the eager-column limit and bound or reserve the fetched-response
memory.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]