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


##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java:
##########
@@ -0,0 +1,266 @@
+// 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.HashMap;
+import java.util.Iterator;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Consumer;
+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. So the pool never holds more connections for a client 
configuration than ranges
+ * with that configuration have been 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. The bound above is per configuration, not for the pool as a 
whole: what is idle under
+ * each configuration adds up until it is closed. One catalog alone has two, 
since its {@code PK_FULL}
+ * ranges leave the log read preference at fluss's default and its other 
ranges set it
+ * ({@code FlussJniScanner#clientConfig}).
+ *
+ * <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, 
FlussConnectionCloser::close);
+
+    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,

Review Comment:
   [P1] Start the reaper from a recoverable path. The first 
`FlussConnectionPool.INSTANCE` access runs this static initializer, and OpenJDK 
17 `scheduleWithFixedDelay` starts its worker before returning. If heap or 
native-thread exhaustion makes that start throw `OutOfMemoryError`, Java marks 
`FlussConnectionPool` erroneous; later Fluss scans in the same cached plugin 
classloader fail with `NoClassDefFoundError` even after resources recover. This 
disables Fluss for the rest of the BE process, beyond the per-range waiter 
failure discussed at 4175774483. Keep fallible thread startup out of class 
initialization, retry it after recovery, and cover a failed initial start 
followed by a successful scan. See the [JDK 
scheduler](https://github.com/openjdk/jdk17u/blob/master/src/java.base/share/classes/java/util/concurrent/ScheduledThreadPoolExecutor.java)
 and [JLS 
12.4.2](https://docs.oracle.com/javase/specs/jls/se17/html/jls-12.html#jls-12.4.2).



##########
be/src/exec/operator/file_scan_operator.cpp:
##########
@@ -327,6 +328,14 @@ Status 
FileScanLocalState::_process_conjuncts(RuntimeState* state) {
 
 Status FileScanOperatorX::prepare(RuntimeState* state) {
     RETURN_IF_ERROR(ScanOperatorX<FileScanLocalState>::prepare(state));
+    // Sharded for the scanners that can reach the cache at once, which are 
now every instance's:
+    // as many as the per-instance caches it replaces had between them. That 
is the fragment's
+    // instances, not this operator's parallelism: a serial operator has one 
instance, and it runs
+    // the scanners of all of them (max_scanners_concurrency).
+    const int shard_num =
+            std::min(ScannerScheduler::default_remote_scan_thread_num(),
+                     state->query_parallel_instance_num() * 
file_scanners_per_instance(state));

Review Comment:
   [P2] Cap the scanner product before multiplying two ints. 
`max_file_scanners_concurrency` is an unrestricted session int; with two 
parallel instances and value 2147483647, this new signed product overflows 
before `std::min` can cap it. On the usual wrapped result, `std::max(shard_num, 
1)` creates a single shared shard. `KVCache::get` holds that shard's mutex 
while loading Iceberg delete files, so independent readers across instances 
become serialized and can time out; the previous per-instance caches remained 
sharded. Compute in 64 bits or clamp before multiplication, then test a high 
legal setting with multiple instances.



##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java:
##########
@@ -0,0 +1,266 @@
+// 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.HashMap;
+import java.util.Iterator;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Consumer;
+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. So the pool never holds more connections for a client 
configuration than ranges
+ * with that configuration have been 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. The bound above is per configuration, not for the pool as a 
whole: what is idle under
+ * each configuration adds up until it is closed. One catalog alone has two, 
since its {@code PK_FULL}
+ * ranges leave the log read preference at fluss's default and its other 
ranges set it
+ * ({@code FlussJniScanner#clientConfig}).
+ *
+ * <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, 
FlussConnectionCloser::close);
+
+    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;
+    /** Closes, without waiting for it, a connection the pool lets go of: 
{@link FlussConnectionCloser#close}. */
+    private final Consumer<Connection> closer;
+
+    /**
+     * 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,
+            Consumer<Connection> closer) {
+        this.factory = factory;
+        this.nanoTime = nanoTime;
+        this.closer = closer;
+    }
+
+    /**
+     * 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:
   [P1] Make the idle insertion exception-safe. `computeIfAbsent` publishes an 
empty deque before `new Idle(...)` succeeds; an OOME on the first return then 
makes later `borrow` and reaper sweeps dereference null even after heap 
recovery. There is a second failure at capacity: OpenJDK 17 
[`ArrayDeque.addLast`](https://github.com/openjdk/jdk17u/blob/master/src/java.base/share/classes/java/util/ArrayDeque.java)
 writes the 17th entry and advances `tail` before growing its array, so an OOME 
during growth leaves `head == tail`; the next same-key return can make all 
older idle connections unreachable to the reaper. Both effects persist after 
the failed scan and exceed the single lost lease in discussion 4175774478. 
Prepare a populated replacement before publishing, or otherwise keep the 
map/deque valid on every failed insertion, and cover first-return and growth 
failures followed by borrow/sweep.



##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionCloser.java:
##########
@@ -0,0 +1,124 @@
+// 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.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.SynchronousQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+
+/**
+ * Closes fluss connections without making the caller wait for them.
+ *
+ * <p>Closing a {@link Connection} shuts its netty client down with netty's 
default graceful period,
+ * which returns only after two quiet seconds, and fluss offers no shorter 
one. Scan ranges no longer
+ * close connections - they give them back to {@link FlussConnectionPool} - 
but the pool closes the
+ * ones nobody has borrowed for a while, and one sweep can expire every 
connection a burst of ranges
+ * opened. Closed one after another they would hold that sweep up for two 
seconds each, so the pool
+ * hands them here and moves on.
+ *
+ * <p><b>Bounded.</b> A connection that is closing still holds its client 
threads (three netty threads
+ * each) until the quiet period passes, so at most {@link #MAX_CLOSING} close 
here at a time; past that
+ * the caller closes its own, which only slows the sweep that brought it.
+ */
+final class FlussConnectionCloser {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(FlussConnectionCloser.class);
+
+    /** Connections closing in the background at once; each is one thread here 
and three of its own. */
+    static final int MAX_CLOSING = 256;
+
+    /** Scanners closing inline again is worth a line in the log, not one per 
connection. */
+    private static final long SATURATION_LOG_INTERVAL_NANOS = 
TimeUnit.MINUTES.toNanos(1);
+
+    private static final AtomicInteger THREAD_COUNTER = new AtomicInteger();
+
+    private static final AtomicLong LAST_SATURATION_LOG_NANOS =
+            new AtomicLong(System.nanoTime() - SATURATION_LOG_INTERVAL_NANOS);
+
+    /**
+     * One thread per closing connection and no queue: a queued connection 
would hold its threads
+     * while waiting for a turn, which is the pile-up the bound exists to 
prevent. A connection that
+     * finds every thread busy is closed by the thread that brought it.
+     */
+    private static final ExecutorService CLOSER = new ThreadPoolExecutor(
+            0, MAX_CLOSING, 60L, TimeUnit.SECONDS, new SynchronousQueue<>(),
+            runnable -> {
+                Thread thread = new Thread(runnable,
+                        "fluss-connection-closer-" + 
THREAD_COUNTER.incrementAndGet());
+                // Never what keeps BE's JVM alive, and never worth waiting 
for at shutdown.
+                thread.setDaemon(true);
+                // Shutting the client down can still load fluss classes; like 
JniScanner does around
+                // open and close, run it under the loader of the plugin that 
can see them.
+                
thread.setContextClassLoader(FlussConnectionCloser.class.getClassLoader());
+                return thread;
+            },
+            (close, executor) -> closeOnCallingThread(close));
+
+    private FlussConnectionCloser() {
+    }
+
+    /**
+     * Every closer thread is busy, so the thread that brought the connection 
closes it, two seconds
+     * a connection, and says so in the log: nothing else would show it.
+     */
+    private static void closeOnCallingThread(Runnable close) {
+        long now = System.nanoTime();
+        long last = LAST_SATURATION_LOG_NANOS.get();
+        if (now - last >= SATURATION_LOG_INTERVAL_NANOS && 
LAST_SATURATION_LOG_NANOS.compareAndSet(last, now)) {
+            LOG.info("{} fluss connections are closing in the background; 
further connections are closed "
+                    + "by the thread that brings them, waiting about two 
seconds for each, until some of "
+                    + "those are done", MAX_CLOSING);
+        }
+        close.run();
+    }
+
+    /**
+     * Closes {@code connection} without making the caller wait for it, unless 
{@link #MAX_CLOSING}
+     * are closing already. A failure to close is logged and nothing else: the 
scans the connection
+     * served have already returned their rows, and nothing may fail over the 
cleanup of one of their
+     * connections.
+     *
+     * <p>Nothing else holds a connection by the time it gets here - the pool 
and the range have let go of
+     * it - so one that cannot be handed to a closer thread is closed by the 
caller. Handing it over
+     * allocates the task and perhaps a thread, which fails with an {@code 
OutOfMemoryError} while BE's JVM
+     * heap is full or no native thread is to be had; dropped there, the 
connection would keep its threads
+     * and sockets for the life of the process.
+     */
+    static void close(Connection connection) {
+        try {
+            CLOSER.execute(() -> closeQuietly(connection));
+        } catch (RuntimeException | Error handoffFailure) {
+            closeQuietly(connection);
+        }
+    }
+
+    private static void closeQuietly(Connection connection) {
+        try {
+            connection.close();
+        } catch (Exception e) {

Review Comment:
   [P2] Retain a retry owner when a background close throws `Error`. This 
catches `Exception` only, so an OOME inside an accepted close task kills that 
worker after the pool and scanner have dropped the connection. Fluss 1.0 
[`NettyClient.close`](https://github.com/apache/fluss/blob/v1.0.0/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java)
 allocates its shutdown-future list before shutting down its event group and 
also catches only `Exception`; an OOME there leaves Netty threads live 
indefinitely, with no later sweep able to reach them. The earlier handoff 
threads cover failures before task acceptance; this is a separate worker-side 
failure. Observe failed close tasks, including `Error`, and retry after heap 
recovery.



##########
fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java:
##########
@@ -348,14 +369,39 @@ 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

Review Comment:
   [P2] Wait for old remote-log downloads before pooling a completed range. A 
bounded LOG, PK_TAIL, or PK_FULL scan can reach its stop offset while Fluss has 
prefetched later remote segments: `poll` sends the next fetch before returning 
progress, and 
[`RemoteLogDownloader.close`](https://github.com/apache/fluss/blob/v1.0.0/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/log/RemoteLogDownloader.java)
 leaves in-flight tasks on the connection's shared download pool. `finished` 
still selects `giveBack` here (`released()` in PK_FULL tracks only the 
snapshot), so the next range can queue its own remote reads behind obsolete 
downloads; a slow object store can stall an otherwise short query. Keep the 
lease out of the pool until those tasks finish or are cancelled, or give each 
range an independent downloader; cover bounded EOF with a slow prefetched 
segment.



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