morningman commented on code in PR #68711:
URL: https://github.com/apache/doris/pull/68711#discussion_r4176191912


##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java:
##########
@@ -0,0 +1,237 @@
+// 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.fluss;
+
+import org.apache.fluss.client.Connection;
+import org.apache.fluss.client.ConnectionFactory;
+import org.apache.fluss.config.Configuration;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayDeque;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Function;
+import java.util.function.LongSupplier;
+
+/**
+ * Fluss connections kept open from one scan range to the next, so that a 
range borrows a connection
+ * instead of opening one of its own.
+ *
+ * <p>A connection is a netty client with its own event loop threads, a 
metadata cache, and - once a
+ * range has read a kv snapshot or a remote log segment through it - a pool of 
download threads. Opening
+ * one per range cost every range the connect, and a close that waits out 
netty's two-second graceful
+ * shutdown. {@link FlussConnectionCloser} took that close off the scanning 
thread, but only up to the
+ * number of closes it runs at once: back-to-back queries over many small 
ranges, or one query over a
+ * partitioned table with hundreds of them, still left ranges closing their 
own connections, two
+ * seconds each. A range read to its end now gives its connection back 
instead; only one that failed or
+ * was closed early (a LIMIT, a cancel) has its connection closed, for the 
reasons {@link #discard} gives.
+ *
+ * <p><b>One range per connection at a time.</b> A fluss connection is 
thread-safe and could serve every
+ * range at once, but then its download threads ({@code 
client.remote-file.download-thread-num}, three by
+ * default) would be shared by every primary-key range copying its kv snapshot 
at the same moment, where
+ * each such range used to have three of its own. Lending a connection to one 
range at a time keeps every
+ * range's client resources what they were; what changes is that a connection 
outlives its range and
+ * serves the next one. The pool never holds more connections than ranges were 
read at once.
+ *
+ * <p><b>Keyed by the whole client configuration.</b> A connection goes back 
to the ranges that would have
+ * opened it with the same settings, so ranges of catalogs with different 
servers, credentials or client
+ * options never share one. The key holds whatever the configuration holds, 
credentials included, and is
+ * never logged.
+ *
+ * <p><b>Idle connections are closed after {@link #IDLE_TIMEOUT_NANOS}.</b> 
What one burst of ranges
+ * opened is there for the next burst and closed, through the closer, once 
nothing has borrowed it for
+ * that long - a BE that stops reading fluss does not keep the threads.
+ */
+final class FlussConnectionPool {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(FlussConnectionPool.class);
+
+    /** How long a connection nobody borrows is kept. */
+    static final long IDLE_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(60);
+
+    /** How often idle connections are looked at: one lives at most this much 
past its timeout. */
+    private static final long REAP_INTERVAL_SECONDS = 10;
+
+    static final FlussConnectionPool INSTANCE =
+            new FlussConnectionPool(ConnectionFactory::createConnection, 
System::nanoTime);
+
+    static {
+        ScheduledExecutorService reaper = 
Executors.newSingleThreadScheduledExecutor(runnable -> {
+            Thread thread = new Thread(runnable, "fluss-connection-reaper");
+            // Never what keeps BE's JVM alive.
+            thread.setDaemon(true);
+            // The closer may hand a close back to this thread, and closing 
can load fluss classes.
+            
thread.setContextClassLoader(FlussConnectionPool.class.getClassLoader());
+            return thread;
+        });
+        reaper.scheduleWithFixedDelay(() -> sweep(INSTANCE), 
REAP_INTERVAL_SECONDS, REAP_INTERVAL_SECONDS,
+                TimeUnit.SECONDS);
+    }
+
+    private final Function<Configuration, Connection> factory;
+    private final LongSupplier nanoTime;
+
+    /**
+     * Idle connections by the configuration they were opened with, the most 
recently returned last. A
+     * configuration with no idle connection has no entry, which is what 
{@link #borrow} relies on.
+     */
+    private final Map<Map<String, String>, ArrayDeque<Idle>> idle = new 
HashMap<>();
+
+    FlussConnectionPool(Function<Configuration, Connection> factory, 
LongSupplier nanoTime) {
+        this.factory = factory;
+        this.nanoTime = nanoTime;
+    }
+
+    /**
+     * A connection opened with {@code config}: the one given back most 
recently if any is idle, else a
+     * new one. The most recent rather than the oldest, so that after a burst 
the connections it no longer
+     * needs are the ones left idle long enough to be closed.
+     */
+    Lease borrow(Configuration config) {
+        Map<String, String> key = config.toMap();
+        synchronized (idle) {
+            ArrayDeque<Idle> connections = idle.get(key);
+            if (connections != null) {
+                Connection connection = connections.pollLast().connection;
+                if (connections.isEmpty()) {
+                    idle.remove(key);
+                }
+                return new Lease(key, connection, false);
+            }
+        }
+        // Outside the lock: opening a connection talks to the cluster.
+        return new Lease(key, factory.apply(config), true);
+    }
+
+    /** Takes back the connection of a range that was read to its end, for the 
next range to borrow. */
+    void giveBack(Lease lease) {
+        long now = nanoTime.getAsLong();
+        synchronized (idle) {
+            idle.computeIfAbsent(lease.key, key -> new 
ArrayDeque<>()).addLast(new Idle(lease.connection, now));
+        }
+    }
+
+    /**
+     * Closes, without waiting for it, the connection of a range that failed 
or was closed before its end,
+     * instead of lending it again. Such a connection may be broken by what 
failed it: an
+     * {@code OutOfMemoryError} on one of its netty threads ends that thread's 
event loop for good. The
+     * most recently returned connection is lent first, so one like that, 
given back, would be lent to
+     * every range that came next. A primary-key range closed before its kv 
snapshot arrived hands its
+     * connection here only once the snapshot copy running on the connection's 
download threads is over
+     * ({@code FlussJniScanner#closeInternal}); closed under the copy, the 
connection would strand it.
+     */
+    void discard(Lease lease) {
+        FlussConnectionCloser.close(lease.connection);
+    }
+
+    /** Closes, without waiting for them, the connections nobody has borrowed 
for {@code idleNanos}. */
+    void closeIdleLongerThan(long idleNanos) {
+        long now = nanoTime.getAsLong();
+        List<Connection> expired = new ArrayList<>();
+        synchronized (idle) {
+            Iterator<ArrayDeque<Idle>> keys = idle.values().iterator();
+            while (keys.hasNext()) {
+                ArrayDeque<Idle> connections = keys.next();
+                // A deque is in the order its connections came back, so the 
expired ones are at its head.
+                while (!connections.isEmpty() && now - 
connections.peekFirst().since >= idleNanos) {
+                    expired.add(connections.pollFirst().connection);

Review Comment:
   Fixed in bb33f27a1ad, for the sweep. It no longer takes the expired 
connections out together: it takes one at a time and hands it to the closer at 
once, with nothing allocated in between (the deque unlinks it, and an emptied 
entry is removed through the iterator, which reuses the hash the map stored), 
and logs after the loop. An error mid-sweep now costs at most the connection in 
hand; the rest stay in the pool for the next sweep. 
`FlussConnectionCloser.close` closes a connection on the calling thread when 
building the task or a closer thread fails, instead of dropping it. Test: 
`FlussConnectionPoolTest.sweepThatFailsMidwayKeepsTheConnectionsItHadNotHandedOver`
 -- a closer that fails its second hand-over; the third connection stays idle 
and the next sweep closes it. It fails with the previous sweep. The existing 
test threw its `OutOfMemoryError` from `nanoTime`, before anything was taken 
out, so it covered only the next sweeps running.
   
   `borrow()` and `giveBack()` stay as they are: each allocates one small 
object while it moves a connection, so an error there loses at most that 
connection. No transfer can be made failure-proof under a full heap -- HotSpot 
can even skip a frame's catch and finally when deoptimization fails to 
reallocate -- and after BE's JVM runs out of heap the fluss client's own netty 
event loops can die as well, leaving requests that never complete. What the 
pool must not do is lose a whole batch, or stop cleaning up for good, and it 
now does neither.
   



##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java:
##########
@@ -0,0 +1,237 @@
+// 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.fluss;
+
+import org.apache.fluss.client.Connection;
+import org.apache.fluss.client.ConnectionFactory;
+import org.apache.fluss.config.Configuration;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayDeque;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Function;
+import java.util.function.LongSupplier;
+
+/**
+ * Fluss connections kept open from one scan range to the next, so that a 
range borrows a connection
+ * instead of opening one of its own.
+ *
+ * <p>A connection is a netty client with its own event loop threads, a 
metadata cache, and - once a
+ * range has read a kv snapshot or a remote log segment through it - a pool of 
download threads. Opening
+ * one per range cost every range the connect, and a close that waits out 
netty's two-second graceful
+ * shutdown. {@link FlussConnectionCloser} took that close off the scanning 
thread, but only up to the
+ * number of closes it runs at once: back-to-back queries over many small 
ranges, or one query over a
+ * partitioned table with hundreds of them, still left ranges closing their 
own connections, two
+ * seconds each. A range read to its end now gives its connection back 
instead; only one that failed or
+ * was closed early (a LIMIT, a cancel) has its connection closed, for the 
reasons {@link #discard} gives.
+ *
+ * <p><b>One range per connection at a time.</b> A fluss connection is 
thread-safe and could serve every
+ * range at once, but then its download threads ({@code 
client.remote-file.download-thread-num}, three by
+ * default) would be shared by every primary-key range copying its kv snapshot 
at the same moment, where
+ * each such range used to have three of its own. Lending a connection to one 
range at a time keeps every
+ * range's client resources what they were; what changes is that a connection 
outlives its range and
+ * serves the next one. The pool never holds more connections than ranges were 
read at once.
+ *
+ * <p><b>Keyed by the whole client configuration.</b> A connection goes back 
to the ranges that would have
+ * opened it with the same settings, so ranges of catalogs with different 
servers, credentials or client
+ * options never share one. The key holds whatever the configuration holds, 
credentials included, and is
+ * never logged.
+ *
+ * <p><b>Idle connections are closed after {@link #IDLE_TIMEOUT_NANOS}.</b> 
What one burst of ranges
+ * opened is there for the next burst and closed, through the closer, once 
nothing has borrowed it for
+ * that long - a BE that stops reading fluss does not keep the threads.
+ */
+final class FlussConnectionPool {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(FlussConnectionPool.class);
+
+    /** How long a connection nobody borrows is kept. */
+    static final long IDLE_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(60);
+
+    /** How often idle connections are looked at: one lives at most this much 
past its timeout. */
+    private static final long REAP_INTERVAL_SECONDS = 10;
+
+    static final FlussConnectionPool INSTANCE =
+            new FlussConnectionPool(ConnectionFactory::createConnection, 
System::nanoTime);
+
+    static {
+        ScheduledExecutorService reaper = 
Executors.newSingleThreadScheduledExecutor(runnable -> {
+            Thread thread = new Thread(runnable, "fluss-connection-reaper");
+            // Never what keeps BE's JVM alive.
+            thread.setDaemon(true);
+            // The closer may hand a close back to this thread, and closing 
can load fluss classes.
+            
thread.setContextClassLoader(FlussConnectionPool.class.getClassLoader());
+            return thread;
+        });
+        reaper.scheduleWithFixedDelay(() -> sweep(INSTANCE), 
REAP_INTERVAL_SECONDS, REAP_INTERVAL_SECONDS,
+                TimeUnit.SECONDS);
+    }
+
+    private final Function<Configuration, Connection> factory;
+    private final LongSupplier nanoTime;
+
+    /**
+     * Idle connections by the configuration they were opened with, the most 
recently returned last. A
+     * configuration with no idle connection has no entry, which is what 
{@link #borrow} relies on.
+     */
+    private final Map<Map<String, String>, ArrayDeque<Idle>> idle = new 
HashMap<>();
+
+    FlussConnectionPool(Function<Configuration, Connection> factory, 
LongSupplier nanoTime) {
+        this.factory = factory;
+        this.nanoTime = nanoTime;
+    }
+
+    /**
+     * A connection opened with {@code config}: the one given back most 
recently if any is idle, else a
+     * new one. The most recent rather than the oldest, so that after a burst 
the connections it no longer
+     * needs are the ones left idle long enough to be closed.
+     */
+    Lease borrow(Configuration config) {
+        Map<String, String> key = config.toMap();
+        synchronized (idle) {
+            ArrayDeque<Idle> connections = idle.get(key);
+            if (connections != null) {
+                Connection connection = connections.pollLast().connection;
+                if (connections.isEmpty()) {
+                    idle.remove(key);
+                }
+                return new Lease(key, connection, false);
+            }
+        }
+        // Outside the lock: opening a connection talks to the cluster.
+        return new Lease(key, factory.apply(config), true);
+    }
+
+    /** Takes back the connection of a range that was read to its end, for the 
next range to borrow. */
+    void giveBack(Lease lease) {
+        long now = nanoTime.getAsLong();
+        synchronized (idle) {
+            idle.computeIfAbsent(lease.key, key -> new 
ArrayDeque<>()).addLast(new Idle(lease.connection, now));

Review Comment:
   The javadoc was wrong across configurations and is corrected in bb33f27a1ad; 
there is no global cap, on purpose.
   
   The bound is per client configuration, and that is the design: connections 
opened with different settings must not be shared. In practice a catalog has at 
most two configurations -- `PK_FULL` ranges leave the log read preference at 
fluss's default and the other ranges set it (`FlussJniScanner#clientConfig`) -- 
so a BE reading one fluss catalog keeps at most the peaks of those two, for up 
to about 70 s (the 60 s idle timeout plus the 10 s reaper interval). An idle 
connection holds one network thread with its selector and the sockets to the 
servers it reached; its buffers come from netty's shared 
`PooledByteBufAllocator.DEFAULT`. While a burst runs there is one connection 
per open range, as on master, where each connection also carried four network 
threads; a cap on idle connections does not bound that, and below the 
per-configuration peak it would only turn pooled connections back into a 
connect and a close per range. Ten catalogs each scanned 128 ranges wide within 
one minute woul
 d hold what the example says until the reaper gets to them; I don't see that 
as an ordinary workload to size a constant for.
   



##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java:
##########
@@ -348,14 +369,36 @@ protected int getNext() throws IOException {
     @Override
     protected void closeInternal() throws IOException {
         IOException failure = null;
+        CompletableFuture<Void> released = fullScanner == null
+                ? CompletableFuture.completedFuture(null) : 
fullScanner.released();
         // Close everything even if an earlier close throws: a leaked fluss 
connection keeps its netty
         // and metadata-updater threads alive for the life of the BE process.
         failure = closeQuietly(scanner, "scanner", failure);
         scanner = null;
+        fullScanner = null;
         failure = closeQuietly(table, "table", failure);
         table = null;
-        failure = closeQuietly(connection, "connection", failure);
-        connection = null;
+        if (lease != null) {
+            FlussConnectionPool.Lease returned = lease;
+            // A range read to its end gives the connection back to serve the 
next one (see
+            // FlussConnectionPool); one that failed or was closed before its 
end (a LIMIT, a cancel) has it
+            // closed instead, for the reasons FlussConnectionPool#discard 
gives.
+            Runnable handBack = finished && failure == null
+                    ? () -> FlussConnectionPool.INSTANCE.giveBack(returned)
+                    : () -> FlussConnectionPool.INSTANCE.discard(returned);
+            if (released.isDone()) {
+                handBack.run();
+            } else {

Review Comment:
   Fixed in bb33f27a1ad, short of a separate owner:
   
   - The waiter keeps waiting after an `OutOfMemoryError` on its own thread, 
parking one poll interval between attempts. It still closes the reader only 
after publication or an initializer failure; initializer failures, an 
`OutOfMemoryError` on fluss's download side included, reach it as exceptions 
from `pollBatch` (`KvSnapshotBatchScanner#ensureNoException`), so it was only 
an error raised on the waiter's own thread that ended it.
   - Building the waiter thread moved inside the `try` that already covered 
starting it, so a failure to construct it also falls back to waiting on the 
closing thread.
   - The deferred hand-back is always a discard (a range whose snapshot never 
arrived was not read to its end), and `FlussConnectionCloser.close` now closes 
a connection on the calling thread when it cannot hand it to a closer thread, 
so the future `thenRun` returns has nothing left to lose.
   
   Test: 
`SafeKvSnapshotAndLogBatchScannerTest.publicationWaiterOutlastsAnOutOfMemoryErrorOnItsOwnThread`
 -- the waiter's first poll throws `OutOfMemoryError`; it polls again, closes 
the reader once after publication, and `released()` completes. Without the 
change the waiter dies and never polls again.
   
   For the record, the comparison with master: there, closing a primary-key 
range before its snapshot arrived closed the connection under the copy, and the 
copy thread, the waiter and the half-copied directory leaked every time. The 
deferral in this PR is what fixes that; what remains needs an error on the 
waiter's own thread. A registry swept by the reaper would also cover a 
deoptimization that skips the waiter's catch, but it would replace the waiter 
rather than harden it, for a case in which the fluss client's own threads do 
not recover either.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to