This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 23b5ee5700f [fix](arrow-flight) Release a finished Flight query's
coordinator instead of holding it until wait_timeout (#67504)
23b5ee5700f is described below
commit 23b5ee5700fdda55857fdf805e6c2a2ec50fde36
Author: Raghvendra Singh <[email protected]>
AuthorDate: Mon Sep 7 07:25:02 2026 +0530
[fix](arrow-flight) Release a finished Flight query's coordinator instead
of holding it until wait_timeout (#67504)
### What problem does this PR solve?
Issue Number: close #67503
Related PR: #64799 (kept the coordinator alive across GetFlightInfo →
DoGet so the BE can fetch external-table splits; this PR narrows that to
the queries that actually need it and bounds how long an abandoned
session can hold the rest)
Problem Summary:
Since #64799 every Arrow Flight SQL query whose results are produced on
the BE kept its coordinator alive after `GetFlightInfo`, until the
session ran its next query or was closed. The coordinator holds the
query's workload-group queue slot and its `active_queries` registration.
Most Flight clients open a session per query and never send
`CloseSession` (the ADBC/JDBC drivers don't do it on connection close,
as the comment in `DorisFlightSqlProducer.closeSession` notes), so each
abandoned session pinned one slot for up to `wait_timeout` (8h by
default).
Observed on a 2-FE cluster with a `max_concurrency=8` workload group:
eight finished Flight queries sat in `information_schema.active_queries`
as RUNNING for over an hour, and every later query in the group failed
with `query queue timeout, timeout: 60000 ms`. Full write-up in #67503.
This PR fixes it in two steps:
1. **Only defer the coordinator when the BE still needs it.** The
deferral added by #64799 is required only for an external-table scan in
batch mode, where the BE lazily fetches splits during `DoGet` from the
split source the coordinator holds. Every other query (internal tables,
external tables in non-batch mode) now closes its coordinator at the end
of `GetFlightInfo` again, releasing the queue slot and the
`active_queries` entry right away. That is safe: the BE buffers its
results independently of the coordinator, and finalizing the FE side
does not cancel BE execution (`QeProcessorImpl.reportExecStatus` accepts
reports for an unregistered query, and `Coordinator.close()` only
releases the queue token and the split sources). New predicates
`ScanNode.hasBatchSplitSource()` / `Coordinator.hasBatchSplitSource()`
drive the gate in `StmtExecutor.executeAndSendResult`.
2. **Bound the remaining deferred queries without killing the session.**
New mutable FE config `arrow_flight_deferred_query_idle_timeout_second`
(default 3600). Once a Flight session has been sleeping for longer than
this since its last query started, the connection timeout checker
finalizes the session's deferred executors (releasing the slot and the
registration) and leaves the session alive; `wait_timeout` still governs
the session itself. The bound is floored at the execution timeout the
deferred query actually ran with, captured when the executor is deferred
(`SET_VAR` hint values are reverted at the end of `execute()`, so the
session value cannot be read later). `0` disables the bound.
Why not kill the session (the first revision of this PR): a killed
Flight session's bearer token stays in the token cache marked as already
used, so the client's next call on it fails with `UserSession expire
after access` and has to re-handshake. With a 1h bound, every pooled or
BI-tool Flight connection that idles for an hour would fail once on its
next use. Reaping only the deferred query releases the leaked resources
with no client-visible change.
### Release note
Arrow Flight SQL: a query no longer holds its workload-group queue slot
and `active_queries` entry after `GetFlightInfo` unless it is an
external-table scan in batch mode, the only case where the BE still
fetches splits from the FE during `DoGet`. For that case a new FE config
`arrow_flight_deferred_query_idle_timeout_second` (default 1h) releases
the coordinator of an idle, never-closed session without killing the
session.
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [x] Regression test
(`arrow_flight_sql_p0/test_arrow_flight_query_release`: a finished
Flight query on an internal table no longer occupies a
`max_concurrency=1` group and is gone from `active_queries`;
`external_table_p0/iceberg/test_iceberg_arrow_flight_split_source`: the
batch-mode scan stays registered after `DoGet` and is released by the
idle reaper while the session survives)
- [x] Unit Test (`FlightSqlDeferredQueryIdleTimeoutTest`: the reaper
through `checkTimeout`, the exec-timeout floor, `0` disables, nothing
deferred, MySQL untouched; `ArrowFlightDeferralGateTest`: the
batch-split-source predicates;
`StmtExecutorTest.testDeferForArrowFlightFreezesExecTimeoutInEffect`)
- [x] Manual test (the first revision was verified on a test cluster
with an 8s bound, see the PR history; the current revision is covered by
the regression tests above)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [x] Yes. (1) A Flight query that is not an external-table batch-mode
scan releases its coordinator, queue slot and `active_queries` entry at
the end of `GetFlightInfo`, as it did before #64799. (2) The deferred
coordinator of a batch-mode scan on an idle session is released after 1h
by default instead of at `wait_timeout`; set
`arrow_flight_deferred_query_idle_timeout_second=0` for the previous
behavior. Sessions are never killed by this change.
- Does this need documentation?
- [ ] No.
- [x] Yes. New FE config
`arrow_flight_deferred_query_idle_timeout_second` — doris-website PR to
follow once this is reviewed.
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- #64799 is in 4.0.8 / 4.1.4 -->
---------
Signed-off-by: Raghvendra Singh <[email protected]>
Co-authored-by: Raghav <[email protected]>
Co-authored-by: Claude Fable 5 <[email protected]>
Co-authored-by: morningman <[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 c0f2a1b76be..a56912a77cd 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
@@ -2644,6 +2644,17 @@ public class Config extends ConfigBase {
@ConfField(description = "Maximum number of connections for the 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 = "In auto
bucketing, the number of buckets is "
+ "estimated based on the partition size. For storage "
+ "and computing integration, a partition size of 5GB " + "is
estimated as one bucket, but for cloud, a "
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 efa5d7e5406..98d9056e1af 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
@@ -134,6 +134,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 428f86c957e..80ac4cb1170 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
@@ -1006,7 +1006,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) {
@@ -1033,6 +1035,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.
*/
@@ -1305,6 +1346,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 4a9e58949ca..7502e4a5753 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
@@ -810,6 +810,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 113985e1cd1..bc0faf28758 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
@@ -208,6 +208,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;
// Optional forward target for cancellations issued on this executor:
statements that
// spawn a nested internal executor with its own query id (e.g. IVM
dry-run delta
@@ -1084,6 +1088,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.
@@ -1563,23 +1582,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..a02de20f14f
--- /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.split.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 fd40fb68dad..dac844f33b0 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
@@ -78,6 +78,39 @@ public class StmtExecutorTest extends TestWithFeService {
Assertions.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]