Gabriel39 commented on code in PR #66360:
URL: https://github.com/apache/doris/pull/66360#discussion_r3702053743


##########
be/src/format_v2/table_reader.cpp:
##########
@@ -765,6 +765,104 @@ Status TableReader::_build_table_filters_from_conjuncts() 
{
     return Status::OK();
 }
 
+namespace {
+
+bool same_scan_projection(const LocalColumnIndex& lhs, const LocalColumnIndex& 
rhs) {
+    if (lhs.index != rhs.index || lhs.project_all_children != 
rhs.project_all_children ||
+        lhs.children.size() != rhs.children.size()) {
+        return false;
+    }
+    for (size_t index = 0; index < lhs.children.size(); ++index) {
+        if (!same_scan_projection(lhs.children[index], rhs.children[index])) {
+            return false;
+        }
+    }
+    return true;
+}
+
+const LocalColumnIndex* find_scan_projection(const FileScanRequest& request,
+                                             LocalColumnId column_id) {
+    const auto find_by_id = [column_id](const std::vector<LocalColumnIndex>& 
projections) {
+        return std::ranges::find_if(projections, [column_id](const 
LocalColumnIndex& projection) {
+            return projection.column_id() == column_id;
+        });
+    };
+    auto it = find_by_id(request.predicate_columns);
+    if (it != request.predicate_columns.end()) {
+        return &*it;
+    }
+    it = find_by_id(request.non_predicate_columns);
+    return it == request.non_predicate_columns.end() ? nullptr : &*it;
+}
+
+bool same_physical_scan_layout(const FileScanRequest& lhs, const 
FileScanRequest& rhs) {
+    if (lhs.local_positions != rhs.local_positions) {
+        return false;
+    }
+    for (const auto& [column_id, _] : lhs.local_positions) {
+        const auto* lhs_projection = find_scan_projection(lhs, column_id);
+        const auto* rhs_projection = find_scan_projection(rhs, column_id);
+        if (lhs_projection == nullptr || rhs_projection == nullptr ||
+            !same_scan_projection(*lhs_projection, *rhs_projection)) {
+            return false;
+        }
+    }
+    return true;
+}
+
+} // namespace
+
+Status TableReader::refresh_conjuncts(VExprContextSPtrs conjuncts) {
+    _conjuncts = std::move(conjuncts);
+    if (_data_reader.reader == nullptr) {
+        // The split is prepared but its physical reader has not opened yet. 
open_reader() will use
+        // this newest snapshot directly, so no pending request is needed.
+        return Status::OK();
+    }
+    if (!_data_reader.reader->supports_scan_request_refresh()) {
+        return Status::OK();
+    }
+
+    RETURN_IF_ERROR(_build_table_filters_from_conjuncts());
+    // create_scan_request() rebuilds mapping projections in place. Build late 
predicates with an
+    // isolated mapper so the active row group cannot observe an unprepared or 
incompatible mapper
+    // before its physical request reaches the reader's safe activation 
boundary.
+    auto refreshed_mapper = 
_data_reader.reader->create_column_mapper(_mapper_options);
+    DORIS_CHECK(refreshed_mapper != nullptr);
+    RETURN_IF_ERROR(refreshed_mapper->create_mapping(_projected_columns, 
_partition_values,
+                                                     
_data_reader.file_schema));
+    auto refreshed_request = std::make_shared<FileScanRequest>();
+    RETURN_IF_ERROR(refreshed_mapper->create_scan_request(
+            _table_filters, _projected_columns, refreshed_request.get(), 
_runtime_state,
+            _file_scan_request == nullptr ? nullptr : 
&_file_scan_request->local_positions));
+    if (_push_down_agg_type == TPushAggOp::type::COUNT && 
_push_down_count_columns.has_value() &&
+        _push_down_count_columns->empty()) {
+        for (const auto& column : refreshed_request->non_predicate_columns) {
+            
refreshed_request->count_star_placeholder_columns.push_back(column.column_id());
+        }
+    }
+    RETURN_IF_ERROR(customize_file_scan_request(refreshed_request.get()));
+    if (_file_scan_request == nullptr ||
+        !same_physical_scan_layout(*refreshed_request, *_file_scan_request)) {

Review Comment:
   Fixed in fb1c5d1a06. COUNT(*) fallback no longer marks retained columns as 
value-less placeholders while runtime filters are pending, including refreshed 
requests. CountStarFallbackKeepsLateRuntimeFilterCarrierValues reads the first 
row group, applies a late id > 4 filter, and verifies that the remaining real 
carrier values are {5, 6}.



##########
be/src/format_v2/parquet/parquet_scan.cpp:
##########
@@ -908,6 +961,29 @@ void ParquetScanScheduler::reset() {
     reset_current_row_group();
 }
 
+void 
ParquetScanScheduler::set_scan_request(std::shared_ptr<format::FileScanRequest> 
request) {
+    DORIS_CHECK(request != nullptr);
+    _active_request = std::move(request);
+    _pending_request.reset();
+    _predicate_schedule_request = nullptr;
+}
+
+void 
ParquetScanScheduler::queue_scan_request(std::shared_ptr<format::FileScanRequest>
 request) {
+    DORIS_CHECK(request != nullptr);
+    _pending_request = std::move(request);
+}
+
+void 
ParquetScanScheduler::activate_pending_scan_request_at_row_group_boundary() {
+    if (_has_current_row_group || !_pending_predicate_selection.empty() ||
+        _pending_request == nullptr) {
+        return;
+    }
+    // Column readers and predicate schedules retain request-derived state for 
one row group. Swap
+    // only after they are gone; the refreshed request may promote a lazy 
column to a predicate.
+    _active_request = std::move(_pending_request);

Review Comment:
   Fixed in fb1c5d1a06. Activating a refreshed request now invalidates the 
adaptive schedule/statistics and marks unopened row groups for replanning. Each 
remaining group reruns current-request footer-statistics pruning before 
expensive metadata probes. 
LateRequestReplansUnopenedRowGroupsWithFooterStatistics verifies that the 
middle group is rejected before data-page decode (RawRowsRead=4 and one 
min/max-filtered group).



##########
be/src/format_v2/table_reader.cpp:
##########
@@ -765,6 +765,104 @@ Status TableReader::_build_table_filters_from_conjuncts() 
{
     return Status::OK();
 }
 
+namespace {
+
+bool same_scan_projection(const LocalColumnIndex& lhs, const LocalColumnIndex& 
rhs) {
+    if (lhs.index != rhs.index || lhs.project_all_children != 
rhs.project_all_children ||
+        lhs.children.size() != rhs.children.size()) {
+        return false;
+    }
+    for (size_t index = 0; index < lhs.children.size(); ++index) {
+        if (!same_scan_projection(lhs.children[index], rhs.children[index])) {
+            return false;
+        }
+    }
+    return true;
+}
+
+const LocalColumnIndex* find_scan_projection(const FileScanRequest& request,
+                                             LocalColumnId column_id) {
+    const auto find_by_id = [column_id](const std::vector<LocalColumnIndex>& 
projections) {
+        return std::ranges::find_if(projections, [column_id](const 
LocalColumnIndex& projection) {
+            return projection.column_id() == column_id;
+        });
+    };
+    auto it = find_by_id(request.predicate_columns);
+    if (it != request.predicate_columns.end()) {
+        return &*it;
+    }
+    it = find_by_id(request.non_predicate_columns);
+    return it == request.non_predicate_columns.end() ? nullptr : &*it;
+}
+
+bool same_physical_scan_layout(const FileScanRequest& lhs, const 
FileScanRequest& rhs) {
+    if (lhs.local_positions != rhs.local_positions) {
+        return false;
+    }
+    for (const auto& [column_id, _] : lhs.local_positions) {
+        const auto* lhs_projection = find_scan_projection(lhs, column_id);
+        const auto* rhs_projection = find_scan_projection(rhs, column_id);
+        if (lhs_projection == nullptr || rhs_projection == nullptr ||
+            !same_scan_projection(*lhs_projection, *rhs_projection)) {
+            return false;
+        }
+    }
+    return true;
+}
+
+} // namespace
+
+Status TableReader::refresh_conjuncts(VExprContextSPtrs conjuncts) {

Review Comment:
   Fixed in fb1c5d1a06. The refresh path now records TableReader 
RefreshConjunctsTime, FileReaderRefreshScanRequestTime, and concrete Parquet 
RefreshScanRequestTime in their owning profile hierarchy; JNI refreshes use the 
corresponding TableReader/FileReader scopes. The open-reader refresh test 
asserts the new counters are registered.



-- 
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]

Reply via email to