This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 765185a5ec5 branch-4.1: [fix](arrow-flight) Release a finished Flight
query's coordinator instead of holding it until wait_timeout #67504 (#67558)
765185a5ec5 is described below
commit 765185a5ec5034808d9a10413f027757d9705eb2
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Thu Sep 10 07:48:18 2026 +0800
branch-4.1: [fix](arrow-flight) Release a finished Flight query's
coordinator instead of holding it until wait_timeout #67504 (#67558)
Cherry-picked from #67504
### Backport notes
Two conflicts, both resolved without changing behaviour:
- `Config`: branch-4.1 still has the bilingual
`autobucket_partition_size_per_bucket_gb` description
that master rewrote, so that entry is kept as it is here and the new
`arrow_flight_deferred_query_idle_timeout_second` is inserted above it,
verbatim.
- `StmtExecutor`: the comment above the deferral gate is taken from
upstream verbatim, so it already
mentions #67368. That fix is being backported separately (see below);
the gate itself does not
depend on it.
One adaptation: `SplitAssignment` is
`org.apache.doris.datasource.SplitAssignment` on this branch, not
`org.apache.doris.datasource.split.SplitAssignment`, so
`ArrowFlightDeferralGateTest` imports it from
its package here.
Everything else -- `ScanNode.hasBatchSplitSource`,
`Coordinator.hasBatchSplitSource`, the
`ConnectContext` idle reaper, `FlightSqlConnectProcessor`,
`StmtExecutorTest`,
`FlightSqlDeferredQueryIdleTimeoutTest` and both regression suites --
applied cleanly.
### Merge order
Touches the same `StmtExecutor` comment as #67557 (the branch-4.1
backport of #67487). Merging that one first keeps this
one conflict-free.
### Local verification
`./build.sh --fe` on this branch: **BUILD SUCCESS** (2:26 min), no
errors, checkstyle clean on every
module -- `fe-common` (the new `Config` entry) and `fe-core` (the five
main-source changes plus
`ArrowFlightDeferralGateTest`, `FlightSqlDeferredQueryIdleTimeoutTest`
and the new `StmtExecutorTest`
case) included. The regression suites were not run locally.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
https://claude.ai/code/session_01DzUYFcGHQH3bnLGCjpncVj
Co-authored-by: Raghvendra Singh <[email protected]>
---
.../main/java/org/apache/doris/common/Config.java | 11 ++
.../java/org/apache/doris/planner/ScanNode.java | 11 ++
.../java/org/apache/doris/qe/ConnectContext.java | 45 +++++-
.../main/java/org/apache/doris/qe/Coordinator.java | 19 +++
.../java/org/apache/doris/qe/StmtExecutor.java | 51 ++++---
.../arrowflight/FlightSqlConnectProcessor.java | 11 +-
.../doris/qe/ArrowFlightDeferralGateTest.java | 69 +++++++++
.../java/org/apache/doris/qe/StmtExecutorTest.java | 33 +++++
.../FlightSqlDeferredQueryIdleTimeoutTest.java | 156 +++++++++++++++++++++
.../test_arrow_flight_query_release.groovy | 88 ++++++++++++
.../test_iceberg_arrow_flight_split_source.groovy | 73 ++++++++++
11 files changed, 544 insertions(+), 23 deletions(-)
diff --git a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
index 2adee1ec6ed..f006447a714 100644
--- a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
+++ b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
@@ -3020,6 +3020,17 @@ public class Config extends ConfigBase {
"Maximal number of connections of Arrow Flight Server per FE."})
public static int arrow_flight_max_connections = 4096;
+ @ConfField(mutable = true, description = "Arrow Flight SQL only. A query
that scans an external table in "
+ + "batch mode keeps its FE coordinator alive after GetFlightInfo,
so the BE can keep fetching splits "
+ + "while the client pulls the results (DoGet); that coordinator is
normally released when the "
+ + "session runs its next query or is closed. Most Flight clients
never close a session, so the "
+ + "coordinator, and with it the query's workload group queue slot
and its active_queries entry, "
+ + "would otherwise stay held until wait_timeout. If the session
stays idle for longer than this "
+ + "many seconds after the query started, the coordinator is
released anyway. The bound is never "
+ + "shorter than the query's own execution timeout, and the session
itself is not killed "
+ + "(wait_timeout still governs that). 0 disables the bound.")
+ public static int arrow_flight_deferred_query_idle_timeout_second = 3600;
+
@ConfField(mutable = true, masterOnly = true, description = {
"Auto Buckets 中按照 partition size 去估算 bucket 数,存算一体 partition size 5G
估算一个 bucket,"
+ "但存算分离下 partition size 10G 估算一个 bucket。若配置小于
0,会在在代码中会自适应存算一体模式默认 5G,在存算分离默认 10G",
diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
index efa50779d64..b3def9bc8c6 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/ScanNode.java
@@ -140,6 +140,17 @@ public abstract class ScanNode extends PlanNode implements
SplitGenerator {
return desc;
}
+ /**
+ * Whether this scan hands out its splits lazily through a batch {@link
SplitSource} that the
+ * BE fetches from the FE while it is scanning (external-table batch mode,
see
+ * {@link SplitGenerator#isBatchMode()}). Such a scan needs its
coordinator alive until the BE
+ * has finished scanning, even after the FE is done dispatching the query:
closing the
+ * coordinator releases the split source ({@link #stop()}) and the BE's
next split fetch fails.
+ */
+ public boolean hasBatchSplitSource() {
+ return splitAssignment != null;
+ }
+
protected abstract void createScanRangeLocations() throws UserException;
/**
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
index 6cec50d4ee5..1d4aac17bd9 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
@@ -1043,7 +1043,9 @@ public class ConnectContext {
// held by the coordinator's scan nodes), so closing the coordinator at
the end of
// GetFlightInfo would release the SplitSource too early and make the BE's
fetchSplitBatch fail
// with "Split source X is released". These executors are finalized when
the next query starts
- // on this connection, or when the connection is torn down. See #62259.
+ // on this connection, when the connection is torn down, or by the idle
reaper in checkTimeout
+ // once the connection has been sleeping for
arrow_flight_deferred_query_idle_timeout_second.
+ // See #62259 and #67503.
private final List<StmtExecutor> flightSqlDeferredExecutors = new
ArrayList<>();
public void addFlightSqlDeferredExecutor(StmtExecutor executor) {
@@ -1070,6 +1072,45 @@ public class ConnectContext {
}
}
+ /**
+ * How long, in seconds, a sleeping connection may keep its deferred Arrow
Flight executors
+ * before the timeout checker finalizes them without killing the connection
+ * (Config.arrow_flight_deferred_query_idle_timeout_second). A Flight
client that opens a
+ * session per query and never closes it would otherwise pin each deferred
query's query queue
+ * slot and query registration until wait_timeout (8h by default). The
bound is never shorter
+ * than the execution timeout the deferred query was run with: the client
may still be pulling
+ * that query's results from the BE, which still needs the batch split
source the coordinator
+ * holds. Returns -1 when the bound is disabled or nothing is deferred.
+ */
+ public long getFlightSqlDeferredExecutorsIdleTimeoutS() {
+ int configTimeoutS =
Config.arrow_flight_deferred_query_idle_timeout_second;
+ if (configTimeoutS <= 0) {
+ return -1;
+ }
+ long execTimeoutS = -1;
+ synchronized (flightSqlDeferredExecutors) {
+ if (flightSqlDeferredExecutors.isEmpty()) {
+ return -1;
+ }
+ for (StmtExecutor deferredExecutor : flightSqlDeferredExecutors) {
+ execTimeoutS = Math.max(execTimeoutS,
deferredExecutor.getDeferredExecTimeoutS());
+ }
+ }
+ return Math.max(configTimeoutS, execTimeoutS);
+ }
+
+ // Called by the timeout checker for a sleeping connection that is not
past wait_timeout yet.
+ private void reapIdleFlightSqlDeferredExecutors(long idleMs) {
+ long timeoutS = getFlightSqlDeferredExecutorsIdleTimeoutS();
+ if (timeoutS < 0 || idleMs <= timeoutS * 1000L) {
+ return;
+ }
+ LOG.warn("release deferred arrow flight query of idle connection,
connectionId: {}, remote: {}, "
+ + "idle: {}ms, idle timeout: {}s",
+ connectionId, getRemoteHostPortString(), idleMs, timeoutS);
+ closeFlightSqlDeferredExecutors();
+ }
+
/**
* This method is idempotent.
*/
@@ -1342,6 +1383,8 @@ public class ConnectContext {
// Need kill this connection.
killFlag = true;
killConnection = true;
+ } else {
+ reapIdleFlightSqlDeferredExecutors(delta);
}
} else {
String timeoutTag = "query";
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
index 2db61ffc92a..c582dccf7f9 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/Coordinator.java
@@ -816,6 +816,25 @@ public class Coordinator implements CoordInterface {
execInternal();
}
+ /**
+ * Whether the BE keeps calling back into this coordinator after {@link
#exec()} returned: an
+ * external-table scan in batch mode fetches its splits lazily from the
split source that its
+ * scan node holds, so the coordinator must not be closed until the BE has
finished scanning.
+ * Arrow Flight SQL uses this to decide whether a query's coordinator has
to outlive
+ * GetFlightInfo, the client pulling the results from the BE later in
DoGet. See #62259.
+ */
+ public boolean hasBatchSplitSource() {
+ if (scanNodes == null) {
+ return false;
+ }
+ for (ScanNode scanNode : scanNodes) {
+ if (scanNode.hasBatchSplitSource()) {
+ return true;
+ }
+ }
+ return false;
+ }
+
@Override
public void close() {
// NOTE: all close method should be no exception
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index b1a80513b43..1639a30f562 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -194,6 +194,10 @@ public class StmtExecutor {
// is finalized later by ConnectContext (see #62259), so the eager close
in executeAndSendResult
// is skipped.
private volatile boolean deferredForArrowFlight = false;
+ // The execution timeout in effect when the coordinator was deferred.
Captured at that moment
+ // because per-statement SET_VAR values are reverted at the end of
execute(), so reading
+ // ConnectContext.getExecTimeoutS() later would report the session value
instead.
+ private volatile int deferredExecTimeoutS = -1;
private MasterOpExecutor masterOpExecutor = null;
private RedirectStatus redirectStatus = null;
private Planner planner;
@@ -980,6 +984,21 @@ public class StmtExecutor {
return deferredForArrowFlight;
}
+ // Execution timeout (seconds) the deferred query was run with; -1 when
the query is not deferred.
+ public int getDeferredExecTimeoutS() {
+ return deferredExecTimeoutS;
+ }
+
+ // Keep this query's coordinator alive past GetFlightInfo (see the gate in
executeAndSendResult)
+ // and hand it to the ConnectContext, which finalizes it later. Records
the execution timeout in
+ // effect right now: it floors the idle reaper's bound and must be the
value the query actually
+ // ran with, not the session value left behind after SET_VAR hints are
reverted.
+ void deferForArrowFlight() {
+ deferredForArrowFlight = true;
+ deferredExecTimeoutS = context.getExecTimeoutS();
+ context.addFlightSqlDeferredExecutor(this);
+ }
+
// Finalize an Arrow Flight query whose coordinator was kept alive across
the
// GetFlightInfo -> DoGet phases: close the coordinator (releasing
external-table batch
// SplitSources and the query queue slot) and then unregister the query.
See #62259.
@@ -1435,23 +1454,21 @@ public class StmtExecutor {
if (context.getConnectType().equals(ConnectType.ARROW_FLIGHT_SQL))
{
Preconditions.checkState(!context.isReturnResultFromLocal());
profile.getSummaryProfile().setTempStartTime();
- // Defer closing the coordinator to ConnectContext (closed on
the next query or
- // connection teardown) instead of in the finally block below.
This gate covers
- // every Arrow Flight query whose results are produced on the
BE (coordBase ==
- // coord) -- internal-table and external, batch or not. It is
REQUIRED only for an
- // external-table scan in batch mode, where the BE lazily
fetches splits from the FE
- // during the later DoGet phase, so closing the coordinator
here would release its
- // batch SplitSource too early and break DoGet. Other
remote-result queries do not
- // need deferral (the BE buffers their result independently)
but are captured by the
- // same gate; the trade-off is their coordinator, query queue
slot and query
- // registration stay held until the next query / teardown
instead of being released
- // at the end of GetFlightInfo. A short-circuit point query is
the one case with a
- // different coordBase, and it can no longer reach here: it
has no Arrow result on
- // either side, so LogicalResultSinkToShortCircuitPointQuery
keeps Arrow Flight SQL
- // on the normal execution path. See #62259 and #67368.
- if (coordBase == coord) {
- deferredForArrowFlight = true;
- context.addFlightSqlDeferredExecutor(this);
+ // The client pulls the results from the BE later (DoGet).
Only an external-table
+ // scan in batch mode still needs the coordinator after this
point: the BE fetches
+ // its splits lazily from the split source the coordinator
holds, so closing the
+ // coordinator here would release that source too early and
break DoGet (#62259).
+ // Such a coordinator is closed later by ConnectContext: on
the session's next
+ // query, on teardown, or by the idle reaper in checkTimeout.
The trade-off is that
+ // its query queue slot and query registration stay held until
then. Every other
+ // query closes its coordinator in the finally block below and
releases both right
+ // away, the BE buffering its results independently of the
coordinator (#67503).
+ // A short-circuit point query is the one case with a
different coordBase, and it
+ // can no longer reach here: it has no Arrow result on either
side, so
+ // LogicalResultSinkToShortCircuitPointQuery keeps Arrow
Flight SQL on the normal
+ // execution path (#67368).
+ if (coordBase == coord && coord.hasBatchSplitSource()) {
+ deferForArrowFlight();
}
return;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
index 24e30c3942f..2296986c535 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
@@ -196,11 +196,12 @@ public class FlightSqlConnectProcessor extends
ConnectProcessor implements AutoC
@Override
public void close() throws Exception {
ctx.setCommand(MysqlCommand.COM_SLEEP);
- // Executors whose results are pulled from the BE keep their
coordinator alive past
- // GetFlightInfo (registered as deferred executors on the
ConnectContext) so the BE can
- // still fetch external-table splits during DoGet. Do NOT finalize
those here; they are
- // finalized when the next query starts or the connection is torn
down. Executors that are
- // not deferred (local results, or a query that already failed) are
finalized now. See #62259.
+ // An external-table scan in batch mode keeps its coordinator alive
past GetFlightInfo
+ // (registered as a deferred executor on the ConnectContext) so the BE
can still fetch its
+ // splits during DoGet. Do NOT finalize those here; they are finalized
when the next query
+ // starts, when the connection is torn down, or by the idle reaper in
+ // ConnectContext.checkTimeout. Every other executor (local results,
results the BE buffers
+ // on its own, or a query that already failed) is finalized now. See
#62259 and #67503.
for (StmtExecutor asynExecutor : returnResultFromRemoteExecutor) {
if (!asynExecutor.isDeferredForArrowFlight()) {
asynExecutor.finalizeQuery();
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/qe/ArrowFlightDeferralGateTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/ArrowFlightDeferralGateTest.java
new file mode 100644
index 00000000000..efd6730fa99
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/qe/ArrowFlightDeferralGateTest.java
@@ -0,0 +1,69 @@
+// 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.
+
+package org.apache.doris.qe;
+
+import org.apache.doris.analysis.DescriptorTable;
+import org.apache.doris.datasource.SplitAssignment;
+import org.apache.doris.planner.PlanFragment;
+import org.apache.doris.planner.ScanNode;
+import org.apache.doris.thrift.TUniqueId;
+
+import com.google.common.collect.Lists;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.lang.reflect.Field;
+import java.util.List;
+
+/**
+ * The predicate behind the Arrow Flight deferral gate in
StmtExecutor.executeAndSendResult (#67503):
+ * a coordinator has to outlive GetFlightInfo only when one of its scans still
hands out splits to
+ * the BE lazily, i.e. an external-table scan in batch mode holding a batch
split source (#62259).
+ */
+public class ArrowFlightDeferralGateTest {
+
+ private static ScanNode scanNode(boolean batchSplitSource) throws
Exception {
+ ScanNode node = Mockito.mock(ScanNode.class,
Mockito.CALLS_REAL_METHODS);
+ if (batchSplitSource) {
+ // FileQueryScanNode.createScanRangeLocations sets this only in
batch mode.
+ Field field = ScanNode.class.getDeclaredField("splitAssignment");
+ field.setAccessible(true);
+ field.set(node, Mockito.mock(SplitAssignment.class));
+ }
+ return node;
+ }
+
+ private static Coordinator coordinator(List<ScanNode> scanNodes) {
+ return new Coordinator(1L, new TUniqueId(1L, 2L), new
DescriptorTable(), Lists.<PlanFragment>newArrayList(),
+ scanNodes, "UTC", false, false);
+ }
+
+ @Test
+ public void
testScanNodeHasBatchSplitSourceOnlyWhenSplitsAreHandedOutLazily() throws
Exception {
+ Assertions.assertFalse(scanNode(false).hasBatchSplitSource());
+ Assertions.assertTrue(scanNode(true).hasBatchSplitSource());
+ }
+
+ @Test
+ public void testCoordinatorHasBatchSplitSourceIfAnyScanDoes() throws
Exception {
+
Assertions.assertFalse(coordinator(Lists.newArrayList()).hasBatchSplitSource());
+ Assertions.assertFalse(coordinator(Lists.newArrayList(scanNode(false),
scanNode(false))).hasBatchSplitSource());
+ Assertions.assertTrue(coordinator(Lists.newArrayList(scanNode(false),
scanNode(true))).hasBatchSplitSource());
+ }
+}
diff --git a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
index 607cea3b40a..76d8e0be6bf 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
@@ -73,6 +73,39 @@ public class StmtExecutorTest extends TestWithFeService {
Assert.assertEquals(QueryState.MysqlStateType.OK,
connectContext.getState().getStateType());
}
+ // The deferral gate (#67503): a coordinator is kept alive past
GetFlightInfo only when the BE
+ // still fetches splits from it (Coordinator.hasBatchSplitSource), and the
execution timeout it
+ // ran with is frozen at that moment. SET_VAR hint values are reverted
when execute() ends, so
+ // the idle reaper must not read the session value later.
+ @Test
+ public void testDeferForArrowFlightFreezesExecTimeoutInEffect() throws
Exception {
+ int savedQueryTimeout =
connectContext.getSessionVariable().getQueryTimeoutS();
+ int savedIdleTimeout =
Config.arrow_flight_deferred_query_idle_timeout_second;
+ connectContext.setQueryId(new TUniqueId(0x67503L, 0x1L));
+ try {
+ Config.arrow_flight_deferred_query_idle_timeout_second = 1;
+ connectContext.getSessionVariable().setQueryTimeoutS(1234);
+ StmtExecutor stmtExecutor = new StmtExecutor(connectContext, "");
+ Assertions.assertFalse(stmtExecutor.isDeferredForArrowFlight());
+ Assertions.assertEquals(-1,
stmtExecutor.getDeferredExecTimeoutS());
+
+ stmtExecutor.deferForArrowFlight();
+
+ Assertions.assertTrue(stmtExecutor.isDeferredForArrowFlight());
+ Assertions.assertEquals(1234,
stmtExecutor.getDeferredExecTimeoutS());
+ // the reaper's bound is floored at the frozen value ...
+ Assertions.assertEquals(1234L,
connectContext.getFlightSqlDeferredExecutorsIdleTimeoutS());
+ // ... even after the session value moved on, as it does when a
SET_VAR hint is reverted
+ connectContext.getSessionVariable().setQueryTimeoutS(5);
+ Assertions.assertEquals(1234,
stmtExecutor.getDeferredExecTimeoutS());
+ Assertions.assertEquals(1234L,
connectContext.getFlightSqlDeferredExecutorsIdleTimeoutS());
+ } finally {
+ connectContext.closeFlightSqlDeferredExecutors();
+
connectContext.getSessionVariable().setQueryTimeoutS(savedQueryTimeout);
+ Config.arrow_flight_deferred_query_idle_timeout_second =
savedIdleTimeout;
+ }
+ }
+
// Arrow Flight SQL keeps a query's coordinator alive across GetFlightInfo
-> DoGet (see #62259);
// it is released later by finalizeArrowFlightQuery(), which closes the
coordinator and then
// unregisters the query. The close and the unregister must be
independent: if coord.close()
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlDeferredQueryIdleTimeoutTest.java
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlDeferredQueryIdleTimeoutTest.java
new file mode 100644
index 00000000000..211f2ab9849
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlDeferredQueryIdleTimeoutTest.java
@@ -0,0 +1,156 @@
+// 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.
+
+package org.apache.doris.service.arrowflight.sessions;
+
+import org.apache.doris.common.Config;
+import org.apache.doris.common.FeConstants;
+import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+/**
+ * The idle reaper for deferred Arrow Flight queries (#67503). A sleeping
Flight session whose last
+ * query kept its coordinator alive (an external-table scan in batch mode, see
#62259) gets that
+ * coordinator finalized by the connection timeout checker once the session
has been idle for
+ * arrow_flight_deferred_query_idle_timeout_second, floored at the execution
timeout the query ran
+ * with. The session itself is not killed, wait_timeout still governs that,
and a MySQL session is
+ * untouched.
+ */
+public class FlightSqlDeferredQueryIdleTimeoutTest {
+ private int savedIdleTimeout;
+ private boolean savedRunningUnitTest;
+
+ @BeforeEach
+ public void setUp() {
+ savedIdleTimeout =
Config.arrow_flight_deferred_query_idle_timeout_second;
+ savedRunningUnitTest = FeConstants.runningUnitTest;
+ // ConnectContext.init() registers the session with Env unless running
as a unit test.
+ FeConstants.runningUnitTest = true;
+ }
+
+ @AfterEach
+ public void tearDown() {
+ Config.arrow_flight_deferred_query_idle_timeout_second =
savedIdleTimeout;
+ FeConstants.runningUnitTest = savedRunningUnitTest;
+ }
+
+ private static StmtExecutor deferredExecutor(int execTimeoutS) {
+ StmtExecutor executor = Mockito.mock(StmtExecutor.class);
+
Mockito.when(executor.getDeferredExecTimeoutS()).thenReturn(execTimeoutS);
+ return executor;
+ }
+
+ // A Flight session that ran a query and has been sleeping since; the
client never closed it.
+ private static FlightSqlConnectContext
sleepingFlightSession(StmtExecutor... deferred) {
+ FlightSqlConnectContext ctx = new
FlightSqlConnectContext("test-peer-identity");
+ ctx.setCommand(MysqlCommand.COM_SLEEP);
+ ctx.setStartTime();
+ for (StmtExecutor executor : deferred) {
+ ctx.addFlightSqlDeferredExecutor(executor);
+ }
+ return ctx;
+ }
+
+ @Test
+ public void testIdleSessionReleasesDeferredQueryButIsNotKilled() {
+ Config.arrow_flight_deferred_query_idle_timeout_second = 7;
+ StmtExecutor deferred = deferredExecutor(5);
+ FlightSqlConnectContext ctx = sleepingFlightSession(deferred);
+ long start = ctx.getStartTime();
+ Assertions.assertEquals(7L,
ctx.getFlightSqlDeferredExecutorsIdleTimeoutS());
+
+ // not idle for long enough yet
+ ctx.checkTimeout(start + 7_000L);
+ Mockito.verify(deferred, Mockito.never()).finalizeArrowFlightQuery();
+ Assertions.assertFalse(ctx.isKilled());
+
+ // past the bound: the deferred coordinator is finalized and the
session survives
+ ctx.checkTimeout(start + 7_001L);
+ Mockito.verify(deferred).finalizeArrowFlightQuery();
+ Assertions.assertFalse(ctx.isKilled());
+ Assertions.assertEquals(-1L,
ctx.getFlightSqlDeferredExecutorsIdleTimeoutS());
+
+ // a later tick has nothing left to release
+ ctx.checkTimeout(start + 60_000L);
+ Mockito.verify(deferred, Mockito.times(1)).finalizeArrowFlightQuery();
+ Assertions.assertFalse(ctx.isKilled());
+ }
+
+ @Test
+ public void testBoundIsFlooredAtTheExecTimeoutTheDeferredQueryRanWith() {
+ Config.arrow_flight_deferred_query_idle_timeout_second = 3;
+ StmtExecutor shortQuery = deferredExecutor(5);
+ StmtExecutor longQuery = deferredExecutor(20);
+ FlightSqlConnectContext ctx = sleepingFlightSession(shortQuery,
longQuery);
+ long start = ctx.getStartTime();
+ // the longest deferred query wins: a client may still be pulling its
results from the BE
+ Assertions.assertEquals(20L,
ctx.getFlightSqlDeferredExecutorsIdleTimeoutS());
+
+ ctx.checkTimeout(start + 19_999L);
+ Mockito.verify(shortQuery, Mockito.never()).finalizeArrowFlightQuery();
+ Mockito.verify(longQuery, Mockito.never()).finalizeArrowFlightQuery();
+
+ ctx.checkTimeout(start + 20_001L);
+ Mockito.verify(shortQuery).finalizeArrowFlightQuery();
+ Mockito.verify(longQuery).finalizeArrowFlightQuery();
+ Assertions.assertFalse(ctx.isKilled());
+ }
+
+ @Test
+ public void testZeroDisablesTheReaper() {
+ Config.arrow_flight_deferred_query_idle_timeout_second = 0;
+ StmtExecutor deferred = deferredExecutor(5);
+ FlightSqlConnectContext ctx = sleepingFlightSession(deferred);
+ Assertions.assertEquals(-1L,
ctx.getFlightSqlDeferredExecutorsIdleTimeoutS());
+
+ // idle for almost the whole wait_timeout: nothing is released and the
session is alive
+ long waitTimeoutMs = ctx.getSessionVariable().getWaitTimeoutS() *
1000L;
+ ctx.checkTimeout(ctx.getStartTime() + waitTimeoutMs - 1);
+ Mockito.verify(deferred, Mockito.never()).finalizeArrowFlightQuery();
+ Assertions.assertFalse(ctx.isKilled());
+ }
+
+ @Test
+ public void testNothingDeferredMeansNoBound() {
+ Config.arrow_flight_deferred_query_idle_timeout_second = 7;
+ FlightSqlConnectContext ctx = sleepingFlightSession();
+ Assertions.assertEquals(-1L,
ctx.getFlightSqlDeferredExecutorsIdleTimeoutS());
+
+ ctx.checkTimeout(ctx.getStartTime() + 3_600_000L);
+ Assertions.assertFalse(ctx.isKilled());
+ }
+
+ @Test
+ public void testMysqlSessionIsUntouched() {
+ Config.arrow_flight_deferred_query_idle_timeout_second = 1;
+ ConnectContext ctx = new ConnectContext();
+ ctx.setCommand(MysqlCommand.COM_SLEEP);
+ ctx.setStartTime();
+ Assertions.assertEquals(-1L,
ctx.getFlightSqlDeferredExecutorsIdleTimeoutS());
+
+ // idle far beyond the Flight bound but within wait_timeout: still
alive
+ ctx.checkTimeout(ctx.getStartTime() + 3_600_000L);
+ Assertions.assertFalse(ctx.isKilled());
+ }
+}
diff --git
a/regression-test/suites/arrow_flight_sql_p0/test_arrow_flight_query_release.groovy
b/regression-test/suites/arrow_flight_sql_p0/test_arrow_flight_query_release.groovy
new file mode 100644
index 00000000000..0913164a1eb
--- /dev/null
+++
b/regression-test/suites/arrow_flight_sql_p0/test_arrow_flight_query_release.groovy
@@ -0,0 +1,88 @@
+// 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.
+
+// Regression for https://github.com/apache/doris/issues/67503
+//
+// Over Arrow Flight SQL a query runs in two phases: GetFlightInfo (plan and
start it on the BE)
+// and DoGet (the client pulls the results from the BE). Only an
external-table scan in batch mode
+// needs its FE coordinator after GetFlightInfo (#62259). Every other query
has to release its
+// coordinator, and with it the workload group queue slot and the
active_queries entry, at the end
+// of GetFlightInfo: most Flight clients never close their session, so a
coordinator that waited
+// for the session's next query kept one queue slot per finished query until
wait_timeout.
+//
+// The framework's Flight session behaves like such a client: it is reused
across statements and
+// never closed.
+suite("test_arrow_flight_query_release", "arrow_flight_sql") {
+ def tableName = "test_arrow_flight_query_release_tbl"
+ def wgName = "test_arrow_flight_query_release_wg"
+
+ def forComputeGroupStr = ""
+ if (isCloudMode()) {
+ def clusters = sql " SHOW CLUSTERS; "
+ assertTrue(!clusters.isEmpty())
+ forComputeGroupStr = " for ${clusters[0][0]} "
+ }
+
+ sql "DROP TABLE IF EXISTS ${tableName}"
+ sql """
+ CREATE TABLE ${tableName} (id int, name varchar(20))
+ DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 1
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql "INSERT INTO ${tableName} VALUES (1, 'a'), (2, 'b'), (3, 'c')"
+
+ sql "ADMIN SET FRONTEND CONFIG ('enable_workload_group' = 'true')"
+ sql "DROP WORKLOAD GROUP IF EXISTS ${wgName} ${forComputeGroupStr}"
+ // One running query at a time and no waiting queue: while a query still
holds the slot, the
+ // next scanning query in the group fails at once with "query waiting
queue is full".
+ sql """
+ CREATE WORKLOAD GROUP ${wgName} ${forComputeGroupStr}
+ PROPERTIES ('max_concurrency' = '1', 'max_queue_size' = '0',
'queue_timeout' = '0')
+ """
+ try {
+ // The Flight session is a session of its own, so it is bound to the
group separately.
+ sql "SET workload_group = '${wgName}'"
+ arrow_flight_sql "SET workload_group = '${wgName}'"
+
+ // A scanning query over Arrow Flight SQL. The session stays open
afterwards.
+ def flightRows = arrow_flight_sql "SELECT id, name FROM ${tableName}
ORDER BY id"
+ assertEquals(3, flightRows.size())
+
+ // Its coordinator was released at the end of GetFlightInfo, so the
query is gone from
+ // active_queries. The LIKE pattern is assembled with CONCAT so that
this statement's own
+ // text does not match it.
+ def registered = sql """
+ SELECT QUERY_ID, SQL FROM information_schema.active_queries
+ WHERE SQL LIKE CONCAT('%FROM ${tableName}', ' ORDER BY id%')
+ """
+ assertTrue(registered.isEmpty(), "finished Arrow Flight query is still
registered: ${registered}")
+
+ // ... and its queue slot is free again: a scanning query in the same
group runs instead
+ // of failing with "query waiting queue is full".
+ def mysqlRows = sql "SELECT id FROM ${tableName} ORDER BY id"
+ assertEquals(3, mysqlRows.size())
+ } finally {
+ sql "SET workload_group = 'normal'"
+ try {
+ arrow_flight_sql "SET workload_group = 'normal'"
+ } catch (Throwable ignore) {
+ // best effort: the Flight session must not keep pointing at the
dropped group
+ }
+ sql "DROP WORKLOAD GROUP IF EXISTS ${wgName} ${forComputeGroupStr}"
+ sql "DROP TABLE IF EXISTS ${tableName}"
+ }
+}
diff --git
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_arrow_flight_split_source.groovy
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_arrow_flight_split_source.groovy
index e721d0e9d8e..67b9a431b82 100644
---
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_arrow_flight_split_source.groovy
+++
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_arrow_flight_split_source.groovy
@@ -81,6 +81,12 @@ suite("test_iceberg_arrow_flight_split_source",
"p0,external") {
"s3.region" = "us-east-1"
);"""
+ // #67503: the idle reaper for a deferred batch-mode scan (see below). Set
the bound low, and
+ // restore the FE's original value afterwards.
+ int idleTimeoutS = 10
+ def origIdleTimeout = sql """ ADMIN SHOW FRONTEND CONFIG LIKE
'arrow_flight_deferred_query_idle_timeout_second' """
+ assert origIdleTimeout.size() == 1 :
"arrow_flight_deferred_query_idle_timeout_second not found in FE config"
+
Connection flightConn = null
try {
// Baseline over the MySQL protocol (works regardless of the bug).
@@ -120,7 +126,74 @@ suite("test_iceberg_arrow_flight_split_source",
"p0,external") {
// deferred coordinator when the next query starts.
def flightLimited = flightSql """ select * from ${table} limit 10 """
assert flightLimited.size() > 0 && flightLimited.size() <= 10 :
"unexpected row count: ${flightLimited.size()}"
+
+ // #67503, the other side of the deferral gate: the SAME external
table scanned WITHOUT
+ // batch mode is not deferred. Its coordinator, and with it the
query's workload group queue
+ // slot and its active_queries entry, is released at the end of
GetFlightInfo, before the
+ // client pulls anything. That is the case the gate actually moved, so
it needs its own
+ // coverage here: every other Flight query in this suite runs in batch
mode.
+ flightSql """ set enable_external_table_batch_mode = false """
+
+ // Negative control, mirroring the batch assertion above:
"(approximate)" is emitted only
+ // when isBatchMode(), so its absence proves this really is the
synchronous split path and
+ // the assertions below cannot silently pass on the batch path.
+ def explainNonBatch = flightSql """ explain select * from ${table} """
+ boolean stillBatch = explainNonBatch.any { row ->
+ row.any { cell -> cell != null &&
cell.toString().contains("approximate") }
+ }
+ assert !stillBatch : "expected the non-batch split path in the Arrow
Flight plan, got: ${explainNonBatch}"
+
+ // The scan must still be complete: the FE closed the coordinator at
the end of
+ // GetFlightInfo, and the BE buffers the result independently of it.
+ def flightNonBatch = flightSql """ select * from ${table} """
+ assertEquals(expectedRows, (flightNonBatch.size() as long))
+
+ // ... and the release really was eager, unlike the batch-mode scan
below. A distinct limit
+ // keeps this query's text apart from the other scans, and the LIKE
pattern is assembled
+ // with CONCAT so that the probe statement's own text does not match
it. No polling is
+ // needed: finalizeQuery() runs inside GetFlightInfo, so it has
already happened by the time
+ // the client has the rows.
+ def flightNonBatchLimited = flightSql """ select * from ${table} limit
17 """
+ assertEquals(17, flightNonBatchLimited.size())
+ def nonBatchRegistered = sql """ select QUERY_ID from
information_schema.active_queries
+ where SQL like CONCAT('%from ${table} limit', ' 17%') """
+ assert nonBatchRegistered.isEmpty() : "a non-batch Flight query must
release its coordinator at the end of GetFlightInfo, still registered:
${nonBatchRegistered}"
+
+ // Back to batch mode: the idle reaper below needs a deferred
coordinator to release.
+ flightSql """ set enable_external_table_batch_mode = true """
+
+ // #67503: a batch-mode scan keeps its coordinator (and with it the
query's workload group
+ // queue slot and its active_queries entry) alive after GetFlightInfo,
until the session
+ // runs its next query or is closed. A client that does neither would
hold them until
+ // wait_timeout, so the FE releases the coordinator once the session
has been idle for
+ // arrow_flight_deferred_query_idle_timeout_second, never before the
query's own execution
+ // timeout, and without killing the session.
+ sql """ ADMIN SET FRONTEND CONFIG
('arrow_flight_deferred_query_idle_timeout_second' = '${idleTimeoutS}') """
+ flightSql """ set query_timeout = ${idleTimeoutS} """
+ def flightReap = flightSql """ select * from ${table} limit 13 """
+ assertEquals(13, flightReap.size())
+
+ // The LIKE pattern is assembled with CONCAT so that this statement's
own text does not
+ // match it.
+ def deferredQuery = { ->
+ sql """ select QUERY_ID from information_schema.active_queries
+ where SQL like CONCAT('%from ${table} limit', ' 13%') """
+ }
+ // Right after the scan the query is still registered: its coordinator
is deferred.
+ assert deferredQuery().size() == 1 : "expected the batch-mode Flight
query to stay registered until the idle reaper releases it"
+
+ // Once the session has been idle for the bound, the reaper releases
it ...
+ long deadline = System.currentTimeMillis() + 60_000L
+ while (!deferredQuery().isEmpty() && System.currentTimeMillis() <
deadline) {
+ Thread.sleep(1000)
+ }
+ assert deferredQuery().isEmpty() : "the idle reaper did not release
the deferred Flight query within 60s"
+
+ // ... and the session survives: it still runs queries.
+ def afterReap = flightSql """ select * from ${table} limit 1 """
+ assertEquals(1, afterReap.size())
} finally {
+ sql """ ADMIN SET FRONTEND CONFIG
('arrow_flight_deferred_query_idle_timeout_second' =
'${origIdleTimeout[0][1]}') """
// Close our own connection (best effort) so a dead endpoint cannot
mask the real failure,
// then drop the catalog over the reliable MySQL connection.
if (flightConn != null) {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]