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 14c957e048b [fix](remote-doris) End a remote Doris scan's Flight SQL 
session with the query; stop the regression framework leaking connections 
(#68338)
14c957e048b is described below

commit 14c957e048bfb3617f8cb23bc39db712089fb98f
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Tue Sep 22 09:10:14 2026 +0800

    [fix](remote-doris) End a remote Doris scan's Flight SQL session with the 
query; stop the regression framework leaking connections (#68338)
    
    ### What problem does this PR solve?
    
    Related PR: #68101 (the one connection pool, which made the leak
    visible), #68266 (the bearer token as the session's credential, which
    removed its last cap), #67503 / #62259 (the Arrow Flight deferral gate
    this generalizes)
    
    Problem Summary:
    
    **Context.** A Remote Doris catalog with `use_arrow_flight = true` reads
    another Doris cluster over Arrow Flight SQL: for every scan,
    `RemoteDorisScanNode` on the local FE performs a Flight SQL handshake
    against a remote FE (`authenticateBasicToken`), runs the query there
    (`GetFlightInfo`), and hands the endpoints - a ticket per remote BE - to
    the local BE, which reads the rows with `DoGet` straight from the remote
    BEs.
    
    The handshake is not free on the remote side: since #68101 a Flight SQL
    session is a connection in the remote FE's one connection pool, counted
    against `qe_max_connection`, the Arrow Flight SQL sub-quota and the
    catalog user's `max_user_connections`; since #68266 the bearer token is
    that session's name in the pool and nothing else, so the session ends
    only on `CloseSession`, `KILL CONNECTION`, `wait_timeout` (8h by
    default) or an FE restart.
    
    The regression framework has a leak of its own of the same shape: a
    suite's Doris connections are `ThreadLocal` to the thread that opened
    them, and only the suite thread and `Suite.thread()` close theirs; a
    `sql` on any other thread opens a connection nobody closes.
    
    **1. The problem, and what it cost**
    
    - `RemoteDorisScanNode.executeFlightSqlQuery` closed the gRPC channel
    and the allocator in a try-with-resources but never sent `CloseSession`.
    Every scan of a Remote Doris table therefore left one Flight SQL session
    behind on the remote FE, under the catalog user, until `wait_timeout`.
    Before #68101 this was invisible: Flight sessions had a pool of their
    own, were not counted per user, and a per-user LRU of
    `max_user_connections / 2` tokens evicted the oldest. After #68101 the
    leaked sessions eat the catalog user's quota on the remote FE; after
    #68266 nothing caps them at all. A hundred scans within 8h and the user
    - MySQL clients included - is refused there with `Reach limit of
    connections`.
    - This is what broke the external regression pipeline on 2026-09-21
    (TeamCity 1053416 on #68308): the `remote_doris` suites point the
    catalog at the FE under test with user `root`; 49 scans left 48 Flight
    sessions (`Arrow Flight SQL: 512 (current: 48)` in the refusal), which
    took half of root's 100; the other half was taken by the framework leak
    below, and 14 MySQL connections were refused in a six-second window.
    - The framework leak: `Awaitility.await()...until { sql ... }` evaluates
    the condition on Awaitility's own thread, which dies with the `await()`.
    Each call leaked one root connection until the client JVM
    garbage-collected it (the FE logs those as `No more data to be read.
    Close connection`). 230 suites call `Awaitility.await()` directly; in
    the failing run one suite opened 22 such connections in 46 seconds, and
    29 of them were collected in one GC at the moment the refusals stopped.
    Suites that call `sql` from threads of their own (`Thread.start {
    streamLoad }`, an `Executors` pool) leak the same way - a P0 run of the
    same day shows 76 Awaitility connections and 89 own-thread connections
    opened by root within one minute, all left to the garbage collector -
    and the two docker helpers dropped the connection their action opened
    without closing it.
    
    **2. What this PR does, and why it helps**
    
    FE:
    
    - `RemoteDorisFlightSession` (new): the session as an object -
    handshake, `execute`, and `close()` = `CloseSession` (bounded to 5s so a
    remote FE that stopped answering cannot hang the local query's teardown;
    the session is then left to its `wait_timeout` as before) followed by
    the channel and the allocator. `open()` leaves nothing behind when the
    handshake is refused; a query that fails closes its session at once, so
    the retry on the next node leaves nothing behind either. Idempotent.
    - `RemoteDorisScanNode` keeps the session from `getSplits` until
    `stop()`, which the coordinator calls when the local query closes or is
    cancelled - i.e. when the local BE is done with the remote query's
    endpoints. It cannot be closed right after `GetFlightInfo`: the remote
    FE cancels whatever a closed session was still running, and when the
    remote table is itself an external table scanned in batch mode the
    remote query is deferred there and still running while the local BE
    reads.
    - For the same reason the local coordinator has to outlive dispatch when
    the local query is itself an Arrow Flight SQL query (otherwise #67503
    closes it right after `exec()`, while the local BE may still be
    reading). The deferral gate's predicate is generalized from "has a batch
    split source" to "the BE still depends on this scan after dispatch":
    `ScanNode.hasBatchSplitSource()` -> `coordinatorMustOutliveDispatch()`,
    `Coordinator.hasBatchSplitSource()` -> `mustOutliveDispatch()`;
    `RemoteDorisScanNode` adds its open session as the second reason.
    Batch-mode external scans behave exactly as before.
    - The statement is the fallback owner. The session is opened while the
    plan is translated, before any coordinator exists, and not every plan
    gets a coordinator or gets one that is closed: a statement that fails
    between planning and dispatch (a SQL block rule on the scan, an `INSERT`
    whose transaction cannot begin - a re-used `WITH LABEL`, the per-db txn
    limit), the plan `INSERT OVERWRITE` and every materialized-view refresh
    run only to locate the sink and discard, a load job created from the
    plan. `keepFlightSession` registers the node with the
    `StatementContext`, whose `close()` (the per-statement finally of
    `ConnectProcessor`, `TaskProcessor`, `MTMVTask`) stops what is still
    registered; the deferral gate hands the nodes over to the deferred
    coordinator before `deferForArrowFlight`, so a Flight query kept alive
    for DoGet is untouched. `INSERT OVERWRITE` releases its probe plan's
    scan nodes as soon as the plan has been read.
    - A same-plan retry must not reuse a plan whose scan node released what
    the BE scans with: `handleQueryWithRetry` re-dispatches the failed
    attempt's plan after its `cancel()` stopped the scan nodes, and the
    endpoints of a remote Doris scan belong to the query of a session that
    is then gone (the remote FE tears down a query it had deferred).
    `ScanNode.cannotBeRedispatched()` (true for a remote Doris scan once
    `stop()` ended its session) makes the retry rethrow the original error
    instead.
    
    Regression framework:
    
    - `Awaitility.pollInSameThread()` at framework start-up: every `until {
    }` now runs on the suite thread and reuses the suite's connection, as
    `Suite.awaitUntil` already did. The trade-off: an `atMost()` no longer
    bounds a condition that blocks (the poll runs to completion before the
    bound is checked). A condition that runs statements is bounded by their
    timeouts; a condition that waits on anything else bounds the wait itself
    - `SuiteCluster` now runs its `doris-compose` subprocess waits on a
    helper thread joined with the command's timeout and destroys the process
    on expiry (they used to rely on `atMost()`).
    - `SuiteContext` records every connection its thread-local accessors
    open, with the thread that opened it. On every statement it closes the
    connections of threads that have finished (a suite that starts a thread
    per step - `Thread.start { streamLoad }; join`, as the mow flexible
    suites do 72 times - now holds at most the connections of the threads
    still running, deterministically, where before the population depended
    on the JVM's next GC), and when the suite ends it closes whatever is
    left, with a warning naming the suite.
    - The two docker helpers (`docker`, `dockers`) close the connection
    their action opened before restoring the original one (and the
    multi-cluster one now types that original as the `ConnectionInfo` it
    is).
    
    Tests:
    
    - `RemoteDorisScanNodeTest`: an in-process Flight SQL server counts the
    sessions it is asked to close. The session lives from the query until
    `stop()`; `stop()` twice closes once; a failed query closes at once; a
    refused handshake opens nothing; a session handed over after `stop()` is
    closed at once; the coordinator of a query with such a scan
    `mustOutliveDispatch()`; a session no coordinator takes ends with the
    statement, one handed to a deferred coordinator does not; and a node
    whose `stop()` ended a session `cannotBeRedispatched()`.
    - `ArrowFlightDeferralGateTest` follows the rename.
    - Regression
    `external_table_p0/remote_doris/test_remote_doris_flight_session`: a
    catalog logging in as a user of its own scans a table five times, then
    `INSERT OVERWRITE`s from it, then runs an `INSERT ... WITH LABEL` twice
    (the second is refused after planning), and asserts after each statement
    that `information_schema.processlist` holds no `ArrowFlightSQL` session
    of that user - assertions that fail on master.
    
    What it buys: a Remote Doris scan costs the remote FE one session for
    exactly the duration of the local query, whatever the protocol of the
    local client; the catalog user's quota on the remote FE is no longer
    consumed by history; and the regression framework no longer manufactures
    the MySQL half of the pressure.
    
    **3. The classes, and how they call each other**
    
    - `RemoteDorisScanNode` (existing): `getSplits` -> `executeQuery` ->
    `executeFlightSqlQuery(host, user, password, sql, timeout)`:
    `RemoteDorisFlightSession.open` + `execute`, then `keepFlightSession`
    (which also registers the node with
    `StatementContext.stopScanNodeAtClose`). `stop()` (from
    `Coordinator.close()` / `cancel()`, or `StatementContext.close()` as the
    fallback) closes the session; `coordinatorMustOutliveDispatch()` is true
    while one is held; `cannotBeRedispatched()` once `stop()` ended one.
    - `RemoteDorisFlightSession` (new): `open` (FlightClient +
    `authenticateBasicToken`), `execute` (`FlightSqlClient.execute`),
    `close` (`closeSession` with a 5s deadline, then client and allocator).
    - `ScanNode.coordinatorMustOutliveDispatch()` (renamed from
    `hasBatchSplitSource`): `splitAssignment != null`, overridable.
    `ScanNode.cannotBeRedispatched()` (new): false by default.
    - `Coordinator.mustOutliveDispatch()` (renamed): any scan node's
    `coordinatorMustOutliveDispatch()`.
    - `StmtExecutor.executeAndSendResult`: the deferral gate now reads
    `coord.mustOutliveDispatch()`; the deferral gate calls
    `StatementContext.handOverScanNodesToDeferredCoordinator` before
    `deferForArrowFlight`; `handleQueryWithRetry` rethrows instead of
    retrying when `planCannotBeRedispatched()`.
    - `StatementContext`: `stopScanNodeAtClose`,
    `handOverScanNodesToDeferredCoordinator`, and `close()` stopping what is
    left (after the table locks, before the connector scope).
    - `InsertOverwriteTableCommand.run`: stops the scan nodes of the plan it
    probes and discards.
    - `SuiteCluster.waitForDorisCompose`: the bounded subprocess wait
    `runCmd` / `runCmdList` use instead of `Awaitility.await().atMost(...)`.
    - Remote FE, untouched: `DorisFlightSqlProducer.closeSession` ->
    `FlightSessionsInConnectPool.closeConnectContext` ->
    `ConnectContext.cleanup()` + `cancelQuery`.
    - `RegressionTest.initGroovyEnv`: `Awaitility.pollInSameThread()`.
    - `SuiteContext`: `openedDorisConnections` (connection -> opening
    thread), `trackDorisConnection`, `closeConnectionsOfFinishedThreads`
    (from `getConnection()`, i.e. every statement), `closeDorisConnection`,
    `closeLeftoverDorisConnections` (from `close()`); `Suite.dockerImpl` /
    `dockers` call `closeDorisConnection`.
    
    ```
    local FE                                                  remote FE         
                remote BE
    RemoteDorisScanNode.getSplits
      '- executeFlightSqlQuery
           |- RemoteDorisFlightSession.open ---- handshake --> openSession 
(pool: +1 for the catalog user)
           |- session.execute ----------------- GetFlightInfo --> runs the 
query ---------------> result buffered
           '- keepFlightSession                                                 
                  (ticket per BE)
    Coordinator.exec  -> local BE ------------------------------ DoGet(ticket) 
--------------------> rows
      (Arrow Flight local client: coordinator kept alive, mustOutliveDispatch() 
== true)
    Coordinator.close / cancel
      '- scanNode.stop
           '- session.close ------------------- CloseSession --> 
closeConnectContext (pool: -1)
    ```
---
 .../doris/source/RemoteDorisFlightSession.java     | 153 ++++++++++++
 .../doris/source/RemoteDorisScanNode.java          | 149 +++++++----
 .../org/apache/doris/nereids/StatementContext.java |  60 +++++
 .../insert/InsertOverwriteTableCommand.java        |   8 +
 .../java/org/apache/doris/planner/ScanNode.java    |  27 +-
 .../main/java/org/apache/doris/qe/Coordinator.java |  17 +-
 .../java/org/apache/doris/qe/StmtExecutor.java     |  59 +++--
 .../doris/source/RemoteDorisScanNodeTest.java      | 271 +++++++++++++++++++++
 .../doris/qe/ArrowFlightDeferralGateTest.java      |  19 +-
 .../java/org/apache/doris/qe/StmtExecutorTest.java |   2 +-
 .../apache/doris/regression/RegressionTest.groovy  |  14 ++
 .../org/apache/doris/regression/suite/Suite.groovy |  14 +-
 .../doris/regression/suite/SuiteCluster.groovy     |  33 ++-
 .../doris/regression/suite/SuiteContext.groovy     |  95 ++++++--
 .../test_remote_doris_flight_session.groovy        | 130 ++++++++++
 15 files changed, 937 insertions(+), 114 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisFlightSession.java
 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisFlightSession.java
new file mode 100644
index 00000000000..ef1557c3216
--- /dev/null
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisFlightSession.java
@@ -0,0 +1,153 @@
+// 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.datasource.doris.source;
+
+import org.apache.doris.common.Pair;
+import org.apache.doris.common.UserException;
+
+import com.google.common.annotations.VisibleForTesting;
+import org.apache.arrow.flight.CallOptions;
+import org.apache.arrow.flight.CloseSessionRequest;
+import org.apache.arrow.flight.FlightClient;
+import org.apache.arrow.flight.FlightInfo;
+import org.apache.arrow.flight.Location;
+import org.apache.arrow.flight.grpc.CredentialCallOption;
+import org.apache.arrow.flight.sql.FlightSqlClient;
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.io.Closeable;
+import java.net.URI;
+import java.util.Optional;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * The Flight SQL session a {@link RemoteDorisScanNode} opens on a remote 
Doris frontend for one
+ * scan, and ends with a CloseSession once the scan is over.
+ *
+ * <p>The handshake ({@code authenticateBasicToken}) opens a session on the 
remote frontend: a
+ * connection in its pool, counted against {@code qe_max_connection}, the 
Arrow Flight SQL sub-quota
+ * and the catalog user's {@code max_user_connections}, that only a 
CloseSession, a KILL or
+ * {@code wait_timeout} ends. Closing the gRPC channel does not. A session 
opened per scan and never
+ * closed therefore stays for hours and, one scan at a time, exhausts the 
catalog user's connection
+ * quota on the remote frontend - refusing that user's MySQL connections there 
as well.
+ *
+ * <p>The session outlives GetFlightInfo on purpose: the query it ran serves 
the BE's DoGet of the
+ * endpoints, and the remote frontend cancels whatever a closed session was 
still running - the
+ * query itself, when the remote table is an external table scanned in batch 
mode and the query is
+ * therefore deferred there. So {@link #close()} is called from {@link 
RemoteDorisScanNode#stop()},
+ * when the coordinator of the local query closes, and that coordinator is 
kept alive until the BE
+ * has finished scanning ({@link 
RemoteDorisScanNode#coordinatorMustOutliveDispatch()}).
+ */
+class RemoteDorisFlightSession implements Closeable {
+    private static final Logger LOG = 
LogManager.getLogger(RemoteDorisFlightSession.class);
+
+    // A bound on the CloseSession round trip. close() runs on the local 
query's teardown path,
+    // which must not hang on a remote frontend that has stopped answering; 
the session is then left
+    // to the remote frontend's wait_timeout, as every session was before this 
class existed.
+    @VisibleForTesting
+    static final int CLOSE_SESSION_TIMEOUT_SECONDS = 5;
+
+    private final Pair<String, Integer> hostAndPort;
+    private final BufferAllocator allocator;
+    private final FlightSqlClient client;
+    private final CredentialCallOption credential;
+    private boolean closed = false;
+
+    private RemoteDorisFlightSession(Pair<String, Integer> hostAndPort, 
BufferAllocator allocator,
+            FlightSqlClient client, CredentialCallOption credential) {
+        this.hostAndPort = hostAndPort;
+        this.allocator = allocator;
+        this.client = client;
+        this.credential = credential;
+    }
+
+    /**
+     * Opens a session on the remote frontend at {@code hostAndPort} with the 
catalog's credentials.
+     * Nothing is left behind when this fails: a handshake that was refused 
opened no session, and
+     * the channel and allocator are released before the exception propagates.
+     */
+    static RemoteDorisFlightSession open(Pair<String, Integer> hostAndPort, 
String user, String password)
+            throws Exception {
+        BufferAllocator allocator = new RootAllocator();
+        FlightClient flightClient = null;
+        try {
+            URI uri = new URI("grpc", null, hostAndPort.first, 
hostAndPort.second, null, null, null);
+            flightClient = FlightClient.builder(allocator, new 
Location(uri)).build();
+            Optional<CredentialCallOption> credential = 
flightClient.authenticateBasicToken(user, password);
+            if (!credential.isPresent()) {
+                throw new UserException("Authenticates with a username and 
password failure");
+            }
+            return new RemoteDorisFlightSession(hostAndPort, allocator, new 
FlightSqlClient(flightClient),
+                    credential.get());
+        } catch (Throwable t) {
+            closeQuietly(flightClient, allocator, hostAndPort);
+            throw t;
+        }
+    }
+
+    /** Runs {@code sql} on the remote frontend; the endpoints of the result 
are where the BE reads it. */
+    FlightInfo execute(String sql, int timeoutSec) {
+        return client.execute(sql, credential, CallOptions.timeout(timeoutSec, 
TimeUnit.SECONDS));
+    }
+
+    Pair<String, Integer> getHostAndPort() {
+        return hostAndPort;
+    }
+
+    /**
+     * Ends the session on the remote frontend (CloseSession), then releases 
the channel and the
+     * allocator. Never throws: this runs on the local query's teardown path, 
and a session the
+     * remote frontend could not be told to close is only left to its 
wait_timeout. Idempotent.
+     */
+    @Override
+    public synchronized void close() {
+        if (closed) {
+            return;
+        }
+        closed = true;
+        try {
+            client.closeSession(new CloseSessionRequest(), credential,
+                    CallOptions.timeout(CLOSE_SESSION_TIMEOUT_SECONDS, 
TimeUnit.SECONDS));
+        } catch (Throwable t) {
+            LOG.warn("failed to close the Arrow Flight SQL session on remote 
Doris frontend {}:{}, it stays open there"
+                    + " until its wait_timeout", hostAndPort.first, 
hostAndPort.second, t);
+        }
+        closeQuietly(client, allocator, hostAndPort);
+    }
+
+    private static void closeQuietly(AutoCloseable client, BufferAllocator 
allocator,
+            Pair<String, Integer> hostAndPort) {
+        try {
+            if (client != null) {
+                client.close();
+            }
+        } catch (Throwable t) {
+            LOG.warn("failed to close the Arrow Flight client to remote Doris 
frontend {}:{}",
+                    hostAndPort.first, hostAndPort.second, t);
+        }
+        try {
+            allocator.close();
+        } catch (Throwable t) {
+            LOG.warn("failed to close the Arrow allocator of the Flight client 
to remote Doris frontend {}:{}",
+                    hostAndPort.first, hostAndPort.second, t);
+        }
+    }
+}
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNode.java
 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNode.java
index 729eb7f91ee..b91bbdccd18 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNode.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNode.java
@@ -44,29 +44,21 @@ import org.apache.doris.thrift.TFileRangeDesc;
 import org.apache.doris.thrift.TRemoteDorisFileDesc;
 import org.apache.doris.thrift.TTableFormatFileDesc;
 
+import com.google.common.annotations.VisibleForTesting;
 import com.google.common.base.Joiner;
 import com.google.common.collect.Lists;
-import org.apache.arrow.flight.CallOptions;
-import org.apache.arrow.flight.FlightClient;
 import org.apache.arrow.flight.FlightEndpoint;
 import org.apache.arrow.flight.FlightInfo;
 import org.apache.arrow.flight.Location;
-import org.apache.arrow.flight.grpc.CredentialCallOption;
-import org.apache.arrow.flight.sql.FlightSqlClient;
-import org.apache.arrow.memory.BufferAllocator;
-import org.apache.arrow.memory.RootAllocator;
 import org.apache.logging.log4j.LogManager;
 import org.apache.logging.log4j.Logger;
 
-import java.net.URI;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
-import java.util.Optional;
 import java.util.Set;
-import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 
 public class RemoteDorisScanNode extends FileQueryScanNode {
@@ -79,6 +71,18 @@ public class RemoteDorisScanNode extends FileQueryScanNode {
 
     private RemoteDorisSource source;
 
+    // The Flight SQL session this scan opened on the remote frontend, from 
getSplits until stop()
+    // closes it (see RemoteDorisFlightSession for why it must live that long 
and no longer). All
+    // three guarded by this: stop() may run on another thread than the one 
that planned the query -
+    // a KILL, the timeout checker - and more than once (cancel, then close).
+    private RemoteDorisFlightSession flightSession;
+    private boolean stopped;
+    // Whether stop() ended a session this scan had opened: the endpoints 
handed to the backend
+    // belong to that session's query, and a plan dispatched again with them 
(the same-plan retry
+    // of StmtExecutor.handleQueryWithRetry) would read what the remote 
frontend may have torn
+    // down with the session.
+    private boolean sessionClosedByStop;
+
     public RemoteDorisScanNode(PlanNodeId id, TupleDescriptor desc, boolean 
needCheckColumnPriv,
                                SessionVariable sv, ScanContext scanContext) {
         super(id, desc, "REMOTE_DORIS_SCAN_NODE", scanContext, 
needCheckColumnPriv, sv);
@@ -177,7 +181,8 @@ public class RemoteDorisScanNode extends FileQueryScanNode {
                     source.nextHostAndArrowPort(),
                     source.getCatalog().getUsername(),
                     source.getCatalog().getPassword(),
-                    queryStr
+                    queryStr,
+                    source.getCatalog().getQueryTimeoutSec()
                 );
             } catch (Exception e) {
                 LOG.warn("arrow request node [{}] failures {}, try next nodes",
@@ -189,17 +194,98 @@ public class RemoteDorisScanNode extends 
FileQueryScanNode {
         throw new RuntimeException("Failed to execute query: " + queryStr, 
lastException);
     }
 
-    private List<Pair<String, ByteBuffer>> executeFlightSqlQuery(Pair<String, 
Integer> hostAndPort,
-                     String user, String psw, String sql) throws Exception {
-        try (
-                BufferAllocator allocatorFE = new RootAllocator();
-                FlightClient clientFE = createFlightClient(allocatorFE, 
hostAndPort);
-                FlightSqlClient sqlClientFE = new FlightSqlClient(clientFE)
-        ) {
-            CredentialCallOption credentialCallOption = authenticate(clientFE, 
user, psw);
-            FlightInfo info = executeSqlWithTimeout(sqlClientFE, sql, 
credentialCallOption);
-
-            return processFlightEndpoints(info.getEndpoints());
+    // Opens a Flight SQL session on the remote frontend and runs the query in 
it. The session is
+    // kept until stop(): its query serves the BE's DoGet of the endpoints 
returned here. A session
+    // whose query failed is closed right away, so a retry on the next node 
leaves nothing behind.
+    @VisibleForTesting
+    List<Pair<String, ByteBuffer>> executeFlightSqlQuery(Pair<String, Integer> 
hostAndPort,
+                     String user, String psw, String sql, int timeoutSec) 
throws Exception {
+        RemoteDorisFlightSession session = 
RemoteDorisFlightSession.open(hostAndPort, user, psw);
+        FlightInfo info;
+        try {
+            info = session.execute(sql, timeoutSec);
+        } catch (Throwable t) {
+            session.close();
+            throw t;
+        }
+        keepFlightSession(session);
+        return processFlightEndpoints(info.getEndpoints());
+    }
+
+    /**
+     * Holds {@code session} until {@link #stop()}. A session handed over 
after stop() already ran,
+     * or on top of one still held, is closed at once instead: this scan owns 
one session at most,
+     * and none once stopped. The statement registers the node as well: stop() 
is the coordinator's
+     * to call, but a plan that never gets one, or whose coordinator nobody 
closes, is stopped when
+     * the statement ends instead ({@link 
org.apache.doris.nereids.StatementContext#stopScanNodeAtClose}).
+     */
+    @VisibleForTesting
+    void keepFlightSession(RemoteDorisFlightSession session) {
+        RemoteDorisFlightSession toClose;
+        synchronized (this) {
+            if (stopped) {
+                toClose = session;
+            } else {
+                toClose = flightSession;
+                flightSession = session;
+            }
+        }
+        if (toClose != null) {
+            toClose.close();
+        }
+        if (toClose != session) {
+            
ConnectContext.get().getStatementContext().stopScanNodeAtClose(this);
+        }
+    }
+
+    /**
+     * True once {@link #stop()} ended the session this scan opened: the 
endpoints in its scan
+     * ranges belong to that session's query on the remote frontend, so the 
same plan must not be
+     * dispatched again (see {@link ScanNode#cannotBeRedispatched()}).
+     */
+    @Override
+    public boolean cannotBeRedispatched() {
+        synchronized (this) {
+            return sessionClosedByStop;
+        }
+    }
+
+    /**
+     * True while this scan holds a Flight SQL session on the remote frontend: 
the local coordinator
+     * has to stay alive until the BE has finished reading the remote query, 
since closing it is what
+     * ends the session ({@link #stop()}) - and the remote frontend cancels 
what a closed session was
+     * still running. Without this, an Arrow Flight SQL query on this frontend 
would close its
+     * coordinator right after dispatch (#67503), while its BE may still be 
reading.
+     */
+    @Override
+    public boolean coordinatorMustOutliveDispatch() {
+        if (super.coordinatorMustOutliveDispatch()) {
+            return true;
+        }
+        synchronized (this) {
+            return flightSession != null;
+        }
+    }
+
+    /**
+     * Ends the Flight SQL session on the remote frontend, in addition to what 
{@code FileQueryScanNode}
+     * releases. Called by the coordinator when the local query closes or is 
cancelled, i.e. when the
+     * BE is done with (or gave up on) the remote query's endpoints.
+     */
+    @Override
+    public void stop() {
+        super.stop();
+        RemoteDorisFlightSession session;
+        synchronized (this) {
+            stopped = true;
+            session = flightSession;
+            flightSession = null;
+            if (session != null) {
+                sessionClosedByStop = true;
+            }
+        }
+        if (session != null) {
+            session.close();
         }
     }
 
@@ -290,27 +376,6 @@ public class RemoteDorisScanNode extends FileQueryScanNode 
{
             .trim().toLowerCase().startsWith("explain");
     }
 
-    private FlightClient createFlightClient(BufferAllocator allocator,
-                                            Pair<String, Integer> hostAndPort) 
throws Exception {
-        URI uri = new URI("grpc", null, hostAndPort.first, hostAndPort.second, 
null, null, null);
-        return FlightClient.builder(allocator, new Location(uri)).build();
-    }
-
-    private CredentialCallOption authenticate(FlightClient client, String 
user, String psw) throws UserException {
-        Optional<CredentialCallOption> credentialCallOption = 
client.authenticateBasicToken(user, psw);
-        if (!credentialCallOption.isPresent()) {
-            throw new UserException("Authenticates with a username and 
password failure");
-        }
-        return credentialCallOption.get();
-    }
-
-    private FlightInfo executeSqlWithTimeout(FlightSqlClient sqlClient, String 
sql,
-                                             CredentialCallOption 
credentialCallOption) {
-        int timeoutSec = source.getCatalog().getQueryTimeoutSec();
-        return sqlClient.execute(sql, credentialCallOption,
-            CallOptions.timeout(timeoutSec, TimeUnit.SECONDS));
-    }
-
     private List<Pair<String, ByteBuffer>> 
processFlightEndpoints(List<FlightEndpoint> endpoints) {
         List<Pair<String, ByteBuffer>> uniquePairs = new ArrayList<>();
         Set<String> seenPairs = new HashSet<>();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
index d6ea7ebb976..4903a00147a 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
@@ -65,6 +65,7 @@ import 
org.apache.doris.nereids.trees.plans.logical.LogicalCTEConsumer;
 import org.apache.doris.nereids.trees.plans.logical.LogicalCTEProducer;
 import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
 import org.apache.doris.nereids.util.RelationUtil;
+import org.apache.doris.planner.ScanNode;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.GlobalVariable;
 import org.apache.doris.qe.OriginStatement;
@@ -96,6 +97,7 @@ import java.util.Collections;
 import java.util.Comparator;
 import java.util.HashMap;
 import java.util.HashSet;
+import java.util.IdentityHashMap;
 import java.util.LinkedHashMap;
 import java.util.LinkedHashSet;
 import java.util.List;
@@ -218,6 +220,19 @@ public class StatementContext implements Closeable {
     // table locks
     private final Stack<CloseableResource> plannerResources = new Stack<>();
 
+    // Scan nodes that hold something on this frontend for the backend (a 
remote Doris scan's Flight
+    // SQL session on the other frontend) and release it in ScanNode.stop(), 
which the coordinator
+    // of the statement calls when it closes. Not every plan gets a 
coordinator, and not every
+    // coordinator is closed: a plan probed and discarded (INSERT OVERWRITE), 
a statement failing
+    // between planning and dispatch (a SQL block rule on the scan, an INSERT 
whose transaction
+    // cannot begin), a load job created from the plan. close() stops what is 
still registered here
+    // as the fallback (stop() is idempotent, so a coordinator that already 
closed costs nothing).
+    // A coordinator that outlives the statement on purpose - an Arrow Flight 
SQL query kept alive
+    // until DoGet, StmtExecutor.deferForArrowFlight - takes its nodes out 
first
+    // (handOverScanNodesToDeferredCoordinator). Guarded by its own monitor: 
registered on the
+    // planning thread, closed on the statement's thread or the 
forwarded-request finally.
+    private final Set<ScanNode> scanNodesToStopAtClose = 
Collections.newSetFromMap(new IdentityHashMap<>());
+
     // placeholder params for prepared statement
     private List<Placeholder> placeholders = new ArrayList<>();
 
@@ -1093,6 +1108,48 @@ public class StatementContext implements Closeable {
         }
     }
 
+    /**
+     * Registers a scan node whose {@link ScanNode#stop()} must have run by 
the time this statement
+     * ends: the coordinator of the statement runs it when it closes, and 
{@link #close()} runs it
+     * for a plan that never got a coordinator or whose coordinator nobody 
closed (see
+     * {@link #scanNodesToStopAtClose}).
+     */
+    public void stopScanNodeAtClose(ScanNode scanNode) {
+        synchronized (scanNodesToStopAtClose) {
+            scanNodesToStopAtClose.add(scanNode);
+        }
+    }
+
+    /**
+     * The coordinator of the statement outlives it on purpose (an Arrow 
Flight SQL query kept alive
+     * until the client has pulled its result, see {@code 
StmtExecutor.deferForArrowFlight}) and
+     * takes over these nodes: their {@link ScanNode#stop()} runs when that 
coordinator closes, not
+     * when this statement ends.
+     */
+    public void handOverScanNodesToDeferredCoordinator(Collection<ScanNode> 
scanNodes) {
+        synchronized (scanNodesToStopAtClose) {
+            scanNodesToStopAtClose.removeAll(scanNodes);
+        }
+    }
+
+    // The fallback of scanNodesToStopAtClose. Never throws: this runs on the 
statement's teardown
+    // path, after the statement's outcome is decided, and one node failing to 
stop must not keep
+    // the next from stopping.
+    private void stopScanNodesLeftBehind() {
+        List<ScanNode> leftBehind;
+        synchronized (scanNodesToStopAtClose) {
+            leftBehind = new ArrayList<>(scanNodesToStopAtClose);
+            scanNodesToStopAtClose.clear();
+        }
+        for (ScanNode scanNode : leftBehind) {
+            try {
+                scanNode.stop();
+            } catch (Throwable t) {
+                LOG.warn("failed to stop scan node {} at the end of the 
statement", scanNode.getId(), t);
+            }
+        }
+    }
+
     // CHECKSTYLE OFF
     @Override
     protected void finalize() throws Throwable {
@@ -1107,6 +1164,9 @@ public class StatementContext implements Closeable {
     @Override
     public void close() {
         releasePlannerResources();
+        // After the table locks: stopping a remote Doris scan's node sends a 
CloseSession to the other
+        // frontend, which must not be waited for under a lock.
+        stopScanNodesLeftBehind();
         // Fallback deterministic close of the per-statement connector scope, 
for statements that never reach the
         // query-finish callback: external DDL / SHOW / DESCRIBE / EXPLAIN / 
foreground ANALYZE run via Command.run
         // with no coordinator, so PluginDrivenScanNode.getSplits never 
registers a primary close for them. close()
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java
index 5d67020025d..fdbc056143e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/insert/InsertOverwriteTableCommand.java
@@ -61,6 +61,7 @@ import 
org.apache.doris.nereids.trees.plans.logical.UnboundLogicalSink;
 import org.apache.doris.nereids.trees.plans.physical.PhysicalOlapTableSink;
 import org.apache.doris.nereids.trees.plans.physical.PhysicalTableSink;
 import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor;
+import org.apache.doris.planner.ScanNode;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.QueryState.MysqlStateType;
 import org.apache.doris.qe.StmtExecutor;
@@ -176,6 +177,13 @@ public class InsertOverwriteTableCommand extends Command
         NereidsPlanner planner = new NereidsPlanner(ctx.getStatementContext());
         
LineageInfoExtractor.registerAnalyzePlanHook(ctx.getStatementContext(), 
planner);
         planner.plan(logicalPlanAdapter, ctx.getSessionVariable().toThrift());
+        // This plan only locates the sink and the partitions; the insert 
below plans again and runs
+        // that plan. No coordinator ever takes this one, so what its scan 
nodes opened for the
+        // backend while planning (a remote Doris scan's Flight SQL session on 
the other frontend, a
+        // batch split source) is released here, before the real insert opens 
its own.
+        for (ScanNode scanNode : planner.getScanNodes()) {
+            scanNode.stop();
+        }
         Plan analyzedPlan = planner.getAnalyzedPlan();
         lineagePlan = Optional.ofNullable(analyzedPlan);
         executor.checkBlockRules();
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 ead7960a5a1..d594fbbdbae 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
@@ -135,16 +135,31 @@ public abstract class ScanNode extends PlanNode 
implements SplitGenerator {
     }
 
     /**
-     * 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.
+     * Whether the BE still depends on something this scan node holds on the 
FE while it is
+     * scanning, so that the coordinator - closing it releases what the node 
holds, through
+     * {@link #stop()} - has to stay alive until the BE has finished scanning, 
even after the FE is
+     * done dispatching the query. Here: a batch {@link SplitSource} the BE 
fetches its splits from
+     * lazily (external-table batch mode, see {@link 
SplitGenerator#isBatchMode()}); the BE's next
+     * split fetch fails once the source is released. A subclass holding 
another such resource
+     * adds its own reason, e.g. the Flight SQL session a remote Doris scan 
keeps open on the other
+     * frontend for the query the BE reads (RemoteDorisScanNode).
      */
-    public boolean hasBatchSplitSource() {
+    public boolean coordinatorMustOutliveDispatch() {
         return splitAssignment != null;
     }
 
+    /**
+     * Whether {@link #stop()} has released something the BE would need again 
if the same plan were
+     * dispatched once more, so that a retry of the query has to plan again 
rather than reuse this
+     * node's scan ranges (StmtExecutor.handleQueryWithRetry re-dispatches the 
plan of a failed
+     * attempt whose coordinator was cancelled, and cancel() stops the scan 
nodes). A remote Doris
+     * scan's ranges are the endpoints of the query its Flight SQL session ran 
on the other frontend,
+     * gone with the session; a batch split source has the same property but 
is left as it is here.
+     */
+    public boolean cannotBeRedispatched() {
+        return false;
+    }
+
     protected abstract void createScanRangeLocations() throws UserException;
 
     /**
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 b3005119418..0acfe3d7635 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
@@ -797,18 +797,21 @@ public class Coordinator implements CoordInterface {
     }
 
     /**
-     * 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.
+     * Whether the BE still depends on this coordinator after {@link #exec()} 
returned, so it must
+     * not be closed until the BE has finished scanning: one of its scan nodes 
holds something on
+     * the FE that the BE scans with and that {@link #close()} releases
+     * ({@link ScanNode#coordinatorMustOutliveDispatch()}) - the split source 
an external-table
+     * scan in batch mode fetches its splits from lazily, or the Flight SQL 
session a remote Doris
+     * scan keeps open on the other frontend. 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() {
+    public boolean mustOutliveDispatch() {
         if (scanNodes == null) {
             return false;
         }
         for (ScanNode scanNode : scanNodes) {
-            if (scanNode.hasBatchSplitSource()) {
+            if (scanNode.coordinatorMustOutliveDispatch()) {
                 return true;
             }
         }
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 2dc45ab298c..8227950bcd2 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
@@ -814,6 +814,21 @@ public class StmtExecutor {
                 originStmt.originStmt, context.getSqlHash(), 
context.getQualifiedUser());
     }
 
+    // Whether a scan node of the current plan released, when the failed 
attempt was cancelled, what
+    // the BE would scan with again if handleQueryWithRetry dispatched the 
same plan once more
+    // (ScanNode.cannotBeRedispatched).
+    private boolean planCannotBeRedispatched() {
+        if (planner == null) {
+            return false;
+        }
+        for (ScanNode scanNode : planner.getScanNodes()) {
+            if (scanNode.cannotBeRedispatched()) {
+                return true;
+            }
+        }
+        return false;
+    }
+
     public void checkBlockRulesByScan(Planner planner) throws 
AnalysisException {
         if (planner == null) {
             return;
@@ -1174,8 +1189,9 @@ public class StmtExecutor {
     }
 
     // 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.
+    // GetFlightInfo -> DoGet phases: close the coordinator (releasing what 
its scan nodes held for
+    // the BE - external-table batch SplitSources, a remote Doris scan's 
Flight SQL session - and
+    // the query queue slot) and then unregister the query. See #62259.
     public void finalizeArrowFlightQuery() {
         try {
             if (coord != null) {
@@ -1273,6 +1289,15 @@ public class StmtExecutor {
                         }
                     }
                 }
+                if (isNeedRetry && planCannotBeRedispatched()) {
+                    // The failed attempt's cancel() stopped the scan nodes, 
and one of them released
+                    // what the BE scans with: a remote Doris scan's session 
on the other frontend,
+                    // whose query the scan ranges point at. The same plan 
cannot be dispatched again.
+                    LOG.warn("not retrying query {} with the same plan: a scan 
node released what the backend"
+                            + " scans with when the failed attempt was 
cancelled. stmt: {}",
+                            DebugUtil.printId(context.queryId()), 
parsedStmt.getOrigStmt().originStmt);
+                    throw e;
+                }
                 if (i != retryTime - 1 && isNeedRetry && 
context.getProtocolAdapter().canRetryQuery(context)) {
                     LOG.warn("retry {} times. stmt: {}", (i + 1), 
parsedStmt.getOrigStmt().originStmt);
                 } else {
@@ -1667,20 +1692,26 @@ public class StmtExecutor {
 
             if (!context.isReturnResultFromLocal()) {
                 profile.getSummaryProfile().setTempStartTime();
-                // The client pulls the results from the BE later (Arrow 
Flight SQL's 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 cannot reach here: it has no Arrow result on either 
side, so
+                // The client pulls the results from the BE later (Arrow 
Flight SQL's DoGet). Only a
+                // scan the BE keeps depending on the FE for still needs the 
coordinator after this
+                // point: an external-table scan in batch mode fetches its 
splits lazily from the
+                // split source the coordinator holds (#62259), and a remote 
Doris scan keeps the
+                // Flight SQL session open on the other frontend whose query 
the BE reads
+                // (RemoteDorisScanNode); closing the coordinator here would 
release either too early
+                // and break DoGet. 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 cannot reach here: it has no Arrow result 
on either side, so
                 // LogicalResultSinkToShortCircuitPointQuery keeps a Flight 
session on the normal
                 // execution path 
(ProtocolAdapter.supportsShortCircuitPointQuery, #67368).
-                if (coordBase == coord && coord.hasBatchSplitSource()) {
+                if (coordBase == coord && coord.mustOutliveDispatch()) {
+                    // The coordinator outlives this statement, and with it 
what its scan nodes hold
+                    // for the BE: the statement's own end must not stop them 
(StatementContext.close
+                    // is the fallback for a plan no coordinator owns), the 
coordinator's close does.
+                    
statementContext.handOverScanNodesToDeferredCoordinator(planner.getScanNodes());
                     deferForArrowFlight();
                 }
                 return;
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNodeTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNodeTest.java
new file mode 100644
index 00000000000..f947afeae2c
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNodeTest.java
@@ -0,0 +1,271 @@
+// 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.datasource.doris.source;
+
+import org.apache.doris.analysis.DescriptorTable;
+import org.apache.doris.common.Pair;
+import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.planner.PlanFragment;
+import org.apache.doris.planner.ScanNode;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.Coordinator;
+import org.apache.doris.qe.OriginStatement;
+import org.apache.doris.thrift.TUniqueId;
+
+import com.google.common.collect.Lists;
+import org.apache.arrow.flight.CallStatus;
+import org.apache.arrow.flight.CloseSessionRequest;
+import org.apache.arrow.flight.CloseSessionResult;
+import org.apache.arrow.flight.FlightDescriptor;
+import org.apache.arrow.flight.FlightEndpoint;
+import org.apache.arrow.flight.FlightInfo;
+import org.apache.arrow.flight.FlightServer;
+import org.apache.arrow.flight.Location;
+import org.apache.arrow.flight.Ticket;
+import org.apache.arrow.flight.auth2.BasicCallHeaderAuthenticator;
+import org.apache.arrow.flight.auth2.GeneratedBearerTokenAuthenticator;
+import org.apache.arrow.flight.sql.NoOpFlightSqlProducer;
+import org.apache.arrow.flight.sql.impl.FlightSql.CommandStatementQuery;
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.arrow.vector.types.pojo.Schema;
+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;
+
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+
+/**
+ * The Flight SQL session a remote Doris scan opens on the remote frontend 
lives exactly as long as
+ * the scan: opened for the query in getSplits, ended with a CloseSession when 
the coordinator stops
+ * the scan node - or when the statement ends, for a plan no coordinator ever 
took - and never left
+ * behind - not by a failed query, not by a stop() that came first. The remote 
frontend is an
+ * in-process Flight SQL server that counts what the scan does to it.
+ */
+public class RemoteDorisScanNodeTest {
+    private static final String USER = "catalog_user";
+    private static final String PASSWORD = "catalog_password";
+
+    /** A remote frontend that records the queries it ran and the sessions it 
was asked to close. */
+    private static class RecordingRemoteFrontend extends NoOpFlightSqlProducer 
{
+        final List<String> queries = new CopyOnWriteArrayList<>();
+        final List<String> closedSessions = new CopyOnWriteArrayList<>();
+
+        @Override
+        public FlightInfo getFlightInfoStatement(CommandStatementQuery 
command, CallContext context,
+                FlightDescriptor descriptor) {
+            String query = command.getQuery();
+            queries.add(query);
+            if (query.contains("boom")) {
+                throw CallStatus.INTERNAL.withDescription("query failed on the 
remote frontend").toRuntimeException();
+            }
+            FlightEndpoint endpoint = new FlightEndpoint(new 
Ticket(query.getBytes(StandardCharsets.UTF_8)),
+                    Location.forGrpcInsecure("127.0.0.1", 9999));
+            return new FlightInfo(new Schema(Collections.emptyList()), 
descriptor,
+                    Collections.singletonList(endpoint), -1, -1);
+        }
+
+        @Override
+        public void closeSession(CloseSessionRequest request, CallContext 
context,
+                StreamListener<CloseSessionResult> listener) {
+            closedSessions.add(context.peerIdentity());
+            listener.onNext(new 
CloseSessionResult(CloseSessionResult.Status.CLOSED));
+            listener.onCompleted();
+        }
+    }
+
+    private BufferAllocator serverAllocator;
+    private RecordingRemoteFrontend remote;
+    private FlightServer server;
+    private Pair<String, Integer> hostAndPort;
+    // The statement the scan plans under: a scan node registers itself with 
it when it keeps a
+    // session, so that a statement no coordinator ever takes the plan of 
still ends the session.
+    private StatementContext statementContext;
+
+    @BeforeEach
+    public void startRemoteFrontend() throws Exception {
+        ConnectContext ctx = new ConnectContext();
+        statementContext = new StatementContext(ctx, new 
OriginStatement("select 1", 0));
+        ctx.setStatementContext(statementContext);
+        ctx.setThreadLocalInfo();
+        serverAllocator = new RootAllocator();
+        remote = new RecordingRemoteFrontend();
+        // The handshake the scan performs (authenticateBasicToken) opens a 
session and issues a bearer
+        // token for it, and the peer identity a later call carries is the one 
authenticated then.
+        server = FlightServer.builder(serverAllocator, 
Location.forGrpcInsecure("127.0.0.1", 0), remote)
+                .headerAuthenticator(new GeneratedBearerTokenAuthenticator(
+                        new BasicCallHeaderAuthenticator((user, password) -> {
+                            if (USER.equals(user) && 
PASSWORD.equals(password)) {
+                                return () -> user;
+                            }
+                            throw 
CallStatus.UNAUTHENTICATED.withDescription("bad 
credentials").toRuntimeException();
+                        })))
+                .build()
+                .start();
+        hostAndPort = Pair.of("127.0.0.1", server.getPort());
+    }
+
+    @AfterEach
+    public void stopRemoteFrontend() throws Exception {
+        // The statement ends while the remote frontend is still up, as it 
does in production; a
+        // session a test left to the statement is closed here, one it already 
ended is a no-op.
+        statementContext.close();
+        ConnectContext.remove();
+        server.close();
+        serverAllocator.close();
+    }
+
+    private static RemoteDorisScanNode scanNode() {
+        // The node under test is only its session bookkeeping; the planner 
state a real node carries
+        // (descriptors, the source, the catalog) plays no part in it.
+        return Mockito.mock(RemoteDorisScanNode.class, 
Mockito.CALLS_REAL_METHODS);
+    }
+
+    private static Coordinator coordinator(ScanNode... scanNodes) {
+        return new Coordinator(1L, new TUniqueId(1L, 2L), new 
DescriptorTable(), Lists.<PlanFragment>newArrayList(),
+                Lists.newArrayList(scanNodes), "UTC", false, false);
+    }
+
+    @Test
+    public void testSessionLivesFromTheQueryUntilTheScanStops() throws 
Exception {
+        RemoteDorisScanNode node = scanNode();
+        List<Pair<String, ByteBuffer>> endpoints = 
node.executeFlightSqlQuery(hostAndPort, USER, PASSWORD,
+                "select 1", 10);
+
+        Assertions.assertEquals(1, endpoints.size());
+        Assertions.assertEquals(Collections.singletonList("select 1"), 
remote.queries);
+        // Open on the remote frontend, so the local coordinator has to 
outlive dispatch: its close is
+        // what ends the session, and the BE may still be reading the remote 
query until then.
+        Assertions.assertTrue(remote.closedSessions.isEmpty());
+        Assertions.assertTrue(node.coordinatorMustOutliveDispatch());
+        Assertions.assertTrue(coordinator(node).mustOutliveDispatch());
+
+        node.stop();
+
+        Assertions.assertEquals(Collections.singletonList(USER), 
remote.closedSessions);
+        Assertions.assertFalse(node.coordinatorMustOutliveDispatch());
+        Assertions.assertFalse(coordinator(node).mustOutliveDispatch());
+
+        // cancel() then close() both stop the scan node; the session is 
closed once.
+        node.stop();
+        Assertions.assertEquals(1, remote.closedSessions.size());
+    }
+
+    @Test
+    public void testFailedQueryClosesItsSessionAtOnce() {
+        RemoteDorisScanNode node = scanNode();
+
+        Assertions.assertThrows(Exception.class,
+                () -> node.executeFlightSqlQuery(hostAndPort, USER, PASSWORD, 
"select boom", 10));
+
+        // The retry on the next node must not leave this node's session 
behind.
+        Assertions.assertEquals(Collections.singletonList(USER), 
remote.closedSessions);
+        Assertions.assertFalse(node.coordinatorMustOutliveDispatch());
+    }
+
+    @Test
+    public void testRefusedHandshakeOpensNoSession() {
+        Assertions.assertThrows(Exception.class,
+                () -> RemoteDorisFlightSession.open(hostAndPort, USER, "wrong 
password"));
+
+        Assertions.assertTrue(remote.queries.isEmpty());
+        Assertions.assertTrue(remote.closedSessions.isEmpty());
+    }
+
+    @Test
+    public void testSessionHandedOverAfterStopIsClosedAtOnce() throws 
Exception {
+        RemoteDorisScanNode node = scanNode();
+        node.stop();
+
+        RemoteDorisFlightSession late = 
RemoteDorisFlightSession.open(hostAndPort, USER, PASSWORD);
+        node.keepFlightSession(late);
+
+        Assertions.assertEquals(Collections.singletonList(USER), 
remote.closedSessions);
+        Assertions.assertFalse(node.coordinatorMustOutliveDispatch());
+    }
+
+    @Test
+    public void testSessionOfAScanNoCoordinatorTakesEndsWithTheStatement() 
throws Exception {
+        RemoteDorisScanNode node = scanNode();
+        node.executeFlightSqlQuery(hostAndPort, USER, PASSWORD, "select 1", 
10);
+        Assertions.assertTrue(remote.closedSessions.isEmpty());
+
+        // The statement fails after planning (a SQL block rule on the scan, 
an INSERT whose
+        // transaction cannot begin) or discards the plan (the INSERT 
OVERWRITE probe): no
+        // coordinator ever calls stop(), the statement's end does.
+        statementContext.close();
+
+        Assertions.assertEquals(Collections.singletonList(USER), 
remote.closedSessions);
+        Assertions.assertFalse(node.coordinatorMustOutliveDispatch());
+        // A coordinator closing afterwards finds nothing left to end.
+        node.stop();
+        Assertions.assertEquals(1, remote.closedSessions.size());
+    }
+
+    @Test
+    public void testSessionHandedToADeferredCoordinatorOutlivesTheStatement() 
throws Exception {
+        RemoteDorisScanNode node = scanNode();
+        node.executeFlightSqlQuery(hostAndPort, USER, PASSWORD, "select 1", 
10);
+
+        // An Arrow Flight SQL query keeps its coordinator past the statement 
(deferForArrowFlight):
+        // the BE reads the remote query after the statement ended, so the 
statement's end must not
+        // close the session; the coordinator's close does, later.
+        
statementContext.handOverScanNodesToDeferredCoordinator(Collections.singletonList(node));
+        statementContext.close();
+        Assertions.assertTrue(remote.closedSessions.isEmpty());
+        Assertions.assertTrue(node.coordinatorMustOutliveDispatch());
+
+        node.stop();
+        Assertions.assertEquals(Collections.singletonList(USER), 
remote.closedSessions);
+    }
+
+    @Test
+    public void testAPlanWhoseSessionStopEndedCannotBeRedispatched() throws 
Exception {
+        RemoteDorisScanNode node = scanNode();
+        Assertions.assertFalse(node.cannotBeRedispatched());
+        node.executeFlightSqlQuery(hostAndPort, USER, PASSWORD, "select 1", 
10);
+        Assertions.assertFalse(node.cannotBeRedispatched());
+
+        // cancel() of a failed attempt stops the node: the endpoints in its 
scan ranges belong to
+        // the query of a session that is gone, so the same-plan retry must 
not dispatch them again.
+        node.stop();
+        Assertions.assertTrue(node.cannotBeRedispatched());
+
+        // A node that never held a session releases nothing when stopped.
+        RemoteDorisScanNode idle = scanNode();
+        idle.stop();
+        Assertions.assertFalse(idle.cannotBeRedispatched());
+    }
+
+    @Test
+    public void testSessionCloseIsIdempotent() throws Exception {
+        RemoteDorisFlightSession session = 
RemoteDorisFlightSession.open(hostAndPort, USER, PASSWORD);
+        session.execute("select 1", 10);
+
+        session.close();
+        session.close();
+
+        Assertions.assertEquals(Collections.singletonList(USER), 
remote.closedSessions);
+    }
+}
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
index a02de20f14f..9f8a06a7093 100644
--- 
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
@@ -33,8 +33,9 @@ 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).
+ * a coordinator has to outlive GetFlightInfo only when the BE still depends 
on one of its scans
+ * after dispatch - here an external-table scan in batch mode holding a batch 
split source (#62259).
+ * The other reason, a remote Doris scan's Flight SQL session, is covered by 
RemoteDorisScanNodeTest.
  */
 public class ArrowFlightDeferralGateTest {
 
@@ -55,15 +56,15 @@ public class ArrowFlightDeferralGateTest {
     }
 
     @Test
-    public void 
testScanNodeHasBatchSplitSourceOnlyWhenSplitsAreHandedOutLazily() throws 
Exception {
-        Assertions.assertFalse(scanNode(false).hasBatchSplitSource());
-        Assertions.assertTrue(scanNode(true).hasBatchSplitSource());
+    public void 
testScanNodeMustOutliveDispatchOnlyWhenSplitsAreHandedOutLazily() throws 
Exception {
+        
Assertions.assertFalse(scanNode(false).coordinatorMustOutliveDispatch());
+        Assertions.assertTrue(scanNode(true).coordinatorMustOutliveDispatch());
     }
 
     @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());
+    public void testCoordinatorMustOutliveDispatchIfAnyScanRequiresIt() throws 
Exception {
+        
Assertions.assertFalse(coordinator(Lists.newArrayList()).mustOutliveDispatch());
+        Assertions.assertFalse(coordinator(Lists.newArrayList(scanNode(false), 
scanNode(false))).mustOutliveDispatch());
+        Assertions.assertTrue(coordinator(Lists.newArrayList(scanNode(false), 
scanNode(true))).mustOutliveDispatch());
     }
 }
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 a37f1b05c05..95b7d17734f 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
@@ -101,7 +101,7 @@ public class StmtExecutorTest extends TestWithFeService {
     }
 
     // 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
+    // still depends on it (Coordinator.mustOutliveDispatch), 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
diff --git 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/RegressionTest.groovy
 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/RegressionTest.groovy
index 64f8d8fa2b5..c4b01932d92 100644
--- 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/RegressionTest.groovy
+++ 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/RegressionTest.groovy
@@ -38,6 +38,7 @@ import org.apache.doris.regression.util.TeamcityUtils
 import groovy.util.logging.Slf4j
 import org.apache.commons.cli.*
 import org.apache.commons.lang3.concurrent.BasicThreadFactory;
+import org.awaitility.Awaitility
 import org.codehaus.groovy.control.CompilerConfiguration
 import org.codehaus.groovy.vmplugin.v8.IndyInterface
 import org.slf4j.LoggerFactory
@@ -162,6 +163,19 @@ class RegressionTest {
 
     static void initGroovyEnv(Config config) {
         log.info("parallel = ${config.parallel}, suiteParallel = 
${config.suiteParallel}, actionParallel = ${config.actionParallel}")
+        // Evaluate every Awaitility condition on the thread that awaits it, 
as Suite.awaitUntil already
+        // does. A suite's connections are ThreadLocal to the thread that 
opened them (see
+        // SuiteContext.getConnection), so a `sql` inside 
`Awaitility.await()...until { }` on Awaitility's
+        // own polling thread opened a fresh connection that nothing closed 
when that thread died with
+        // the await(): one connection leaked on the frontend per await(), 
held until the client JVM
+        // garbage-collected it, and enough of them at once reach the user's 
max_user_connections.
+        // Polled on the suite thread, the condition reuses the suite's 
connection. The trade-off: an
+        // atMost() no longer bounds a condition that blocks - the poll runs 
to completion before the
+        // bound is checked. A condition that runs statements is bounded by 
their timeouts (the
+        // frontend's query_timeout, the framework's socketTimeout); a 
condition that waits on
+        // anything else has to bound that wait itself, as SuiteCluster does 
for its doris-compose
+        // subprocesses.
+        Awaitility.pollInSameThread()
         classloader = new GroovyClassLoader()
         compileConfig = new CompilerConfiguration()
         compileConfig.setScriptBaseClass((SuiteScript as Class).name)
diff --git 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/Suite.groovy
 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/Suite.groovy
index 6ef166d9459..7ea4fffc180 100644
--- 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/Suite.groovy
+++ 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/Suite.groovy
@@ -403,6 +403,13 @@ class Suite implements GroovyInterceptable {
             context.threadLocalConn.remove()
             actionSupplier.call()
         } finally {
+            // The connection the action opened to the docker cluster is 
unreachable once the original
+            // one is put back, so close it rather than leave it to the 
suite's end. (Still the
+            // original one when the cluster failed to start before the action 
ran.)
+            ConnectionInfo dockerConnection = context.threadLocalConn.get()
+            if (dockerConnection != null && 
!dockerConnection.is(originConnection)) {
+                context.closeDorisConnection(dockerConnection.conn, "docker 
cluster connection")
+            }
             if (originConnection == null) {
                 context.threadLocalConn.remove()
             } else {
@@ -536,13 +543,18 @@ class Suite implements GroovyInterceptable {
             // Wait for BE to report
             Thread.sleep(5000)
 
-            Connection originConnection = context.threadLocalConn.get()
+            ConnectionInfo originConnection = context.threadLocalConn.get()
             context.threadLocalConn.remove()
             context.isMultiDockerClusterRunning = true
             try {
                 actionSupplier.call(clusters)
             } finally {
                 context.isMultiDockerClusterRunning = false
+                // As in dockerImpl: the action's connection to a docker 
cluster is closed here.
+                ConnectionInfo dockerConnection = context.threadLocalConn.get()
+                if (dockerConnection != null && 
!dockerConnection.is(originConnection)) {
+                    context.closeDorisConnection(dockerConnection.conn, 
"docker cluster connection")
+                }
                 if (originConnection == null) {
                     context.threadLocalConn.remove()
                 } else {
diff --git 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteCluster.groovy
 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteCluster.groovy
index 689547f78fe..354528af0eb 100644
--- 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteCluster.groovy
+++ 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteCluster.groovy
@@ -23,14 +23,12 @@ import org.apache.doris.regression.util.JdbcUtils
 import org.apache.doris.regression.util.NodeType
 
 import com.google.common.collect.Maps
-import org.awaitility.Awaitility
 import org.slf4j.Logger
 import org.slf4j.LoggerFactory
 
 import groovy.json.JsonSlurper
 import groovy.transform.CompileStatic
 import groovy.util.logging.Slf4j
-import static java.util.concurrent.TimeUnit.SECONDS
 import java.util.stream.Collectors
 import java.sql.Connection
 
@@ -912,16 +910,34 @@ class SuiteCluster {
         runCmd(cmd, timeoutSecond)
     }
 
+    // Waits for a doris-compose process to exit and its output to be read, 
for at most timeoutSecond;
+    // one that outlives that is destroyed and the command fails, the way 
atMost() failed it when the
+    // wait ran on Awaitility's own thread. The framework polls Awaitility 
conditions on the calling
+    // thread (RegressionTest.initGroovyEnv), so an atMost() around a blocking 
call no longer bounds
+    // it: a condition that waits on something other than a statement bounds 
the wait itself. The
+    // wait runs on a helper thread so that the timeout can be enforced from 
here; waitForProcessOutput
+    // joins the two stream readers, so the buffers are complete once the 
thread has ended.
+    private static void waitForDorisCompose(Process proc, StringBuilder 
outBuf, StringBuilder errBuf,
+                                            int timeoutSecond) throws 
Exception {
+        Thread waiter = Thread.start('doris-compose-wait') {
+            proc.waitForProcessOutput(outBuf, errBuf)
+        }
+        waiter.join(timeoutSecond * 1000L)
+        if (waiter.isAlive()) {
+            proc.destroyForcibly()
+            waiter.join(10 * 1000L)
+            throw new Exception(String.format('doris compose cmd did not 
finish within %d seconds and was killed,'
+                    + ' stdout: %s, stderr: %s', timeoutSecond, 
outBuf.toString(), errBuf.toString()))
+        }
+    }
+
     private Object runCmd(String cmd, int timeoutSecond = 60) throws Exception 
{
         def fullCmd = String.format('python -W ignore %s %s -v --output-json', 
config.dorisComposePath, cmd)
         logger.info('Run doris compose cmd: {}', fullCmd)
         def proc = fullCmd.execute()
         def outBuf = new StringBuilder()
         def errBuf = new StringBuilder()
-        Awaitility.await().atMost(timeoutSecond, SECONDS).until({
-            proc.waitForProcessOutput(outBuf, errBuf)
-            return true
-        })
+        waitForDorisCompose(proc, outBuf, errBuf, timeoutSecond)
         if (proc.exitValue() != 0) {
             throw new Exception(String.format('Exit value: %s != 0, stdout: 
%s, stderr: %s',
                                               proc.exitValue(), 
outBuf.toString(), errBuf.toString()))
@@ -969,10 +985,7 @@ class SuiteCluster {
         def proc = fullCmdList.execute()
         def outBuf = new StringBuilder()
         def errBuf = new StringBuilder()
-        Awaitility.await().atMost(timeoutSecond, SECONDS).until({
-            proc.waitForProcessOutput(outBuf, errBuf)
-            return true
-        })
+        waitForDorisCompose(proc, outBuf, errBuf, timeoutSecond)
         if (proc.exitValue() != 0) {
             throw new Exception(String.format('Exit value: %s != 0, stdout: 
%s, stderr: %s',
                                               proc.exitValue(), 
outBuf.toString(), errBuf.toString()))
diff --git 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteContext.groovy
 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteContext.groovy
index fcd59cb3a73..019b82a89bb 100644
--- 
a/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteContext.groovy
+++ 
b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteContext.groovy
@@ -30,6 +30,7 @@ import groovy.util.logging.Slf4j
 import java.lang.reflect.UndeclaredThrowableException
 import java.sql.Connection
 import java.sql.DriverManager
+import java.util.concurrent.ConcurrentHashMap
 import java.util.concurrent.ExecutorService
 import java.util.function.Function
 import org.apache.doris.regression.util.JdbcUtils
@@ -55,6 +56,13 @@ class SuiteContext implements Closeable {
     public final ThreadLocal<Connection> threadHiveRemoteConn = new 
ThreadLocal<>()
     public final ThreadLocal<Connection> threadSparkIcebergConn = new 
ThreadLocal<>()
     public final ThreadLocal<Connection> threadDB2DockerConn = new 
ThreadLocal<>()
+    // Every Doris connection the thread-local accessors above opened, with 
the thread that opened it.
+    // Only the suite thread and the threads Suite.thread() runs end with 
closeThreadLocal(); a thread
+    // the suite created itself (a Thread, an Executors pool) never does, so 
its connection stayed open
+    // on the frontend until the client JVM garbage-collected it - seconds or 
minutes later, at the
+    // JVM's whim. Now the next connection the suite opens closes the 
connections of threads that have
+    // finished (closeConnectionsOfFinishedThreads), and close() closes 
whatever is left.
+    private final Map<Connection, Thread> openedDorisConnections = new 
ConcurrentHashMap<>()
     private final ThreadLocal<Syncer> syncer = new ThreadLocal<>()
     public final Config config
     public final File dataPath
@@ -148,10 +156,11 @@ class SuiteContext implements Closeable {
 
     // jdbc:mysql
     Connection getConnection() {
+        closeConnectionsOfFinishedThreads()
         def threadConnInfo = threadLocalConn.get()
         if (threadConnInfo == null) {
             threadConnInfo = new ConnectionInfo()
-            threadConnInfo.conn = getConnectionByDbName(dbName)
+            threadConnInfo.conn = 
trackDorisConnection(getConnectionByDbName(dbName))
             threadConnInfo.username = config.jdbcUser
             threadConnInfo.password = config.jdbcPassword
             threadLocalConn.set(threadConnInfo)
@@ -159,6 +168,59 @@ class SuiteContext implements Closeable {
         return threadConnInfo.conn
     }
 
+    private Connection trackDorisConnection(Connection conn) {
+        openedDorisConnections.put(conn, Thread.currentThread())
+        return conn
+    }
+
+    // Closes a connection one of the thread-local accessors opened, and 
forgets it (see openedDorisConnections).
+    void closeDorisConnection(Connection conn, String what) {
+        openedDorisConnections.remove(conn)
+        closeQuietly(conn, what)
+    }
+
+    private static void closeQuietly(Connection conn, String what) {
+        try {
+            conn.close()
+        } catch (Throwable t) {
+            log.warn("Close ${what} failed".toString(), t)
+        }
+    }
+
+    // A thread the suite created itself took its thread-local connection to 
the grave: nothing on that
+    // thread runs closeThreadLocal() once it has finished. Called on every 
statement, this closes those
+    // connections, so a suite that starts a thread per step (Thread.start { 
streamLoad ... }; join) holds
+    // at most the connections of the threads still running, not one per step 
until the suite ends.
+    private void closeConnectionsOfFinishedThreads() {
+        int closed = 0
+        for (Map.Entry<Connection, Thread> entry : 
openedDorisConnections.entrySet()) {
+            if (!entry.value.isAlive() && 
openedDorisConnections.remove(entry.key, entry.value)) {
+                closeQuietly(entry.key, "connection of finished thread 
${entry.value.name}".toString())
+                closed++
+            }
+        }
+        if (closed > 0) {
+            log.info("Closed ${closed} connection(s) opened on threads of 
suite ${suiteName} that have finished"
+                    .toString())
+        }
+    }
+
+    // The connections still open once the suite is over, whichever thread 
opened them (see
+    // openedDorisConnections). The warning names the suite: a `sql` on a 
thread the suite created
+    // itself and left running (an Executors pool it never shut down) is what 
leaves them behind.
+    private void closeLeftoverDorisConnections() {
+        List<Connection> leftover = new 
ArrayList<>(openedDorisConnections.keySet())
+        openedDorisConnections.clear()
+        if (leftover.isEmpty()) {
+            return
+        }
+        log.warn("Suite ${suiteName} left ${leftover.size()} connection(s) 
open on threads of its own, "
+                + "closing them now".toString())
+        for (Connection conn : leftover) {
+            closeQuietly(conn, "leftover connection")
+        }
+    }
+
     Connection getConnectionByDbName(String dbName) {
         def jdbcUrl = getJdbcUrl()
         def jdbcConn = DriverManager.getConnection(jdbcUrl, config.jdbcUser, 
config.jdbcPassword)
@@ -192,7 +254,7 @@ class SuiteContext implements Closeable {
         def threadConnInfo = threadLocalMasterConn.get()
         if (threadConnInfo == null) {
             threadConnInfo = new ConnectionInfo()
-            threadConnInfo.conn = getMasterConnectionByDbName(dbName)
+            threadConnInfo.conn = 
trackDorisConnection(getMasterConnectionByDbName(dbName))
             threadConnInfo.username = config.jdbcUser
             threadConnInfo.password = config.jdbcPassword
             threadLocalMasterConn.set(threadConnInfo)
@@ -204,7 +266,7 @@ class SuiteContext implements Closeable {
         def threadConnInfo = threadArrowFlightSqlConn.get()
         if (threadConnInfo == null) {
             threadConnInfo = new ConnectionInfo()
-            threadConnInfo.conn = 
config.getConnectionByArrowFlightSqlDbName(dbName)
+            threadConnInfo.conn = 
trackDorisConnection(config.getConnectionByArrowFlightSqlDbName(dbName))
             threadConnInfo.username = config.jdbcUser
             threadConnInfo.password = config.jdbcPassword
             threadArrowFlightSqlConn.set(threadConnInfo)
@@ -482,15 +544,11 @@ class SuiteContext implements Closeable {
         ConnectionInfo oldConn = threadLocalConn.get()
         if (oldConn != null) {
             threadLocalConn.remove()
-            try {
-                oldConn.conn.close()
-            } catch (Throwable t) {
-                log.warn("Close connection failed", t)
-            }
+            closeDorisConnection(oldConn.conn, "connection")
         }
 
         def newConnInfo = new ConnectionInfo()
-        newConnInfo.conn = DriverManager.getConnection(url, username, password)
+        newConnInfo.conn = 
trackDorisConnection(DriverManager.getConnection(url, username, password))
         newConnInfo.username = username
         newConnInfo.password = password
         threadLocalConn.set(newConnInfo)
@@ -585,31 +643,19 @@ class SuiteContext implements Closeable {
         ConnectionInfo conn = threadLocalConn.get()
         if (conn != null) {
             threadLocalConn.remove()
-            try {
-                conn.conn.close()
-            } catch (Throwable t) {
-                log.warn("Close connection failed", t)
-            }
+            closeDorisConnection(conn.conn, "connection")
         }
 
         ConnectionInfo master_conn = threadLocalMasterConn.get()
         if (master_conn != null) {
             threadLocalMasterConn.remove()
-            try {
-                master_conn.conn.close()
-            } catch (Throwable t) {
-                log.warn("Close master connection failed", t)
-            }
+            closeDorisConnection(master_conn.conn, "master connection")
         }
 
         ConnectionInfo arrow_flight_sql_conn = threadArrowFlightSqlConn.get()
         if (arrow_flight_sql_conn != null) {
             threadArrowFlightSqlConn.remove()
-            try {
-                arrow_flight_sql_conn.conn.close()
-            } catch (Throwable t) {
-                log.warn("Close connection failed", t)
-            }
+            closeDorisConnection(arrow_flight_sql_conn.conn, "arrow flight sql 
connection")
         }
 
         Connection hive2_docker_conn = threadHive2DockerConn.get()
@@ -677,6 +723,7 @@ class SuiteContext implements Closeable {
     @Override
     void close() {
         closeThreadLocal()
+        closeLeftoverDorisConnections()
 
         if (outputBlocksWriter != null) {
             outputBlocksWriter.close()
diff --git 
a/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_flight_session.groovy
 
b/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_flight_session.groovy
new file mode 100644
index 00000000000..b1d7cc7fd53
--- /dev/null
+++ 
b/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_flight_session.groovy
@@ -0,0 +1,130 @@
+// 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.
+
+// A remote Doris scan opens an Arrow Flight SQL session on the remote 
frontend for the query the BE
+// reads, and has to close it once the local query is over. A session left 
behind stays in the remote
+// frontend's connection pool until wait_timeout and counts against the 
catalog user's
+// max_user_connections there, so that user - MySQL clients included - is 
refused after a hundred
+// scans. The remote frontend here is this one, and the catalog logs in as a 
user of its own, so the
+// Flight sessions of that user in the processlist are exactly the ones this 
suite's scans opened.
+suite("test_remote_doris_flight_session", "p0,external") {
+    String host = context.config.otherConfigs.get("extArrowFlightSqlHost")
+    def frontends = sql "show frontends"
+    String arrowPort = frontends[0][6]
+    String httpPort = frontends[0][3]
+    String thriftPort = frontends[0][5]
+    log.info("show frontends = ${frontends}, arrow: ${arrowPort}, http: 
${httpPort}, thrift: ${thriftPort}")
+
+    String user = "test_remote_doris_flight_session_user"
+    String pwd = "C123_567p"
+    String db = "test_remote_doris_flight_session_db"
+    String table = "test_remote_doris_flight_session_t"
+    String catalog = "test_remote_doris_flight_session_catalog"
+
+    sql """DROP CATALOG IF EXISTS `${catalog}`"""
+    sql """DROP USER IF EXISTS '${user}'@'%'"""
+    sql """DROP DATABASE IF EXISTS `${db}`"""
+    sql """CREATE DATABASE `${db}`"""
+    sql """
+        CREATE TABLE `${db}`.`${table}` (
+          `id` int NOT NULL,
+          `v` varchar(16) NULL
+        ) ENGINE=OLAP
+        DUPLICATE KEY(`id`)
+        DISTRIBUTED BY HASH(`id`) BUCKETS 1
+        PROPERTIES (
+        "replication_allocation" = "tag.location.default: 1"
+        );
+    """
+    sql """INSERT INTO `${db}`.`${table}` VALUES (1, 'a'), (2, 'b'), (3, 
'c')"""
+
+    // What the remote frontend asks of the catalog user: the Flight handshake 
takes any user, the
+    // metadata REST calls need SHOW on the database, the query it runs needs 
SELECT on the table.
+    sql """CREATE USER '${user}'@'%' IDENTIFIED BY '${pwd}'"""
+    sql """GRANT SELECT_PRIV ON internal.`${db}`.* TO '${user}'@'%'"""
+    if (isCloudMode()) {
+        def clusters = sql " SHOW CLUSTERS; "
+        assertTrue(!clusters.isEmpty())
+        def validCluster = clusters[0][0]
+        sql """GRANT USAGE_PRIV ON CLUSTER `${validCluster}` TO 
'${user}'@'%'"""
+    }
+
+    sql """
+        CREATE CATALOG `${catalog}` PROPERTIES (
+            'type' = 'doris',
+            'fe_thrift_hosts' = '${host}:${thriftPort}',
+            'fe_http_hosts' = 'http://${host}:${httpPort}',
+            'fe_arrow_hosts' = '${host}:${arrowPort}',
+            'user' = '${user}',
+            'password' = '${pwd}',
+            'use_arrow_flight' = 'true'
+        );
+    """
+
+    def flightSessionsOfCatalogUser = { ->
+        def rows = sql """
+            SELECT COUNT(*) FROM information_schema.processlist
+            WHERE User = '${user}' AND Protocol = 'ArrowFlightSQL'
+        """
+        return rows[0][0] as long
+    }
+    assertEquals(0L, flightSessionsOfCatalogUser())
+
+    try {
+        // Every scan opens one session on the remote frontend. It is closed 
when the local query's
+        // coordinator closes, which happens before the result reaches the 
client, so none is left by
+        // the time the next statement runs - however many scans in a row.
+        for (int i = 0; i < 5; i++) {
+            def rows = sql """SELECT id, v FROM 
`${catalog}`.`${db}`.`${table}` ORDER BY id"""
+            assertEquals([[1, 'a'], [2, 'b'], [3, 'c']], rows)
+            assertEquals(0L, flightSessionsOfCatalogUser())
+        }
+
+        // A plan no coordinator ever takes is released when the statement 
ends instead. INSERT
+        // OVERWRITE plans the query once only to find the target partitions 
and discards that plan
+        // (its scan opened a session on the remote frontend) before the 
insert plans again ...
+        sql """DROP TABLE IF EXISTS `${db}`.`${table}_sink`"""
+        sql """
+            CREATE TABLE `${db}`.`${table}_sink` (
+              `id` int NOT NULL,
+              `v` varchar(16) NULL
+            ) ENGINE=OLAP
+            DUPLICATE KEY(`id`)
+            DISTRIBUTED BY HASH(`id`) BUCKETS 1
+            PROPERTIES (
+            "replication_allocation" = "tag.location.default: 1"
+            );
+        """
+        sql """INSERT OVERWRITE TABLE `${db}`.`${table}_sink` SELECT id, v 
FROM `${catalog}`.`${db}`.`${table}`"""
+        assertEquals([[3L]], sql("""SELECT COUNT(*) FROM 
`${db}`.`${table}_sink`"""))
+        assertEquals(0L, flightSessionsOfCatalogUser())
+
+        // ... and a statement that fails after planning - here an INSERT 
whose label was already
+        // used, refused when its transaction begins - has built no 
coordinator to close the session.
+        sql """INSERT INTO `${db}`.`${table}_sink` WITH LABEL 
test_remote_doris_flight_session_label SELECT id, v FROM 
`${catalog}`.`${db}`.`${table}`"""
+        assertEquals(0L, flightSessionsOfCatalogUser())
+        test {
+            sql """INSERT INTO `${db}`.`${table}_sink` WITH LABEL 
test_remote_doris_flight_session_label SELECT id, v FROM 
`${catalog}`.`${db}`.`${table}`"""
+            exception "already been used"
+        }
+        assertEquals(0L, flightSessionsOfCatalogUser())
+    } finally {
+        sql """DROP CATALOG IF EXISTS `${catalog}`"""
+        sql """DROP USER IF EXISTS '${user}'@'%'"""
+        sql """DROP DATABASE IF EXISTS `${db}`"""
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to