From 9b5ff496ed517a202b7262f1f0b70c351544db33 Mon Sep 17 00:00:00 2001 From: Mryange Date: Mon, 20 Jul 2026 15:04:21 +0800 Subject: [PATCH 1/2] [fix](be) Finish pending scanners after shared limit ### What problem does this PR solve? Issue Number: None Related PR: #62222 Problem Summary: Queries with a shared scan LIMIT could leave pending scan tasks unscheduled after the shared budget reached zero. With no scanner in flight, the scanner context never reached EOS and the pipeline waited until query timeout. Submit pending scanners so they observe the exhausted limit, return EOS, and complete the scanner context lifecycle. ### Release note Fix query timeout when a LIMIT is exhausted with pending parallel scan tasks. ### Check List (For Author) - Test: Regression test - Regression test: test_shared_scan_limit_pending_tasks - Manual test: case 500078775 passed 10/10 with local shuffle disabled - Behavior changed: Yes, affected LIMIT queries now finish instead of timing out - Does this need documentation: No --- be/src/exec/scan/scanner_context.cpp | 5 - .../test_shared_scan_limit_pending_tasks.out | 7 + ...est_shared_scan_limit_pending_tasks.groovy | 175 ++++++++++++++++++ 3 files changed, 182 insertions(+), 5 deletions(-) create mode 100644 regression-test/data/correctness_p0/test_shared_scan_limit_pending_tasks.out create mode 100644 regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy diff --git a/be/src/exec/scan/scanner_context.cpp b/be/src/exec/scan/scanner_context.cpp index d6c7347e4a54d9..abd3dce06f31cf 100644 --- a/be/src/exec/scan/scanner_context.cpp +++ b/be/src/exec/scan/scanner_context.cpp @@ -735,11 +735,6 @@ std::shared_ptr ScannerContext::_pull_next_scan_task( } if (!_pending_tasks.empty()) { - // Skip submitting more pending scanners once the LIMIT budget is - // exhausted; they would only open and immediately EOF. - if (_shared_scan_limit->load(std::memory_order_acquire) == 0) { - return nullptr; - } std::shared_ptr next_scan_task; next_scan_task = _pending_tasks.top(); _pending_tasks.pop(); diff --git a/regression-test/data/correctness_p0/test_shared_scan_limit_pending_tasks.out b/regression-test/data/correctness_p0/test_shared_scan_limit_pending_tasks.out new file mode 100644 index 00000000000000..ae79c1d44f40e2 --- /dev/null +++ b/regression-test/data/correctness_p0/test_shared_scan_limit_pending_tasks.out @@ -0,0 +1,7 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !shared_scan_limit_parallel_2 -- + +-- !shared_scan_limit_parallel_6 -- + +-- !shared_scan_limit_parallel_8 -- + diff --git a/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy b/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy new file mode 100644 index 00000000000000..6b5f49469dddb7 --- /dev/null +++ b/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy @@ -0,0 +1,175 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +suite("test_shared_scan_limit_pending_tasks") { + sql "DROP TABLE IF EXISTS shared_scan_limit_t1" + sql "DROP TABLE IF EXISTS shared_scan_limit_t3" + sql "DROP TABLE IF EXISTS shared_scan_limit_t4" + sql "DROP TABLE IF EXISTS shared_scan_limit_t5" + + sql """ + CREATE TABLE shared_scan_limit_t1 ( + pk INT, + v1 BIGINT, + v2 BIGINT + ) ENGINE=OLAP + DUPLICATE KEY(pk) + DISTRIBUTED BY HASH(pk) BUCKETS 10 + PROPERTIES ("replication_num" = "1") + """ + sql """ + CREATE TABLE shared_scan_limit_t3 ( + pk INT, + v1 BIGINT, + v2 BIGINT + ) ENGINE=OLAP + DUPLICATE KEY(pk) + DISTRIBUTED BY HASH(pk) BUCKETS 10 + PROPERTIES ("replication_num" = "1") + """ + sql """ + CREATE TABLE shared_scan_limit_t4 ( + pk INT, + v1 BIGINT, + v2 BIGINT + ) ENGINE=OLAP + DUPLICATE KEY(pk) + DISTRIBUTED BY HASH(pk) BUCKETS 10 + PROPERTIES ("replication_num" = "1") + """ + sql """ + CREATE TABLE shared_scan_limit_t5 ( + pk INT, + v1 BIGINT, + v2 BIGINT + ) ENGINE=OLAP + DUPLICATE KEY(pk) + DISTRIBUTED BY HASH(pk) BUCKETS 10 + PROPERTIES ("replication_num" = "1") + """ + + sql """ + INSERT INTO shared_scan_limit_t1 VALUES + (0,70,32163), (1,-71,NULL), (2,-85,-222125), (3,26,22867), + (4,17,NULL), (5,NULL,6236439), (6,-16907,-24720), (7,26877,1893), + (8,30405,9), (9,5873343,66), (10,-17,-47), (11,NULL,20987), + (12,NULL,-4672), (13,NULL,3551614), (14,28521,-1372), (15,NULL,NULL), + (16,38,NULL), (17,NULL,-54), (18,NULL,-4686447), (19,NULL,-8420) + """ + sql """ + INSERT INTO shared_scan_limit_t3 VALUES + (0,28862,-25225), (1,5516858,-4609390), (2,-7815300,NULL), + (3,NULL,-7685824), (4,22293,26373), (5,NULL,NULL), (6,1237,-31), + (7,1606,3318374), (8,NULL,-21355), (9,NULL,-20127), (10,5413995,-3), + (11,-4718995,NULL), (12,3179854,3733421), (13,18867,76), + (14,15517,-4932405), (15,-30778,NULL), (16,NULL,-7258), + (17,NULL,7520313), (18,6936629,NULL), (19,-45,6231596) + """ + sql """ + INSERT INTO shared_scan_limit_t4 VALUES + (0,5649134,4812327), (1,8249417,4670054), (2,106,6714155), + (3,NULL,NULL), (4,-52850,-5113499), (5,-23911,-9637), (6,-26168,9), + (7,-47,-28398), (8,-29518,-2365317), (9,20685,22956), (10,97,25099), + (11,-32617,6143808), (12,NULL,NULL), (13,-52,NULL), (14,7680925,NULL), + (15,10848,NULL), (16,64,-3394), (17,113,-12488), (18,-87,-12093), + (19,NULL,22) + """ + sql """ + INSERT INTO shared_scan_limit_t5 VALUES + (0,-3018453,-5763927), (1,-60,-6233027), (2,78,-5304), (3,63,NULL), + (4,23287,NULL), (5,98,7989209), (6,-61,-30493), (7,-6781665,-22321), + (8,-5165806,-75), (9,NULL,31290), (10,1311226,7754904), (11,-72,54), + (12,114,3741372), (13,NULL,NULL), (14,-40,2014555), (15,69,-4555977), + (16,118,-8708), (17,NULL,NULL), (18,-123,NULL), (19,NULL,19843) + """ + sql "sync" + + sql "SET enable_sql_cache = false" + sql "SET experimental_enable_local_shuffle = false" + sql "SET query_timeout = 5" + + // A shared LIMIT can be exhausted while another ScannerContext still owns pending tasks. + // Every context must still reach EOS instead of leaving its pipeline task blocked forever. + sql "SET parallel_pipeline_task_num = 2" + qt_shared_scan_limit_parallel_2 """ + SELECT v1, pk + FROM shared_scan_limit_t1 AS t1 + WHERE v2 > ( + SELECT SUM(v1) + FROM shared_scan_limit_t5 AS t3 + WHERE EXISTS ( + SELECT pk + FROM shared_scan_limit_t3 AS t4 + WHERE NOT EXISTS ( + SELECT pk, v1, v1, v2, v1 + FROM shared_scan_limit_t4 AS t5 + WHERE v2 NOT IN ( + SELECT pk FROM shared_scan_limit_t1 AS t6 ORDER BY pk + ) + ) AND v1 < 7 + ) + ) + ORDER BY v1, pk DESC + LIMIT 6 + """ + + sql "SET parallel_pipeline_task_num = 6" + qt_shared_scan_limit_parallel_6 """ + SELECT v1, pk + FROM shared_scan_limit_t1 AS t1 + WHERE v2 > ( + SELECT SUM(v1) + FROM shared_scan_limit_t5 AS t3 + WHERE EXISTS ( + SELECT pk + FROM shared_scan_limit_t3 AS t4 + WHERE NOT EXISTS ( + SELECT pk, v1, v1, v2, v1 + FROM shared_scan_limit_t4 AS t5 + WHERE v2 NOT IN ( + SELECT pk FROM shared_scan_limit_t1 AS t6 ORDER BY pk + ) + ) AND v1 < 7 + ) + ) + ORDER BY v1, pk DESC + LIMIT 6 + """ + + sql "SET parallel_pipeline_task_num = 8" + qt_shared_scan_limit_parallel_8 """ + SELECT v1, pk + FROM shared_scan_limit_t1 AS t1 + WHERE v2 > ( + SELECT SUM(v1) + FROM shared_scan_limit_t5 AS t3 + WHERE EXISTS ( + SELECT pk + FROM shared_scan_limit_t3 AS t4 + WHERE NOT EXISTS ( + SELECT pk, v1, v1, v2, v1 + FROM shared_scan_limit_t4 AS t5 + WHERE v2 NOT IN ( + SELECT pk FROM shared_scan_limit_t1 AS t6 ORDER BY pk + ) + ) AND v1 < 7 + ) + ) + ORDER BY v1, pk DESC + LIMIT 6 + """ +} From 0d1809173d0ffeedc74d33616b3fd887d945c4fe Mon Sep 17 00:00:00 2001 From: Mryange Date: Mon, 20 Jul 2026 20:14:41 +0800 Subject: [PATCH 2/2] [fix](be) Handle exhausted shared scan limits ### What problem does this PR solve? Issue Number: None Related PR: #62222 Problem Summary: Shared scan-limit completion only recognized an exact remaining value of zero. Concurrent scanners can overshoot the finite budget below zero, while -1 still represents a query without LIMIT. Add a single finite-exhaustion predicate for ScannerContext completion and increase the regression timeout to avoid load-sensitive failures. ### Release note Fix ScannerContext completion when a finite shared scan limit becomes negative. ### Check List (For Author) - Test: No new test run for this follow-up - Previous regression and manual reproduction passed before the timeout-only adjustment - Behavior changed: Yes, negative exhausted shared limits now complete the scanner context - Does this need documentation: No --- be/src/exec/scan/scanner_context.cpp | 6 +++++- be/src/exec/scan/scanner_context.h | 1 + .../test_shared_scan_limit_pending_tasks.groovy | 2 +- 3 files changed, 7 insertions(+), 2 deletions(-) diff --git a/be/src/exec/scan/scanner_context.cpp b/be/src/exec/scan/scanner_context.cpp index abd3dce06f31cf..42d379bfe49424 100644 --- a/be/src/exec/scan/scanner_context.cpp +++ b/be/src/exec/scan/scanner_context.cpp @@ -404,7 +404,7 @@ Status ScannerContext::get_block_from_queue(RuntimeState* state, Block* block, b if (_completed_tasks.empty() && (_num_finished_scanners == _all_scanners.size() || - (_shared_scan_limit->load(std::memory_order_acquire) == 0 && _in_flight_tasks_num == 0))) { + (_is_shared_scan_limit_exhausted() && _in_flight_tasks_num == 0))) { _set_scanner_done(); _is_finished = true; } @@ -545,6 +545,10 @@ void ScannerContext::_set_scanner_done() { _dependency->set_always_ready(); } +bool ScannerContext::_is_shared_scan_limit_exhausted() const { + return limit >= 0 && _shared_scan_limit->load(std::memory_order_acquire) <= 0; +} + void ScannerContext::update_peak_running_scanner(int num) { #ifndef BE_TEST _local_state->_peak_running_scanner->add(num); diff --git a/be/src/exec/scan/scanner_context.h b/be/src/exec/scan/scanner_context.h index 0f97e82126d87b..013deef35443a2 100644 --- a/be/src/exec/scan/scanner_context.h +++ b/be/src/exec/scan/scanner_context.h @@ -267,6 +267,7 @@ class ScannerContext : public std::enable_shared_from_this, /// 3. `_free_blocks_memory_usage` < `_max_bytes_in_queue`, remains enough memory to scale up /// 4. At most scale up `MAX_SCALE_UP_RATIO` times to `_max_thread_num` void _set_scanner_done(); + bool _is_shared_scan_limit_exhausted() const; RuntimeState* _state = nullptr; ScanLocalStateBase* _local_state = nullptr; diff --git a/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy b/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy index 6b5f49469dddb7..72f47c16f66444 100644 --- a/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy +++ b/regression-test/suites/correctness_p0/test_shared_scan_limit_pending_tasks.groovy @@ -100,7 +100,7 @@ suite("test_shared_scan_limit_pending_tasks") { sql "SET enable_sql_cache = false" sql "SET experimental_enable_local_shuffle = false" - sql "SET query_timeout = 5" + sql "SET query_timeout = 60" // A shared LIMIT can be exhausted while another ScannerContext still owns pending tasks. // Every context must still reach EOS instead of leaving its pipeline task blocked forever.