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]