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]