From 8608b8ecd11a511766383cdcefb4f892a80cba1e Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 21 Sep 2026 19:03:18 +0800 Subject: [PATCH 1/4] [fix](remote-doris) Close the Flight SQL session a remote Doris scan opens on the remote frontend A Remote Doris catalog with use_arrow_flight reads the other cluster over Arrow Flight SQL: for every scan RemoteDorisScanNode performs a handshake against a remote frontend, runs the query there and hands the endpoints to the BE, which reads the rows from the remote BEs with DoGet. The handshake opens a session on the remote frontend, and the scan closed the gRPC channel but never sent CloseSession, so every scan left one session behind until wait_timeout. Since #68101 a Flight session is a connection of the remote frontend's one pool, counted against the catalog user's max_user_connections; since #68266 nothing caps them any more. A hundred scans within eight hours and the catalog user, MySQL clients included, is refused there with "Reach limit of connections". The external regression hit exactly that on 2026-09-21: the remote_doris suites point the catalog at the frontend under test as root, and 48 leaked sessions took half of root's quota. RemoteDorisFlightSession is the session as an object: handshake, execute, and close() = CloseSession (bounded to five seconds, so a remote frontend that stopped answering cannot hang the local query's teardown) followed by the channel and the allocator. RemoteDorisScanNode keeps it from getSplits until stop(), which the coordinator calls when the local query closes or is cancelled - when the BE is done with the remote query's endpoints. It cannot be closed right after GetFlightInfo: the remote frontend cancels whatever a closed session was still running, and when the remote table is itself an external table in batch mode the remote query is deferred there while the BE reads. A query that fails closes its session at once, so the retry on the next node leaves nothing behind either. For the same reason the local coordinator has to outlive dispatch when the local query is an Arrow Flight SQL query, or #67503 would close it, and the session with 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.coordinatorMustOutliveDispatch() and Coordinator.mustOutliveDispatch(), with the open session as the second reason. Batch-mode external scans behave as before. RemoteDorisScanNodeTest drives the scan node against an in-process Flight SQL server that counts the sessions it is asked to close; the regression suite test_remote_doris_flight_session scans through a catalog logging in as a user of its own and asserts after each scan that the processlist holds no ArrowFlightSQL session of that user. Co-Authored-By: Claude Opus 5 (1M context) --- .../source/RemoteDorisFlightSession.java | 153 +++++++++++++ .../doris/source/RemoteDorisScanNode.java | 124 +++++++---- .../org/apache/doris/planner/ScanNode.java | 15 +- .../java/org/apache/doris/qe/Coordinator.java | 17 +- .../org/apache/doris/qe/StmtExecutor.java | 31 +-- .../doris/source/RemoteDorisScanNodeTest.java | 203 ++++++++++++++++++ .../doris/qe/ArrowFlightDeferralGateTest.java | 19 +- .../org/apache/doris/qe/StmtExecutorTest.java | 2 +- .../test_remote_doris_flight_session.groovy | 101 +++++++++ 9 files changed, 586 insertions(+), 79 deletions(-) create mode 100644 fe/fe-core/src/main/java/org/apache/doris/datasource/doris/source/RemoteDorisFlightSession.java create mode 100644 fe/fe-core/src/test/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNodeTest.java create mode 100644 regression-test/suites/external_table_p0/remote_doris/test_remote_doris_flight_session.groovy 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 00000000000000..ef1557c321613a --- /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. + * + *

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. + * + *

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 hostAndPort; + private final BufferAllocator allocator; + private final FlightSqlClient client; + private final CredentialCallOption credential; + private boolean closed = false; + + private RemoteDorisFlightSession(Pair 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 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 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 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 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 729eb7f91ee853..83b7d049fd495f 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.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,13 @@ 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). Both + // 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; + public RemoteDorisScanNode(PlanNodeId id, TupleDescriptor desc, boolean needCheckColumnPriv, SessionVariable sv, ScanContext scanContext) { super(id, desc, "REMOTE_DORIS_SCAN_NODE", scanContext, needCheckColumnPriv, sv); @@ -177,7 +176,8 @@ private List> executeQuery() { 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 +189,78 @@ private List> executeQuery() { throw new RuntimeException("Failed to execute query: " + queryStr, lastException); } - private List> executeFlightSqlQuery(Pair 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> executeFlightSqlQuery(Pair 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. + */ + @VisibleForTesting + void keepFlightSession(RemoteDorisFlightSession session) { + RemoteDorisFlightSession toClose; + synchronized (this) { + if (stopped) { + toClose = session; + } else { + toClose = flightSession; + flightSession = session; + } + } + if (toClose != null) { + toClose.close(); + } + } + + /** + * 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) { + session.close(); } } @@ -290,27 +351,6 @@ private boolean isExplainStatement() { .trim().toLowerCase().startsWith("explain"); } - private FlightClient createFlightClient(BufferAllocator allocator, - Pair 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 = 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> processFlightEndpoints(List endpoints) { List> uniquePairs = new ArrayList<>(); Set seenPairs = new HashSet<>(); 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 ead7960a5a106e..4ef9ffca2ae5a8 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,13 +135,16 @@ public TupleDescriptor getTupleDesc() { } /** - * 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; } 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 b3005119418039..0acfe3d7635087 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 void exec() throws Exception { } /** - * 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 2dc45ab298c2f1..7b06babdf3b160 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 @@ -1174,8 +1174,9 @@ void deferForArrowFlight() { } // 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) { @@ -1667,20 +1668,22 @@ public void executeAndSendResult(boolean isOutfileQuery, boolean isSendFields, 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()) { 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 00000000000000..29a9fe4b39ba59 --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/doris/source/RemoteDorisScanNodeTest.java @@ -0,0 +1,203 @@ +// 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.planner.PlanFragment; +import org.apache.doris.planner.ScanNode; +import org.apache.doris.qe.Coordinator; +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, 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 queries = new CopyOnWriteArrayList<>(); + final List 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 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 hostAndPort; + + @BeforeEach + public void startRemoteFrontend() throws Exception { + 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 { + 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.newArrayList(), + Lists.newArrayList(scanNodes), "UTC", false, false); + } + + @Test + public void testSessionLivesFromTheQueryUntilTheScanStops() throws Exception { + RemoteDorisScanNode node = scanNode(); + List> 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 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 a02de20f14f635..9f8a06a70939bb 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 @@ /** * 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 @@ private static Coordinator coordinator(List scanNodes) { } @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 a37f1b05c05b25..95b7d17734f1c8 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 void testShowNull() throws Exception { } // 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/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 00000000000000..fed4f67c027b50 --- /dev/null +++ b/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_flight_session.groovy @@ -0,0 +1,101 @@ +// 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()) + } + } finally { + sql """DROP CATALOG IF EXISTS `${catalog}`""" + sql """DROP USER IF EXISTS '${user}'@'%'""" + sql """DROP DATABASE IF EXISTS `${db}`""" + } +} From 821fb3d3d617959a86baf07903718849a3ba9bc1 Mon Sep 17 00:00:00 2001 From: morningman Date: Mon, 21 Sep 2026 19:03:34 +0800 Subject: [PATCH 2/4] [fix](regression-test) Poll Awaitility conditions on the suite thread and close the connections a suite's own threads leave behind A suite's Doris connections are ThreadLocal to the thread that opened them, and only the suite thread and the threads Suite.thread() runs close theirs in closeThreadLocal(). Awaitility.await()...until { sql ... } evaluated the condition on Awaitility's own thread, which dies with the await(), so every call opened a connection nothing closed: it stayed on the frontend until the client JVM garbage-collected it (the frontend logs those as "No more data to be read. Close connection"). 230 suites call Awaitility.await() directly; in the external regression of 2026-09-21 one of them opened 22 such connections in 46 seconds, and with the Flight sessions of the remote_doris scans counted into the same quota since #68101, root reached max_user_connections and 14 connections were refused. A P0 run of the same day shows the scale: 76 such connections in one minute, next to 89 opened on threads suites start themselves (Thread.start { streamLoad }; join, one pair per step), all of them left to the garbage collector. Awaitility.pollInSameThread() at framework start-up makes every until { } run on the suite thread and reuse the suite's connection, as Suite.awaitUntil already did. The trade-off, that atMost() cannot interrupt a condition that blocks, is bounded by the timeouts of the statements a condition runs. SuiteContext also 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, so a suite that starts a thread per step holds at most the connections of the threads still running - deterministically, where before the population depended on when the JVM next collected garbage - and when the suite ends it closes whatever is left, with a warning naming the suite. The docker helpers close the connection their action opened before putting the original one back (they dropped it before), and the multi-cluster helper types that original connection as the ConnectionInfo it is. Co-Authored-By: Claude Opus 5 (1M context) --- .../doris/regression/RegressionTest.groovy | 11 +++ .../doris/regression/suite/Suite.groovy | 14 ++- .../regression/suite/SuiteContext.groovy | 95 ++++++++++++++----- 3 files changed, 95 insertions(+), 25 deletions(-) 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 64f8d8fa2b5ec5..4c211ce5bb926c 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,16 @@ 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() cannot interrupt a condition that blocks - is bounded by the timeouts of the + // statements a condition runs. + 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 6ef166d94598eb..7ea4fffc180c53 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/SuiteContext.groovy b/regression-test/framework/src/main/groovy/org/apache/doris/regression/suite/SuiteContext.groovy index fcd59cb3a7337e..019b82a89bbdb2 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 threadHiveRemoteConn = new ThreadLocal<>() public final ThreadLocal threadSparkIcebergConn = new ThreadLocal<>() public final ThreadLocal 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 openedDorisConnections = new ConcurrentHashMap<>() private final ThreadLocal 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 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 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() From ed2ad48a2de76b0efc9b5e3e1a63f8009e2751a0 Mon Sep 17 00:00:00 2001 From: morningman Date: Tue, 22 Sep 2026 00:42:50 +0800 Subject: [PATCH 3/4] [fix](remote-doris) Release a remote Doris scan's Flight SQL session when no coordinator ever owns its plan The session a RemoteDorisScanNode opens on the remote frontend was released only by the coordinator's stop(), and it is opened while the plan is translated, before any coordinator exists. Every plan that never got a coordinator, or got one nobody closed, leaked it exactly as before: a statement failing 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 then discard, a load job created from the plan. And the same-plan retry of handleQueryWithRetry re-dispatched the endpoints of a session the failed attempt's cancel() had just closed - the remote frontend tears down a query it had deferred with the session. The statement is now the fallback owner: keepFlightSession registers the node with the StatementContext, whose close() (ConnectProcessor's per-statement finally, TaskProcessor, MTMVTask) stops what is still registered; the deferral gate hands the nodes over to the deferred coordinator first, so an Arrow Flight query kept alive for DoGet is untouched, and stop() is idempotent, so a coordinator that already closed costs nothing. INSERT OVERWRITE releases its probe plan's scan nodes as soon as the plan has been read, before the real insert opens its own. A retry is refused, with the original error, when a scan node reports cannotBeRedispatched() - a remote Doris scan once stop() ended its session. RemoteDorisScanNodeTest covers the statement-end release, the handover to a deferred coordinator and the predicate; test_remote_doris_flight_session adds an INSERT OVERWRITE and a re-used WITH LABEL and asserts that neither leaves a session behind. Co-Authored-By: Claude Opus 5 (1M context) --- .../doris/source/RemoteDorisScanNode.java | 33 +++++++-- .../doris/nereids/StatementContext.java | 60 ++++++++++++++++ .../insert/InsertOverwriteTableCommand.java | 8 +++ .../org/apache/doris/planner/ScanNode.java | 12 ++++ .../org/apache/doris/qe/StmtExecutor.java | 28 ++++++++ .../doris/source/RemoteDorisScanNodeTest.java | 72 ++++++++++++++++++- .../test_remote_doris_flight_session.groovy | 29 ++++++++ 7 files changed, 236 insertions(+), 6 deletions(-) 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 83b7d049fd495f..b91bbdccd1851f 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 @@ -72,11 +72,16 @@ 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). Both - // 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). + // 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) { @@ -210,7 +215,9 @@ List> executeFlightSqlQuery(Pair hostA /** * 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. + * 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) { @@ -226,6 +233,21 @@ void keepFlightSession(RemoteDorisFlightSession 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; + } } /** @@ -258,6 +280,9 @@ public void stop() { stopped = true; session = flightSession; flightSession = null; + if (session != null) { + sessionClosedByStop = true; + } } if (session != null) { session.close(); 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 d6ea7ebb9767cc..4903a00147ac23 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.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.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 enum TableFrom { // table locks private final Stack 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 scanNodesToStopAtClose = Collections.newSetFromMap(new IdentityHashMap<>()); + // placeholder params for prepared statement private List placeholders = new ArrayList<>(); @@ -1093,6 +1108,48 @@ public synchronized void releasePlannerResources() { } } + /** + * 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 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 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 @@ protected void finalize() throws Throwable { @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 5d67020025d797..fdbc056143e0ee 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.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 void run(ConnectContext ctx, StmtExecutor executor) throws Exception { 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 4ef9ffca2ae5a8..d594fbbdbaec6d 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 @@ -148,6 +148,18 @@ 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/StmtExecutor.java b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java index 7b06babdf3b160..8227950bcd2b26 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 void checkBlockRulesByRegex(OriginStatement originStmt) throws AnalysisEx 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; @@ -1274,6 +1289,15 @@ private void handleQueryWithRetry(TUniqueId queryId) throws Exception { } } } + 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 { @@ -1684,6 +1708,10 @@ public void executeAndSendResult(boolean isOutfileQuery, boolean isSendFields, // LogicalResultSinkToShortCircuitPointQuery keeps a Flight session on the normal // execution path (ProtocolAdapter.supportsShortCircuitPointQuery, #67368). 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 index 29a9fe4b39ba59..f947afeae2cb04 100644 --- 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 @@ -19,9 +19,12 @@ 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; @@ -56,8 +59,9 @@ /** * 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, 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. + * 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"; @@ -95,9 +99,16 @@ public void closeSession(CloseSessionRequest request, CallContext context, private RecordingRemoteFrontend remote; private FlightServer server; private Pair 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 @@ -117,6 +128,10 @@ public void startRemoteFrontend() throws Exception { @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(); } @@ -190,6 +205,59 @@ public void testSessionHandedOverAfterStopIsClosedAtOnce() throws Exception { 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); 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 index fed4f67c027b50..b1d7cc7fd533cd 100644 --- 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 @@ -93,6 +93,35 @@ suite("test_remote_doris_flight_session", "p0,external") { 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}'@'%'""" From aaba58f2e6ea8f5a79145353902c359ddc2da670 Mon Sep 17 00:00:00 2001 From: morningman Date: Tue, 22 Sep 2026 00:42:56 +0800 Subject: [PATCH 4/4] [fix](regression-test) Bound the doris-compose waits of SuiteCluster themselves now that Awaitility polls on the calling thread Awaitility.pollInSameThread() runs a condition to completion before atMost() is checked, so the two SuiteCluster waits that put the whole doris-compose subprocess wait inside a condition (runCmd, runCmdList) had lost their only time bound: a hung docker daemon would hang the docker suite thread for as long as the subprocess lived instead of failing the suite after timeoutSecond, and with dockerSuiteParallel = 1 every later docker suite of the run with it. They now wait on a helper thread joined with the command's timeout and destroy the process on expiry; waitForProcessOutput joins the stream readers, so the output is complete once the thread has ended. The comment at the global default now says what the trade-off really is: a condition that runs statements is bounded by their timeouts, one that waits on anything else has to bound that wait itself. Co-Authored-By: Claude Opus 5 (1M context) --- .../doris/regression/RegressionTest.groovy | 9 +++-- .../regression/suite/SuiteCluster.groovy | 33 +++++++++++++------ 2 files changed, 29 insertions(+), 13 deletions(-) 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 4c211ce5bb926c..c4b01932d9276f 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 @@ -169,9 +169,12 @@ class RegressionTest { // 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() cannot interrupt a condition that blocks - is bounded by the timeouts of the - // statements a condition runs. + // 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() 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 689547f78fec2f..354528af0eb2ef 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()))