github-actions[bot] commented on code in PR #68711:
URL: https://github.com/apache/doris/pull/68711#discussion_r4175774476


##########
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);

Review Comment:
   [P1] Preserve each catalog's filesystem settings when reusing a connection. 
Fluss 1.0's `FlussConnection` constructor calls process-wide 
`FileSystem.initialize` with `client.fs.*`; snapshot downloads later resolve 
files through that static configuration. After catalog A returns a connection, 
catalog B opens one with different filesystem credentials, and A borrows its 
old lease here, A skips initialization and reads its snapshot using B's 
settings (often an access failure or wrong storage account). The full-config 
pool key isolates leases but not Fluss's static filesystem state. Scope 
filesystem resolution to the scan's config, or avoid this cross-catalog reuse, 
and cover an A/B/A primary-key snapshot sequence.



##########
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:
   [P2] Bound idle connections across configuration keys. Every successful 
range is appended here, with no per-key or global limit, and another catalog 
cannot borrow that deque. For example, ten distinct catalogs scanned 
sequentially with 128 ranges each leave 1,280 connections, Netty threads, and 
sockets idle despite at most 128 ranges running together. They remain until the 
60-second timeout and a reaper pass, potentially longer when closing saturates. 
This can exhaust BE resources in ordinary multi-catalog workloads; cap total 
idle clients and close excess leases promptly.



##########
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:
   [P2] Retain the lease until deferred PK cleanup actually succeeds. On an 
early close, this `thenRun` is the only handback after `lease` and 
`fullScanner` are cleared. The sole publication waiter can die on an 
`OutOfMemoryError` from Fluss's `pollBatch()` (or fail to be constructed), 
leaving `released` incomplete; if it does complete, an OOME submitting 
`discard` to the closer is captured by this ignored dependent future. Neither 
failure has a retry after heap recovery, so the SDK reader or connection can 
remain live indefinitely. Keep a recoverable owner and observe the handback 
outcome; test failures on both sides of late publication.



##########
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:
   [P2] Keep connections owned until their close or borrow handoff succeeds. 
`pollFirst()` runs before `expired.add()` can allocate; an `OutOfMemoryError` 
there, or while logging after all dequeues, reaches `sweep()`'s catch with 
connections absent from `idle` and never passed to the closer. Later sweeps 
cannot recover them, leaving their Fluss client threads and sockets live. 
`borrow()` similarly allocates a `Lease` after removing an idle connection or 
creating a fresh one. Make these ownership transfers failure-safe and test an 
OOME after removal.



-- 
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