This is an automated email from the ASF dual-hosted git repository.

Gabriel39 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 7ac0eb2fdb4 [improvement](fluss) Pool fluss connections and share a 
scan node's reader cache across instances (#68711)
7ac0eb2fdb4 is described below

commit 7ac0eb2fdb426595f810e94ab77b7aa2eece12d5
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Sun Oct 4 22:00:06 2026 +0800

    [improvement](fluss) Pool fluss connections and share a scan node's reader 
cache across instances (#68711)
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: #66399 (fluss catalog)
    
    Problem Summary:
    
    **In short.** Two fixed costs made fluss scans slow regardless of how
    much data they read. Every scan range spent about two seconds closing a
    fluss connection of its own, invisible in the profile. A union read of a
    primary-key table also read each bucket's log tail once per pipeline
    instance instead of once per backend. Together they made a union read
    with any log tail slower than not using union read at all: `COUNT(*)`
    took 8.6 s against 4.9 s with union read disabled on a log table, and
    21.9 s against 11.7 s on a primary-key table. This PR pools fluss
    connections across scan ranges and shares a file scan node's reader
    cache across its instances. A union-read `COUNT(*)` with a tail is now
    2-15x faster than union read disabled, and a `COUNT(*)` on a three-row
    table that took 2 s takes 24 ms.
    
    **Background**
    
    - BE reads a fluss table through JNI: `FlussJniScanner` (plugin
    `fe/be-java-extensions/fluss-scanner`) reads one scan range, which is
    one bucket of a log table, one bucket of a primary-key table (`PK_FULL`,
    a kv snapshot plus the change log after it), or the log tail of a union
    read (`LOG` / `PK_TAIL`).
    - A union read splits a table into the part already tiered into the lake
    (paimon splits, read natively or through paimon JNI) and the log tail
    not tiered yet (fluss ranges). For a primary-key table every lake split
    must drop the rows whose keys the tail updated, so
    `FlussUnionLakeReader` first reads the keys of that bucket's tail (a JNI
    read) and caches them in the scan node's `ShardedKVCache`, where the
    other lake splits of the bucket find them. The same cache holds iceberg
    delete files and deletion vectors.
    - A scan node runs as several pipeline instances per backend (up to
    `parallel_pipeline_task_num`), and splits are dealt to instances with no
    regard for the bucket or delete file they share.
    
    **The problem, and what it cost**
    
    1. **Two seconds per range in close.** Every range opened a fluss
    `Connection` in `openInternal` and closed it in `closeInternal`, on the
    scan thread. Closing shuts the connection's netty client down with
    netty's default graceful shutdown, which returns only after two quiet
    seconds (`NettyClient#close`); fluss has no shorter close. The profile
    does not show it: BE publishes the JNI scanner's counters before it
    calls the Java `close()`. Measured on master (Release BE, local fluss
    cluster, 16-bucket tables of 30M rows):
    
       | `COUNT(*)` on master | Time | Where it went |
       |---|---|---|
       | 3-row log table | 2069 ms | 2044 ms closing its one connection |
    | 30M-row log table, 16 ranges | 4480 ms | about 32 s of closes, summed
    over the ranges |
    | log lake table with a 1% log tail, union read | 8320 ms | 32 ms
    reading the tail; nearly all the rest closing connections |
       | the same table with no tail, union read | 130 ms | |
       | the same table, union read disabled | 4610 ms | |
    
       A small tail made union read slower than not using it.
    2. **The tail keys read once per instance.** The cache lived in
    `FileScanLocalState`, one per pipeline instance. On a 16-bucket
    primary-key table whose lake has 80 splits spread over 9 instances,
    `COUNT(*)` read the tails 69 times where 16 would do (cache hits 11 of
    80), each read paying the two seconds above, and kept 6.6M keys in the
    JVM heap for 1.6M distinct ones: 21.9 s, against about 11.7 s with union
    read disabled. Iceberg delete files and deletion vectors were likewise
    parsed once per instance.
    3. **A connection per range costs more than its close.** Every range
    also paid a connect, and every connection started fluss's default of
    four netty network threads, each with its own selector (on macOS a
    kqueue and a pipe pair). Eight concurrent `SELECT *` over a 16-bucket
    table held over a hundred connections and up to 2400 file descriptors in
    BE; past 1024 descriptors, libcurl on macOS (which waits with
    `select()`) fails BE's S3 reads with "A libcurl function was given a bad
    argument".
    4. **A bug in the same lifecycle (also on master).** A `PK_FULL` range
    starts copying its bucket's kv snapshot on the connection's download
    threads when it opens. If the range is closed before the snapshot
    arrives (the query was cancelled or failed elsewhere), closing the
    connection runs `shutdownNow()` on its download pool, which drops the
    queued files, and the copy waits for them without a timeout. The copy
    thread, the thread waiting to close the snapshot reader, and a
    half-copied snapshot directory under `java.io.tmpdir/fluss` then stay
    for the life of BE.
    
    **How this PR fixes it**
    
    Four commits, and a fifth from review that keeps an `OutOfMemoryError`
    in the cleanup of the first two from losing a connection (described with
    them):
    
    1. **Pool connections; close the rest off the scan thread.** Ranges
    borrow connections from the new `FlussConnectionPool`:
    - one range at a time per connection, so a primary-key range keeps three
    download threads to itself, as now; a connection just outlives its range
    and serves the next one;
    - keyed by the whole client configuration, so catalogs with different
    servers, credentials or options never share one;
    - a range read to its end gives the connection back, and the most
    recently returned connection is lent next. A range that failed or was
    closed early (a LIMIT, a cancel) has its connection closed instead,
    since it may be broken by what failed it;
    - a reaper closes connections nobody borrowed for 60 s. Its sweep
    catches `Throwable`, because a periodic task that throws once is never
    run again, and takes expired connections out one at a time, handing each
    to the closer at once, so an error mid-sweep costs at most the
    connection in hand.
    
    Discarded and idle connections are closed by `FlussConnectionCloser` on
    daemon threads, at most 256 at a time; past that the caller closes its
    own and the log says so once a minute. A connection that cannot be
    handed to a closer thread at all (an `OutOfMemoryError` building the
    task or the thread) is closed by the caller as well. A failure to close
    is logged and does not fail a query whose rows were already returned.
    The profile gains `FlussJniConnectionsOpened`.
    2. **Keep a `PK_FULL` range's connection until its snapshot reader lets
    go.** `SafeKvSnapshotAndLogBatchScanner#released()` completes once the
    fluss snapshot reader is closed. `closeInternal` gives back or discards
    the connection after that: at once in every case except a range closed
    before its snapshot arrived, whose connection is handed over by the
    waiter thread after the late publication. The waiter keeps waiting
    through an `OutOfMemoryError` on its own thread, and a failure to build
    it falls back to waiting on the closing thread, as a failure to start it
    already did.
    3. **One network thread per connection.** A pooled connection serves one
    range at a time, one stream of requests, so BE opens connections with
    `netty.client.num-network-threads = 1` unless the catalog sets
    `fluss.netty.client.num-network-threads`, in which case its value wins
    and those connections are pooled apart. FE's connection, one per
    catalog, keeps the fluss default.
    4. **One reader cache per scan node per backend (BE).** `ShardedKVCache`
    moves from `FileScanLocalState` to `FileScanOperatorX` and is created in
    `prepare()`, so every instance of the node reads through one cache. It
    is safe to share: every cached value is plain data (roaring bitmaps,
    delete-row vectors, key blocks with an immutable hash index) referring
    to no instance's profile, IO context or runtime state; the cache locks
    per shard and was already shared by the scanners of an instance; nothing
    is erased and the operator outlives every local state and scanner. The
    shard count stays what the per-instance caches had between them: the
    fragment's instance count times the per-instance scanner limit, capped
    by the remote scan thread pool.
    
    What it buys:
    
    - Union read of fluss tables is faster than union read disabled at every
    tail size.
    - Small tables and back-to-back queries no longer cost two seconds;
    repeated queries open no connection at all.
    - Iceberg delete files, deletion vectors and fluss tail keys are parsed
    once per scan node per backend.
    - Fewer netty threads and file descriptors in BE under concurrent fluss
    scans.
    - A cancelled primary-key read no longer leaks threads and partly copied
    snapshot files.
    
    **Results**
    
    Release BE, 2 GB JVM heap, local fluss 1.0.0 cluster, 30M-row tables of
    16 buckets; master against this PR on the same data:
    
    | Query | master | this PR |
    |---|---|---|
    | 3-row log table, `COUNT(*)` | 2069 ms | 24 ms |
    | 30M-row log table, `COUNT(*)` | 4480 ms | 400 ms |
    | log table with a 15M-row tail, union read, `COUNT(*)` | 8612 ms | 405
    ms |
    | the same with union read disabled | 4864 ms | 568 ms |
    | primary-key table (compacted lake, 100K updated keys per bucket in the
    tail), union read, `COUNT(*)` | 21929 ms | 1038 ms |
    | same, tail-key cache hits / time in tail reads summed over scanners |
    11 of 80 / 139.8 s | 61 of 80 / 0.9 s |
    | same table with an uncompacted lake (32 JNI merge splits), union read,
    `COUNT(*)` | 13.8 s | 1.6 s |
    
    Union read against union read disabled after this PR, by tail size
    (`COUNT(*)`, ms):
    
    | Log table tail | 0.1% | 1% | 10% | 50% |
    |---|---|---|---|---|
    | union read | 175 | 177 | 213 | 375 |
    | union read disabled | 633 | 621 | 568 | 839 |
    
    | Primary-key table tail (keys per bucket) | 1K | 10K | 100K |
    |---|---|---|---|
    | union read, compacted lake | 490 | 641 | 1079 |
    | union read, uncompacted lake | 1229 | 1357 | 1660 |
    | union read disabled | about 7600-8000 | | |
    
    What the pool adds over closing in the background alone, which stops at
    its 256 closes in flight:
    
    | Workload | background close only | pooled (this PR) |
    |---|---|---|
    | 64-bucket log table, `COUNT(*)` 20 times back to back | median 133 ms,
    every fourth query 2.2 s | median 103 ms, slowest 128 ms |
    | connections a repeated query opens (`FlussJniConnectionsOpened`) | one
    per range | 0 |
    | partitioned primary-key lake table (30 partitions x 4 buckets, 1%
    tail), union `COUNT(*)`, 240 JNI reads per query | 2090 ms | 1667 ms |
    | the same, four at once | 11.3 s each, 3326 BE threads | 4.9 s each,
    2067 BE threads |
    | 16-bucket log table, `SELECT *` with 16 scanners | 10.8 s | 10.5 s |
    | 16-bucket primary-key table, `SELECT *` with 16 scanners | 6.67 s |
    6.46 s |
    
    One network thread per connection (16-bucket log and primary-key tables
    of 10M rows, only the plugin swapped, four A/B rounds):
    
    | Measure | 4 network threads (fluss default) | 1 network thread (this
    PR) |
    |---|---|---|
    | after one `SELECT *` (16 connections left in the pool): netty client
    threads / file descriptors | 32 / 625 | 16 / 479 |
    | 4 and 8 concurrent `SELECT *`: peak file descriptors in BE | 1238-2415
    | 811-1167 |
    | 4 and 8 concurrent `SELECT *`: peak threads in BE | 2186-2361 |
    2085-2245 |
    | single and concurrent queries: median latency | | within 5% either way
    |
    
    Snapshot copy (commit 2), after six out-of-memory `SELECT *` queries in
    a row over a 16-bucket primary-key table:
    
    | Left behind | before the fix | with the fix |
    |---|---|---|
    | copy threads parked in `FileDownloadUtils.transferAllDataToDirectory`
    | 2 | 0 |
    | `fluss-snapshot-publication-close` threads blocked behind them | 2 | 0
    |
    | half-copied snapshot directories under `java.io.tmpdir/fluss` | 2 (17
    MB, 4.3 MB) | 0 |
    
    Union read, union read disabled and the expected values agree on
    `COUNT(*)`, `SUM(id)` and `SUM(price)` in every state measured.
    
    **Classes, and how they call each other**
    
    - `FlussJniScanner` (existing): `openInternal` borrows a `Lease` from
    the pool and sets one network thread unless the catalog set a number;
    `closeInternal` closes scanner and table, then gives the connection back
    (read to the end) or discards it (failed or closed early), for a
    `PK_FULL` range only after
    `SafeKvSnapshotAndLogBatchScanner#released()`; reports
    `FlussJniConnectionsOpened`.
    - `FlussConnectionPool` (new): `borrow` / `giveBack` / `discard` /
    `closeIdleLongerThan`, idle connections per client configuration, most
    recent first; the `fluss-connection-reaper` thread sweeps every 10 s.
    - `FlussConnectionCloser` (new): closes connections on up to 256 daemon
    threads, on the caller's thread past that or when one cannot be handed
    over.
    - `SafeKvSnapshotAndLogBatchScanner` / `PublicationGuardedBatchScanner`
    (existing): `released()` completes when the fluss snapshot reader has
    been closed, by whichever path closed it.
    - `FileScanOperatorX` (BE, existing): owns `_kv_cache` now, created in
    `prepare()`; `FileScanLocalState` hands it to every `FileScanner` /
    `FileScannerV2` it creates.
    - `FlussUnionLakeReader`, the iceberg delete-file readers, the
    deletion-vector readers (BE, untouched): use the cache through the
    scanner, as before.
    
    ```
    BE, one file scan node                    instance 1 .. N (pipeline 
instances of the node on this BE)
      FileScanOperatorX::_kv_cache  <---------- every scanner of every instance 
(was: one cache per instance)
        '- FlussUnionLakeReader: keys of bucket b's log tail, read once per BE
        '- iceberg delete files, deletion vectors
    
      scanner --JNI--> FlussJniScanner (one per range)
                         openInternal:  FlussConnectionPool.borrow(client 
config) -> an idle connection, or a new one
                                        (netty.client.num-network-threads = 1 
unless the catalog sets it)
                         ... read ...
                         closeInternal: close scanner, table
                                        read to its end ? pool.giveBack(lease) 
: pool.discard(lease)
                                        PK_FULL closed before its snapshot 
arrived: after SafeKvSnapshotAndLogBatchScanner.released()
    
      FlussConnectionPool --discard, or idle for 60 s (reaper, every 10 s)--> 
FlussConnectionCloser
                                                                                
<= 256 daemon threads, each waiting out
                                                                                
netty's 2 s graceful shutdown
    ```
    
    Not touched: FE planning, which ranges and splits a union read produces,
    and how a range reads its records.
    
    **Merge order.** No dependency on other PRs. Best merged before #68712:
    once a scan runs up to 16 scanners per instance by default again, every
    scan opens up to 16 fluss connections at once, which this PR reuses and
    slims.
    
    ### Release note
    
    Fluss catalog scans reuse fluss client connections across scan ranges
    and no longer spend about two seconds per scan range closing one, which
    speeds up aggregations, union reads with a log tail, back-to-back
    queries and queries on small or heavily partitioned fluss tables. Fluss
    union reads of primary-key tables, iceberg delete files and deletion
    vectors are read once per scan node per backend instead of once per
    pipeline instance. BE opens fluss connections with one network thread
    unless the catalog sets `fluss.netty.client.num-network-threads`. A
    cancelled read of a fluss primary-key table no longer leaves threads and
    partly copied snapshot files behind.
    
    ### Check List (For Author)
    
    - Test <!-- At least one of them must be included. -->
        - [x] Regression test
        - [x] Unit Test
        - [x] Manual test (add detailed scripts or steps below)
        - [ ] No need to test or manual test. Explain why:
    - [ ] This is a refactor/code format and no logic has been changed.
            - [ ] Previous test can cover this change.
            - [ ] No code files have been changed.
            - [ ] Other reason <!-- Add your reason?  -->
    
    **Unit tests.** `FlussConnectionPoolTest` (new, 9): a range borrows what
    the range before it gave back; ranges reading at once each get their own
    connection and none is lent twice (8 threads x 2000 borrows); other
    settings never share; discarded and idle connections are closed; a sweep
    that throws `OutOfMemoryError` leaves the next sweeps coming, and one
    that fails while handing connections to the closer keeps the ones it had
    not reached. `FlussConnectionCloserTest` (new, 2): the caller does not
    wait for the close; past 256 the caller closes its own.
    `FlussJniScannerLogTest` (3 new): closing a range does not wait for its
    connection; a range closed early or failing does not give its connection
    back; one network thread unless the catalog sets the number.
    `FlussJniScannerPkTest` (1 new): the first snapshot file is swapped for
    a FIFO that holds the copy, the range is closed mid-copy and the
    snapshot directory must go (still there after 60 s without the fix, gone
    in 1.3 s with it). `SafeKvSnapshotAndLogBatchScannerTest` (3 new): the
    connection is released only once the snapshot reader is closed, at once
    if the snapshot had arrived, and the publication waiter outlasts an
    `OutOfMemoryError` on its own thread. BE
    `FileScanOperatorFlussTest.EveryInstanceOfTheNodeReadsThroughTheNodesCache`
    (new). On this branch: fluss-scanner 79 tests pass; the two tests added
    in review fail without their fix. BE: the file scan, file scanner,
    scanner scheduling, fluss union and iceberg reader suites (168 tests)
    pass.
    
    **Regression.** `external_table_p0/fluss`, 17 suites, pass against a
    local fluss 1.0.0 docker environment.
    
      **Manual.** The benchmarks above on a Release BE.
    
    - Behavior changed:
        - [ ] No.
    - [x] Yes. <!-- Explain the behavior change --> BE keeps idle fluss
    connections for up to 60 s and opens them with one network thread unless
    the catalog sets `fluss.netty.client.num-network-threads`; a failure to
    close a fluss connection is logged instead of failing the query; the
    scan node's reader cache is shared by its instances on a backend; new
    profile counter `FlussJniConnectionsOpened`.
    
    - Does this need documentation?
        - [x] No.
    - [ ] Yes. <!-- Add document PR link here. eg:
    https://github.com/apache/doris-website/pull/1214 -->
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label <!-- Add branch pick label that this PR
    should merge into -->
    
    ---------
    
    Co-authored-by: Claude Opus 5.5 (1M context) <[email protected]>
---
 be/src/exec/operator/file_scan_operator.cpp        |  33 ++-
 be/src/exec/operator/file_scan_operator.h          |  11 +-
 .../operator/file_scan_operator_fluss_test.cpp     | 108 +++++++
 .../apache/doris/fluss/FlussConnectionCloser.java  | 124 ++++++++
 .../apache/doris/fluss/FlussConnectionPool.java    | 323 +++++++++++++++++++++
 .../org/apache/doris/fluss/FlussJniScanner.java    |  60 +++-
 .../fluss/SafeKvSnapshotAndLogBatchScanner.java    |  55 +++-
 .../doris/fluss/FlussConnectionCloserTest.java     | 106 +++++++
 .../doris/fluss/FlussConnectionPoolTest.java       | 298 +++++++++++++++++++
 .../apache/doris/fluss/FlussJniScannerLogTest.java | 130 ++++++++-
 .../apache/doris/fluss/FlussJniScannerPkTest.java  |  89 ++++++
 .../SafeKvSnapshotAndLogBatchScannerTest.java      |  98 +++++++
 .../org/apache/doris/fluss/StubConnection.java     |  64 ++++
 13 files changed, 1462 insertions(+), 37 deletions(-)

diff --git a/be/src/exec/operator/file_scan_operator.cpp 
b/be/src/exec/operator/file_scan_operator.cpp
index 608bde2aa6c..8eb03873b76 100644
--- a/be/src/exec/operator/file_scan_operator.cpp
+++ b/be/src/exec/operator/file_scan_operator.cpp
@@ -22,6 +22,7 @@
 #include <algorithm>
 #include <memory>
 
+#include "common/cast_set.h"
 #include "core/assert_cast.h"
 #include "core/data_type/data_type_array.h"
 #include "core/data_type/data_type_map.h"
@@ -60,6 +61,12 @@ bool contains_variant_type(const DataTypePtr& input) {
     }
 }
 
+// The most file scanners one instance of a file scan node runs: the 
max_file_scanners_concurrency
+// session variable, 16 when it is not set.
+int file_scanners_per_instance(RuntimeState* state) {
+    return state->max_file_scanners_concurrency() > 0 ? 
state->max_file_scanners_concurrency() : 16;
+}
+
 } // namespace
 
 PushDownType FileScanLocalState::_should_push_down_binary_predicate(
@@ -107,8 +114,7 @@ int 
FileScanLocalState::max_scanners_concurrency(RuntimeState* state) const {
      *
      * If this is a serial operator, the max concurrency should multiply by 
the number of parallel instances of the operator.
      */
-    return (state->max_file_scanners_concurrency() > 0 ? 
state->max_file_scanners_concurrency()
-                                                       : 16) *
+    return file_scanners_per_instance(state) *
            (state->query_parallel_instance_num() / 
_parent->parallelism(state));
 }
 
@@ -198,12 +204,7 @@ Status 
FileScanLocalState::_init_scanners(std::list<ScannerSPtr>* scanners) {
     }
 
     auto& p = _parent->cast<FileScanOperatorX>();
-    // There's only one scan range for each backend in batch split mode. Each 
backend only starts up one ScanNode instance.
-    uint32_t shard_num =
-            std::min(ScannerScheduler::default_remote_scan_thread_num() / 
p.parallelism(state()),
-                     _max_scanners);
-    shard_num = std::max(shard_num, 1U);
-    _kv_cache = std::make_unique<ShardedKVCache>(shard_num);
+    DORIS_CHECK(p._kv_cache != nullptr);
     const TFileScanRangeParams* scan_params = nullptr;
     if (state()->get_query_ctx() != nullptr &&
         
state()->get_query_ctx()->file_scan_range_params_map.count(parent_id()) > 0) {
@@ -235,11 +236,11 @@ Status 
FileScanLocalState::_init_scanners(std::list<ScannerSPtr>* scanners) {
         ScannerSPtr scanner;
         if (use_file_scanner_v2) {
             scanner = FileScannerV2::create_shared(state(), this, p._limit, 
_split_source,
-                                                   _scanner_profile.get(), 
_kv_cache.get(),
+                                                   _scanner_profile.get(), 
p._kv_cache.get(),
                                                    &p._colname_to_slot_id);
         } else {
             scanner = FileScanner::create_shared(state(), this, p._limit, 
_split_source,
-                                                 _scanner_profile.get(), 
_kv_cache.get(),
+                                                 _scanner_profile.get(), 
p._kv_cache.get(),
                                                  &p._colname_to_slot_id);
         }
         RETURN_IF_ERROR(scanner->init(state(), _conjuncts));
@@ -327,6 +328,18 @@ 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). Counted in 64 
bits: the scanners per
+    // instance are a session variable with no upper bound, and in an int the 
product overflows
+    // before std::min can cap it, leaving the whole node one shard to load 
under.
+    const int64_t shard_num =
+            
std::min<int64_t>(ScannerScheduler::default_remote_scan_thread_num(),
+                              
static_cast<int64_t>(state->query_parallel_instance_num()) *
+                                      file_scanners_per_instance(state));
+    _kv_cache =
+            
std::make_unique<ShardedKVCache>(cast_set<uint32_t>(std::max<int64_t>(shard_num,
 1)));
     if (state->get_query_ctx() != nullptr &&
         
state->get_query_ctx()->file_scan_range_params_map.contains(node_id())) {
         TFileScanRangeParams& params =
diff --git a/be/src/exec/operator/file_scan_operator.h 
b/be/src/exec/operator/file_scan_operator.h
index 45901d82858..5250ee35ec3 100644
--- a/be/src/exec/operator/file_scan_operator.h
+++ b/be/src/exec/operator/file_scan_operator.h
@@ -86,12 +86,6 @@ private:
             const std::set<std::string> fn_name) const override;
     std::shared_ptr<SplitSourceConnector> _split_source = nullptr;
     int _max_scanners;
-    // A in memory cache to save some common components
-    // of the this scan node. eg:
-    // 1. iceberg delete file
-    // 2. parquet file meta
-    // KVCache<std::string> _kv_cache;
-    std::unique_ptr<ShardedKVCache> _kv_cache;
     TupleId _output_tuple_id = -1;
 };
 
@@ -119,6 +113,11 @@ private:
 
     const std::string _table_name;
     bool _batch_split_mode = false;
+    // What this scan node's readers parse once and reuse across splits: 
iceberg delete files,
+    // deletion vectors, the keys of a fluss union read's log tails. One cache 
for every instance of
+    // the node on this backend, so a split reuses what a split in another 
instance already read;
+    // splits are dealt to instances with no regard for which delete file or 
bucket they share.
+    std::unique_ptr<ShardedKVCache> _kv_cache;
 };
 
 /// Instantiated once in scan_operator.cpp; suppresses per-TU implicit 
instantiation.
diff --git a/be/test/exec/operator/file_scan_operator_fluss_test.cpp 
b/be/test/exec/operator/file_scan_operator_fluss_test.cpp
index e73a7319bad..751c1e8ec49 100644
--- a/be/test/exec/operator/file_scan_operator_fluss_test.cpp
+++ b/be/test/exec/operator/file_scan_operator_fluss_test.cpp
@@ -17,15 +17,52 @@
 
 #include <gtest/gtest.h>
 
+#include <cstdint>
+#include <limits>
+#include <list>
+#include <map>
+#include <memory>
 #include <string>
+#include <utility>
+#include <vector>
 
+#include "common/object_pool.h"
 #include "exec/operator/file_scan_operator.h"
+#include "exec/pipeline/dependency.h"
 #include "exec/scan/file_scanner_v2.h"
+#include "exec/scan/scanner.h"
+#include "exec/scan/scanner_scheduler.h"
 #include "gen_cpp/PlanNodes_types.h"
+#include "runtime/descriptors.h"
+#include "testutil/mock/mock_runtime_state.h"
 
 namespace doris::pipeline {
 namespace {
 
+// One tuple with no slots: all a file scan node needs to be initialized and 
prepared.
+Status create_one_tuple_descriptors(ObjectPool* pool, DescriptorTbl** 
descriptors) {
+    TDescriptorTable thrift_descriptors;
+    TTupleDescriptor tuple_descriptor;
+    tuple_descriptor.id = 0;
+    tuple_descriptor.byteSize = 0;
+    tuple_descriptor.numNullBytes = 0;
+    thrift_descriptors.tupleDescriptors.push_back(tuple_descriptor);
+    return DescriptorTbl::create(pool, thrift_descriptors, descriptors);
+}
+
+// A file scan node over that tuple.
+TPlanNode file_scan_plan_node() {
+    TPlanNode plan_node;
+    plan_node.node_id = 0;
+    plan_node.node_type = TPlanNodeType::FILE_SCAN_NODE;
+    plan_node.num_children = 0;
+    plan_node.limit = -1;
+    plan_node.row_tuples.push_back(0);
+    plan_node.file_scan_node.tuple_id = 0;
+    plan_node.__isset.file_scan_node = true;
+    return plan_node;
+}
+
 TFileScanRangeParams scan_level_format(const std::string& table_format) {
     TFileScanRangeParams params;
     if (!table_format.empty()) {
@@ -115,4 +152,75 @@ TEST(FileScanOperatorFlussTest, 
ScannerV2SupportsEveryRangeKindOfAForcedNode) {
     EXPECT_TRUE(FileScannerV2::is_supported(params, lake_native));
 }
 
+// Every instance of a file scan node on a backend reads through the node's 
one cache. A primary-key
+// union read keeps the keys of each bucket's log tail there, as iceberg keeps 
its parsed delete
+// files, and splits are dealt to instances with no regard for the bucket or 
delete file they share.
+// With a cache per instance, each instance holding a split of a bucket read 
that bucket's tail again
+// over JNI: 69 reads where 16 would do on a 16-bucket table spread over 9 
instances.
+TEST(FileScanOperatorFlussTest, 
EveryInstanceOfTheNodeReadsThroughTheNodesCache) {
+    ObjectPool pool;
+    DescriptorTbl* descriptors = nullptr;
+    ASSERT_TRUE(create_one_tuple_descriptors(&pool, &descriptors).ok());
+    MockRuntimeState state;
+    state.set_desc_tbl(descriptors);
+
+    const TPlanNode plan_node = file_scan_plan_node();
+    FileScanOperatorX node(&pool, plan_node, 0, *descriptors, 
/*parallel_tasks=*/1);
+    ASSERT_TRUE(node.init(plan_node, &state).ok());
+    ASSERT_TRUE(node.prepare(&state).ok());
+    ASSERT_NE(node._kv_cache, nullptr);
+
+    TFileRangeDesc log_range = jni_range("fluss");
+    log_range.table_format_params.__set_fluss_params({{"fluss.range_type", 
"LOG"}});
+    TFileScanRange file_scan_range;
+    file_scan_range.__set_params(scan_level_format("fluss"));
+    // A query reads no source tuple; tuple 0 there would make this scan a 
load.
+    file_scan_range.params.__set_src_tuple_id(1);
+    file_scan_range.ranges.push_back(log_range);
+    TScanRangeParams scan_range;
+    
scan_range.scan_range.ext_scan_range.__set_file_scan_range(file_scan_range);
+    scan_range.scan_range.__isset.ext_scan_range = true;
+    const std::vector<TScanRangeParams> scan_ranges {scan_range};
+
+    RuntimeProfile profile("EveryInstanceOfTheNodeReadsThroughTheNodesCache");
+    const std::map<int, std::pair<std::shared_ptr<BasicSharedState>,
+                                  std::vector<std::shared_ptr<Dependency>>>>
+            shared_state_map;
+    for (int instance = 0; instance < 2; ++instance) {
+        auto local_state = FileScanLocalState::create_unique(&state, &node);
+        LocalStateInfo info {&profile, scan_ranges, nullptr, shared_state_map, 
instance};
+        ASSERT_TRUE(local_state->init(&state, info).ok());
+        std::list<ScannerSPtr> scanners;
+        ASSERT_TRUE(local_state->_init_scanners(&scanners).ok());
+        ASSERT_FALSE(scanners.empty());
+        for (const auto& scanner : scanners) {
+            const auto* file_scanner = 
dynamic_cast<FileScannerV2*>(scanner.get());
+            ASSERT_NE(file_scanner, nullptr);
+            EXPECT_EQ(file_scanner->_kv_cache, node._kv_cache.get()) << 
"instance " << instance;
+        }
+    }
+}
+
+// The node's cache is sharded for the scanners of all its instances: the 
instance count times
+// max_file_scanners_concurrency, a session variable with no upper bound. A 
user may set it to
+// INT32_MAX to mean "no limit", and with two instances that product overflows 
an int, which would
+// leave the node a single shard, behind whose lock every scanner of every 
instance loads its delete
+// files one at a time. The count has to come out capped by the remote scan 
threads instead.
+TEST(FileScanOperatorFlussTest, 
CacheShardsOfAnUnlimitedScannerSettingAreCappedNotOverflowed) {
+    ObjectPool pool;
+    DescriptorTbl* descriptors = nullptr;
+    ASSERT_TRUE(create_one_tuple_descriptors(&pool, &descriptors).ok());
+    MockRuntimeState state;
+    state.set_desc_tbl(descriptors);
+    state._query_options.__set_parallel_instance(2);
+    
state._query_options.__set_max_file_scanners_concurrency(std::numeric_limits<int32_t>::max());
+
+    const TPlanNode plan_node = file_scan_plan_node();
+    FileScanOperatorX node(&pool, plan_node, 0, *descriptors, 
/*parallel_tasks=*/2);
+    ASSERT_TRUE(node.init(plan_node, &state).ok());
+    ASSERT_TRUE(node.prepare(&state).ok());
+    EXPECT_EQ(node._kv_cache->_num_shards,
+              
static_cast<uint32_t>(ScannerScheduler::default_remote_scan_thread_num()));
+}
+
 } // namespace doris::pipeline
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionCloser.java
 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionCloser.java
new file mode 100644
index 00000000000..6d0a03fbf78
--- /dev/null
+++ 
b/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) {
+            LOG.warn("Failed to close a fluss connection", e);
+        }
+    }
+}
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java
 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java
new file mode 100644
index 00000000000..5e4aa46194f
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java
@@ -0,0 +1,323 @@
+// 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.Deque;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.LinkedList;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+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, 
FlussConnectionPool::startReaper);
+
+    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;
+    /** Starts the thread that runs the sweep it is given on this pool: {@link 
#startReaper}. */
+    private final Consumer<Runnable> reaperStarter;
+    /** Set while the reaper is being started and once it has; see {@link 
#startReaperOnce}. */
+    private final AtomicBoolean reaperStarted = new AtomicBoolean();
+
+    /**
+     * 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} and the reaper rely
+     * on: {@link #giveBack} puts an entry in only once it holds its 
connection.
+     */
+    private final Map<Map<String, String>, Deque<Idle>> idle = new HashMap<>();
+
+    FlussConnectionPool(Function<Configuration, Connection> factory, 
LongSupplier nanoTime,
+            Consumer<Connection> closer, Consumer<Runnable> reaperStarter) {
+        this.factory = factory;
+        this.nanoTime = nanoTime;
+        this.closer = closer;
+        this.reaperStarter = reaperStarter;
+    }
+
+    /**
+     * 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) {
+            Deque<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.
+     *
+     * <p>An {@code OutOfMemoryError} here may cost this connection, never the 
ones the pool holds: wherever
+     * it strikes, the pool is left whole. An entry goes into {@link #idle} 
already holding its connection -
+     * put in empty and filled after, an error in between would leave an entry 
that every later
+     * {@link #borrow} and every sweep take a connection from and find none. 
And an entry is a
+     * {@link LinkedList}, which allocates a node before it links it in: an 
{@code ArrayDeque} stores an
+     * element first and grows its array after, and an error in that growth 
leaves it looking empty while it
+     * holds every connection - the next one given back overwrites the oldest, 
and the rest are out of reach
+     * of {@link #borrow} and of the reaper for good.
+     */
+    void giveBack(Lease lease) {
+        Idle returned = new Idle(lease.connection, nanoTime.getAsLong());
+        synchronized (idle) {
+            Deque<Idle> connections = idle.get(lease.key);
+            if (connections == null) {
+                Deque<Idle> first = new LinkedList<>();
+                first.addLast(returned);
+                idle.put(lease.key, first);
+            } else {
+                connections.addLast(returned);
+            }
+        }
+        startReaperOnce();
+    }
+
+    /**
+     * Starts the reaper with the first connection given back, and with a 
later one again if starting it
+     * failed. Not while the class initializes: starting a thread fails with 
an {@code OutOfMemoryError}
+     * while BE's JVM heap is full or no native thread is to be had, and an 
error in a static initializer
+     * leaves the class unusable in its classloader, which BE keeps for the 
life of the process - every
+     * fluss read of the BE would fail with {@code NoClassDefFoundError} until 
it restarted. Until the reaper
+     * runs, ranges go on borrowing and giving back; only idle connections 
wait for it.
+     */
+    private void startReaperOnce() {
+        if (reaperStarted.get() || !reaperStarted.compareAndSet(false, true)) {
+            return;
+        }
+        try {
+            reaperStarter.accept(() -> sweep(this));
+        } catch (Throwable startFailure) {
+            reaperStarted.set(false);
+            try {
+                LOG.warn("Failed to start the thread that closes idle fluss 
connections; the next connection"
+                        + " given back tries again", startFailure);
+            } catch (Throwable logFailure) {
+                // Logging allocates too, and the heap may still be full; the 
next connection given back
+                // tries again regardless.
+            }
+        }
+    }
+
+    /**
+     * 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) {
+        closer.accept(lease.connection);
+    }
+
+    /**
+     * Closes, without waiting for them, the connections nobody has borrowed 
for {@code idleNanos}.
+     *
+     * <p>One at a time: a connection leaves the pool only to be handed to the 
closer at once. Nothing else
+     * holds it then, so an {@code OutOfMemoryError} between taking it and 
handing it over would leave it
+     * open, threads and sockets, for the life of the process - and while BE's 
JVM heap is full, the error
+     * strikes whatever allocates. Taken out together into a list and handed 
over after a log line, every
+     * expired connection was out of the pool while the list grew, the line 
was logged and each close was
+     * submitted, which is where nearly all of a sweep's allocation is, and an 
error there lost all of them.
+     */
+    void closeIdleLongerThan(long idleNanos) {
+        long now = nanoTime.getAsLong();
+        int closed = 0;
+        for (Connection connection = takeExpired(now, idleNanos); connection 
!= null;
+                connection = takeExpired(now, idleNanos)) {
+            // Outside the lock: a close the closer hands back to this thread 
takes two seconds.
+            closer.accept(connection);
+            closed++;
+        }
+        if (closed > 0) {
+            LOG.info("Closing {} fluss connections nobody has borrowed for {} 
seconds", closed,
+                    TimeUnit.NANOSECONDS.toSeconds(idleNanos));
+        }
+    }
+
+    /**
+     * Takes one connection idle for {@code idleNanos} at {@code now} out of 
the pool, or returns null if
+     * none is. Nothing may allocate between unlinking it and returning it 
({@link #closeIdleLongerThan}):
+     * the list only unlinks it, and the iterator removes an emptied entry by 
the hash the map stored
+     * rather than hashing the configuration again.
+     */
+    private Connection takeExpired(long now, long idleNanos) {
+        synchronized (idle) {
+            Iterator<Deque<Idle>> keys = idle.values().iterator();
+            while (keys.hasNext()) {
+                Deque<Idle> connections = keys.next();
+                // An entry is in the order its connections came back, so an 
expired one is at its head.
+                if (now - connections.peekFirst().since >= idleNanos) {
+                    Connection connection = connections.pollFirst().connection;
+                    if (connections.isEmpty()) {
+                        keys.remove();
+                    }
+                    return connection;
+                }
+            }
+            return null;
+        }
+    }
+
+    /**
+     * One run of the reaper. Nothing may escape it: a scheduled task that 
throws is never run again, and
+     * idle connections would then keep their threads for the life of the 
process. That holds for an
+     * {@code Error} as much as for an exception - above all the {@code 
OutOfMemoryError} of a scan that
+     * filled BE's JVM heap, which strikes whatever allocates while the heap 
is full. Caught as an
+     * exception only, one such error ended the sweeps of a BE for good, and 
the 128 connections idle at
+     * that moment stayed open, threads and sockets, until the BE restarted.
+     */
+    static void sweep(FlussConnectionPool pool) {
+        try {
+            pool.closeIdleLongerThan(IDLE_TIMEOUT_NANOS);
+        } catch (Throwable t) {
+            try {
+                LOG.warn("Failed to close idle fluss connections", t);
+            } catch (Throwable logFailure) {
+                // Logging allocates too, and the heap may still be full; the 
next sweep comes regardless.
+            }
+        }
+    }
+
+    /**
+     * Runs {@code sweep} every {@link #REAP_INTERVAL_SECONDS} on a daemon 
thread of its own. A new executor
+     * on every call, never one kept from a call that failed: that one already 
holds the sweep, queued
+     * before the thread to run it failed to start, and a later start of its 
thread would run it twice.
+     */
+    private static void startReaper(Runnable sweep) {
+        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, REAP_INTERVAL_SECONDS, 
REAP_INTERVAL_SECONDS, TimeUnit.SECONDS);
+    }
+
+    int idleCount() {
+        synchronized (idle) {
+            int count = 0;
+            for (Deque<Idle> connections : idle.values()) {
+                count += connections.size();
+            }
+            return count;
+        }
+    }
+
+    /** A connection lent to one range, and what it has to be given back 
under. */
+    static final class Lease {
+        private final Map<String, String> key;
+        private final Connection connection;
+        private final boolean opened;
+
+        private Lease(Map<String, String> key, Connection connection, boolean 
opened) {
+            this.key = key;
+            this.connection = connection;
+            this.opened = opened;
+        }
+
+        Connection connection() {
+            return connection;
+        }
+
+        /** Whether the range had to open this connection, rather than borrow 
one that was idle. */
+        boolean opened() {
+            return opened;
+        }
+    }
+
+    private static final class Idle {
+        private final Connection connection;
+        private final long since;
+
+        private Idle(Connection connection, long since) {
+            this.connection = connection;
+            this.since = since;
+        }
+    }
+}
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java
 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java
index dc049771128..f7cbcd8c9e8 100644
--- 
a/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java
+++ 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java
@@ -23,7 +23,6 @@ import org.apache.doris.jni.toolkit.vec.JniSchemaParams;
 import org.apache.doris.jni.toolkit.vec.NestedProjection;
 
 import org.apache.fluss.client.Connection;
-import org.apache.fluss.client.ConnectionFactory;
 import org.apache.fluss.client.FlussConnection;
 import org.apache.fluss.client.table.Table;
 import org.apache.fluss.client.table.scanner.batch.BatchScanner;
@@ -46,6 +45,7 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.TimeZone;
+import java.util.concurrent.CompletableFuture;
 
 /**
  * Reads one fluss scan range — one bucket of one partition, bounded by log 
offsets.
@@ -116,11 +116,20 @@ public class FlussJniScanner extends JniScanner {
     /** {@code null} on an unpartitioned table, which fluss subscribes to by 
bucket alone. */
     private final Long partitionId;
 
+    /** Borrowed from {@link FlussConnectionPool} for this range, and given 
back when it is closed. */
+    private FlussConnectionPool.Lease lease;
+    /** Whether no connection was idle, so this range opened the one it 
borrowed. */
+    private boolean openedConnection;
     private Connection connection;
     private Table table;
     private BatchScanner scanner;
     /** The same object as {@link #scanner} on a {@code PK_TAIL} range, for 
what it counted. */
     private PkTailBatchScanner tailScanner;
+    /**
+     * The same object as {@link #scanner} on a {@code PK_FULL} range, for 
when it stops using the
+     * connection; see {@link #closeInternal()}.
+     */
+    private SafeKvSnapshotAndLogBatchScanner fullScanner;
 
     /** Fluss types of the projected columns, positionally aligned with {@link 
#fields}. */
     private List<DataType> projectedTypes;
@@ -194,9 +203,12 @@ public class FlussJniScanner extends JniScanner {
         // The fluss client spawns its own threads (netty IO, metadata 
updater) while connecting, and a
         // thread inherits the context classloader of whoever created it. 
Started under BE's loader they
         // would not see fluss at all - which is why JniScanner.open() 
installs this plugin's loader as
-        // the context classloader around this call and restores the caller's 
on the way out.
+        // the context classloader around this call and restores the caller's 
on the way out. A borrowed
+        // connection was opened under it too, by whichever range opened it 
first.
         try {
-            connection = ConnectionFactory.createConnection(clientConfig());
+            lease = FlussConnectionPool.INSTANCE.borrow(clientConfig());
+            openedConnection = lease.opened();
+            connection = lease.connection();
             table = connection.getTable(TablePath.of(required(DB_NAME), 
required(TABLE_NAME)));
 
             RowType rowType = table.getTableInfo().getRowType();
@@ -263,8 +275,9 @@ public class FlussJniScanner extends JniScanner {
      * snapshot is far behind costs the most.
      */
     private BatchScanner primaryKeyScanner(TableBucket tableBucket, int[] 
projection) {
-        return new SafeKvSnapshotAndLogBatchScanner(
+        fullScanner = new SafeKvSnapshotAndLogBatchScanner(
                 table, tableBucket, kvSnapshotId, logStartOffset, 
logStopOffset, projection);
+        return fullScanner;
     }
 
     /**
@@ -306,6 +319,14 @@ public class FlussJniScanner extends JniScanner {
                 
config.setString(entry.getKey().substring(CLIENT_PREFIX.length()), 
entry.getValue());
             }
         }
+        if (!config.contains(ConfigOptions.NETTY_CLIENT_NUM_NETWORK_THREADS)) {
+            // A connection serves one range at a time (FlussConnectionPool): 
one stream of requests,
+            // which one network thread carries. Fluss's default is four per 
connection - four selectors
+            // opened up front, and a thread for each server the connection 
reaches - while a scan reads
+            // up to sixteen ranges at once, each on a connection of its own, 
so eight concurrent scans
+            // held over a hundred connections. A catalog that sets the option 
keeps its number.
+            config.set(ConfigOptions.NETTY_CLIENT_NUM_NETWORK_THREADS, 1);
+        }
         if (!RANGE_TYPE_PK_FULL.equals(rangeType)) {
             // Local-first mistakes lake-covered offsets for readable local 
log. Remote-first first
             // tries a retained remote segment, then falls back to local; a 
missing segment is
@@ -348,14 +369,39 @@ public class FlussJniScanner extends JniScanner {
     @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 {
+                // A primary-key range closed before its kv snapshot arrived: 
fluss goes on copying the
+                // snapshot on the connection's download threads until the 
reader can be closed. Closed
+                // under that copy, the connection would drop the files still 
queued, and the reader would
+                // wait for them forever - with the pool thread it runs on, 
the thread waiting to close it,
+                // and the half-copied snapshot directory. So the connection 
waits for the reader instead.
+                // handBack is a discard here, since a range whose snapshot 
never arrived was not read to
+                // its end, and a connection no closer thread can take is 
closed on the thread completing
+                // released (FlussConnectionCloser#close): the future thenRun 
returns has nothing to report.
+                released.thenRun(handBack);
+            }
+            lease = null;
+            connection = null;
+        }
         currentBatch = null;
         if (failure != null) {
             throw failure;
@@ -385,6 +431,8 @@ public class FlussJniScanner extends JniScanner {
         Map<String, String> statistics = new HashMap<>();
         statistics.put("counter:FlussJniRowsRead", String.valueOf(rowsRead));
         statistics.put("gauge:FlussJniRequiredFieldCount", 
String.valueOf(fields.length));
+        // Summed over a query's ranges: the connections it had to open 
because none was idle.
+        statistics.put("counter:FlussJniConnectionsOpened", openedConnection ? 
"1" : "0");
         if (tailScanner != null) {
             // What the tail cost and what it hid: the records replayed, and 
the keys it ended deleted —
             // those are lake rows that disappear with nothing returned in 
their place.
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScanner.java
 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScanner.java
index c74946bea05..d8e44fbb3f9 100644
--- 
a/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScanner.java
+++ 
b/fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScanner.java
@@ -47,7 +47,9 @@ import java.util.List;
 import java.util.Map;
 import java.util.NoSuchElementException;
 import java.util.TreeMap;
+import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.locks.LockSupport;
 import java.util.stream.Collectors;
 import java.util.stream.IntStream;
 import javax.annotation.Nullable;
@@ -69,7 +71,7 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
     private final Comparator<InternalRow> primaryKeyComparator;
     private final Map<InternalRow, KeyValueRow> logRows;
 
-    @Nullable private final BatchScanner snapshotScanner;
+    @Nullable private final PublicationGuardedBatchScanner snapshotScanner;
     @Nullable private final LogScanner logScanner;
 
     private boolean logScanFinished;
@@ -190,6 +192,17 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
         IOUtils.closeQuietly(logScanner);
     }
 
+    /**
+     * Completes once this reader has stopped using the connection it was 
created on. For the log reader
+     * that is {@link #close()}; the snapshot reader, closed before its 
snapshot arrived, goes on copying
+     * the snapshot on the connection's download threads until {@link 
PublicationGuardedBatchScanner} can
+     * close it. The connection must stay open until then: shut down, fluss's 
download pool drops the
+     * files still queued, and the snapshot reader waits for them forever.
+     */
+    CompletableFuture<Void> released() {
+        return snapshotScanner == null ? 
CompletableFuture.completedFuture(null) : snapshotScanner.released();
+    }
+
     private static ProjectionPlan createProjectionPlan(
             TableInfo tableInfo, @Nullable int[] projectedFields) {
         return ProjectionPlan.create(
@@ -286,6 +299,10 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
      * null} means ready-but-empty; a non-empty iterator means 
ready-with-data). An early close starts
      * one daemon waiter only for that cancelled scanner, observes the same 
publication boundary, and
      * then performs the delegate's first and only close.</p>
+     *
+     * <p>Until then the initializer is still copying the snapshot on the 
download threads of the
+     * connection the delegate was created on, so that connection has to stay 
open; {@link #released()}
+     * says when it may close.</p>
      */
     static final class PublicationGuardedBatchScanner implements BatchScanner {
         private static final Duration PUBLICATION_POLL = 
Duration.ofMillis(100);
@@ -294,6 +311,8 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
         private final AtomicBoolean closed = new AtomicBoolean();
         private final AtomicBoolean publicationObserved = new AtomicBoolean();
         private final AtomicBoolean delegateClosed = new AtomicBoolean();
+        /** Completed once the delegate has been closed, by whichever path 
closed it. */
+        private final CompletableFuture<Void> released = new 
CompletableFuture<>();
 
         PublicationGuardedBatchScanner(BatchScanner delegate) {
             this.delegate = delegate;
@@ -323,14 +342,17 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
             // Only early cancellation needs a waiter. A dedicated daemon 
avoids deadlocking the
             // ForkJoin common pool that Fluss 1.0 also uses for its 
initializer, while normal scans
             // create no extra thread at all.
-            Thread closeAfterPublication = new Thread(
-                    this::awaitPublicationAndClose, 
"fluss-snapshot-publication-close");
-            closeAfterPublication.setDaemon(true);
             try {
+                Thread closeAfterPublication = new Thread(
+                        this::awaitPublicationAndClose, 
"fluss-snapshot-publication-close");
+                closeAfterPublication.setDaemon(true);
                 closeAfterPublication.start();
             } catch (RuntimeException | Error startFailure) {
-                // Losing the waiter would recreate the native leak. Fall back 
to waiting on this
-                // cancellation thread; publication/failure is the only safe 
point for the SDK close.
+                // Losing the waiter would recreate the native leak, and keep 
the connection, which
+                // FlussJniScanner hands back only once released() completes, 
open for good. A thread
+                // that could not be built - an OutOfMemoryError while the 
heap is full - is as lost as
+                // one that could not start. Fall back to waiting on this 
cancellation thread;
+                // publication/failure is the only safe point for the SDK 
close.
                 awaitPublicationAndClose();
                 throw startFailure;
             }
@@ -348,6 +370,13 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
                     // registry has been closed, so the delegate is safe to 
close at this point too.
                     closeDelegateQuietly();
                     return;
+                } catch (OutOfMemoryError pollFailure) {
+                    // This thread allocated while BE's JVM heap was full, as 
it may well be: a range is
+                    // often closed early because its query filled the heap. 
That is no initialization
+                    // failure - those arrive above, as exceptions - so the 
delegate must not be closed
+                    // yet, and nothing but this thread will close it, nor 
hand back the connection that
+                    // waits for it. Keep waiting; the failed query's memory 
comes back as it closes.
+                    LockSupport.parkNanos(PUBLICATION_POLL.toNanos());
                 }
             }
         }
@@ -365,15 +394,25 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
             return observed;
         }
 
+        /** Completes once the delegate has been closed: from then on it no 
longer uses its connection. */
+        CompletableFuture<Void> released() {
+            return released;
+        }
+
         private void closeDelegate() throws IOException {
             if (delegateClosed.compareAndSet(false, true)) {
-                delegate.close();
+                try {
+                    delegate.close();
+                } finally {
+                    released.complete(null);
+                }
             }
         }
 
         private void closeDelegateQuietly() {
             if (delegateClosed.compareAndSet(false, true)) {
                 IOUtils.closeQuietly(delegate);
+                released.complete(null);
             }
         }
     }
@@ -386,7 +425,7 @@ final class SafeKvSnapshotAndLogBatchScanner implements 
BatchScanner {
     }
 
     static final class ScannerResources {
-        @Nullable BatchScanner snapshotScanner;
+        @Nullable PublicationGuardedBatchScanner snapshotScanner;
         @Nullable LogScanner logScanner;
 
         private ScannerResources() {
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussConnectionCloserTest.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussConnectionCloserTest.java
new file mode 100644
index 00000000000..36cff8a4aab
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussConnectionCloserTest.java
@@ -0,0 +1,106 @@
+// 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.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+/**
+ * What the closer promises whoever hands it a connection: it does not wait 
for the connection to close,
+ * and connections do not pile up behind it. Every step here waits on a latch 
rather than on the clock;
+ * the clock only bounds how long a closer that never does its part is waited 
for.
+ */
+public class FlussConnectionCloserTest {
+
+    /** Far longer than any step takes; reached only when a close never 
starts. */
+    private static final long PATIENCE_SECONDS = 60;
+
+    @Test
+    public void theCallerDoesNotWaitForItsConnectionToClose() throws Exception 
{
+        Thread caller = Thread.currentThread();
+        CountDownLatch mayFinish = new CountDownLatch(1);
+        try {
+            // A connection closed off this thread stays blocked in close() 
until the end of this
+            // test, so close() below can only return if something else is 
running it. The test
+            // before this one leaves every closer thread on its way back from 
a connection, and
+            // until one of them is free again the caller closes its own - 
which is what the closer
+            // promises, and says nothing about this test. So hand connections 
over until one is
+            // taken; only a closer that never takes any runs into the 
deadline.
+            long deadline = System.nanoTime() + 
TimeUnit.SECONDS.toNanos(PATIENCE_SECONDS);
+            Thread closedBy = caller;
+            while (closedBy == caller) {
+                Assertions.assertTrue(System.nanoTime() < deadline,
+                        "every connection was closed on the thread that handed 
it over");
+                CountDownLatch closing = new CountDownLatch(1);
+                AtomicReference<Thread> closer = new AtomicReference<>();
+                FlussConnectionCloser.close(new StubConnection(() -> {
+                    closer.set(Thread.currentThread());
+                    closing.countDown();
+                    if (Thread.currentThread() != caller) {
+                        mayFinish.await();
+                    }
+                }));
+                Assertions.assertTrue(closing.await(PATIENCE_SECONDS, 
TimeUnit.SECONDS),
+                        "the connection was never closed");
+                closedBy = closer.get();
+            }
+        } finally {
+            mayFinish.countDown();
+        }
+    }
+
+    @Test
+    public void onceTooManyAreClosingTheCallerClosesItsOwn() throws Exception {
+        Thread caller = Thread.currentThread();
+        CountDownLatch mayFinish = new CountDownLatch(1);
+        try {
+            // Hold one closer thread per connection until the closer has none 
left. Connections of
+            // other tests may be closing as well, so that can happen before 
MAX_CLOSING of these are
+            // in - but never after.
+            boolean closedByCaller = false;
+            for (int i = 0; i <= FlussConnectionCloser.MAX_CLOSING && 
!closedByCaller; i++) {
+                CountDownLatch closing = new CountDownLatch(1);
+                AtomicReference<Thread> closedBy = new AtomicReference<>();
+                FlussConnectionCloser.close(new StubConnection(() -> {
+                    closedBy.set(Thread.currentThread());
+                    closing.countDown();
+                    if (Thread.currentThread() != caller) {
+                        mayFinish.await();
+                    }
+                }));
+                Assertions.assertTrue(closing.await(PATIENCE_SECONDS, 
TimeUnit.SECONDS),
+                        "connection " + i + " was never closed");
+                closedByCaller = closedBy.get() == caller;
+            }
+            Assertions.assertTrue(closedByCaller,
+                    "more than " + FlussConnectionCloser.MAX_CLOSING + " 
connections were closing at once");
+
+            // A connection that fails to close must not fail whoever handed 
it over, whichever thread
+            // ends up closing it.
+            Assertions.assertDoesNotThrow(() -> 
FlussConnectionCloser.close(new StubConnection(() -> {
+                throw new IllegalStateException("this connection refuses to 
close");
+            })));
+        } finally {
+            mayFinish.countDown();
+        }
+    }
+}
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussConnectionPoolTest.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussConnectionPoolTest.java
new file mode 100644
index 00000000000..5acfd1a1c86
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussConnectionPoolTest.java
@@ -0,0 +1,298 @@
+// 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.config.Configuration;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Consumer;
+
+/**
+ * What the pool promises a scan range: it borrows the connection a range 
before it gave back, never one
+ * another range is still reading through, never one opened with other 
settings; and what nobody borrows
+ * is closed once it has been idle long enough. The connections are stubs and 
the clock is the test's, so
+ * none of this waits for a cluster or for a minute to pass.
+ */
+public class FlussConnectionPoolTest {
+
+    private static final long SECOND = TimeUnit.SECONDS.toNanos(1);
+
+    /** Far longer than a close handed to the closer takes to start. */
+    private static final long PATIENCE_SECONDS = 60;
+
+    /** The tests run each sweep themselves, on their own clock; no reaper 
thread runs one. */
+    private static final Consumer<Runnable> NO_REAPER = sweep -> { };
+
+    private final AtomicLong now = new AtomicLong();
+    /** One per connection opened, counted down when it is closed. Opened 
outside the pool's lock. */
+    private final List<CountDownLatch> closes = 
Collections.synchronizedList(new ArrayList<>());
+    private final FlussConnectionPool pool =
+            new FlussConnectionPool(config -> open(), now::get, 
FlussConnectionCloser::close, NO_REAPER);
+
+    @Test
+    public void rangeBorrowsTheConnectionTheRangeBeforeItGaveBack() {
+        FlussConnectionPool.Lease first = pool.borrow(config("server-a:9123"));
+        Assertions.assertTrue(first.opened(), "the first range found nothing 
idle");
+        pool.giveBack(first);
+
+        FlussConnectionPool.Lease second = 
pool.borrow(config("server-a:9123"));
+        Assertions.assertSame(first.connection(), second.connection());
+        Assertions.assertFalse(second.opened(), "the second range opened a 
connection of its own");
+        Assertions.assertEquals(1, closes.size(), "connections opened");
+    }
+
+    /**
+     * A connection serves one range at a time: its download threads are what 
a primary-key range copies
+     * its kv snapshot with, and ranges reading at the same time must not 
queue for them.
+     */
+    @Test
+    public void rangesReadingAtTheSameTimeEachHaveAConnectionOfTheirOwn() {
+        FlussConnectionPool.Lease first = pool.borrow(config("server-a:9123"));
+        FlussConnectionPool.Lease second = 
pool.borrow(config("server-a:9123"));
+        Assertions.assertNotSame(first.connection(), second.connection());
+        Assertions.assertTrue(second.opened());
+
+        pool.giveBack(first);
+        pool.giveBack(second);
+        Assertions.assertEquals(2, pool.idleCount());
+        // The one given back last is lent first, so that after a burst the 
ones it no longer needs are
+        // the ones that stay idle long enough to be closed.
+        Assertions.assertSame(second.connection(), 
pool.borrow(config("server-a:9123")).connection());
+    }
+
+    /** Servers, credentials and client options are part of what a connection 
is. */
+    @Test
+    public void rangesWithOtherSettingsNeverBorrowTheConnection() {
+        FlussConnectionPool.Lease first = pool.borrow(config("server-a:9123"));
+        pool.giveBack(first);
+
+        Configuration otherServer = config("server-b:9123");
+        Assertions.assertTrue(pool.borrow(otherServer).opened());
+        Configuration otherOption = config("server-a:9123");
+        otherOption.setString("client.scanner.log.read-preference", 
"REMOTE_FIRST");
+        Assertions.assertTrue(pool.borrow(otherOption).opened());
+
+        Assertions.assertSame(first.connection(), 
pool.borrow(config("server-a:9123")).connection());
+    }
+
+    /**
+     * Ranges borrow and give back from every scanner thread at once. None may 
ever hold a connection
+     * another one holds, and the pool opens no more than were held at once.
+     */
+    @Test
+    public void connectionIsNeverLentToTwoRangesAtOnce() throws Exception {
+        int threads = 8;
+        int rounds = 2000;
+        Set<Connection> held = ConcurrentHashMap.newKeySet();
+        AtomicReference<String> clash = new AtomicReference<>();
+        CyclicBarrier start = new CyclicBarrier(threads);
+        ExecutorService ranges = Executors.newFixedThreadPool(threads);
+        try {
+            List<Future<?>> done = new ArrayList<>();
+            for (int t = 0; t < threads; t++) {
+                done.add(ranges.submit(() -> {
+                    start.await();
+                    for (int r = 0; r < rounds; r++) {
+                        FlussConnectionPool.Lease lease = 
pool.borrow(config("server-a:9123"));
+                        if (!held.add(lease.connection())) {
+                            clash.set("a connection was lent to a second range 
while the first held it");
+                        }
+                        Thread.yield();
+                        held.remove(lease.connection());
+                        pool.giveBack(lease);
+                    }
+                    return null;
+                }));
+            }
+            for (Future<?> range : done) {
+                range.get(PATIENCE_SECONDS, TimeUnit.SECONDS);
+            }
+        } finally {
+            ranges.shutdownNow();
+        }
+        Assertions.assertNull(clash.get(), clash.get());
+        Assertions.assertTrue(closes.size() <= threads,
+                closes.size() + " connections opened for " + threads + " 
ranges reading at once");
+        Assertions.assertEquals(closes.size(), pool.idleCount());
+    }
+
+    /** A connection given up on is closed, not kept for the next range. */
+    @Test
+    public void discardedConnectionIsClosedAndNotLent() throws Exception {
+        FlussConnectionPool.Lease lease = pool.borrow(config("server-a:9123"));
+        pool.discard(lease);
+        Assertions.assertTrue(closes.get(0).await(PATIENCE_SECONDS, 
TimeUnit.SECONDS), "never closed");
+        Assertions.assertEquals(0, pool.idleCount());
+        Assertions.assertTrue(pool.borrow(config("server-a:9123")).opened());
+    }
+
+    @Test
+    public void connectionIdleForTheTimeoutIsClosed() throws Exception {
+        FlussConnectionPool.Lease lease = pool.borrow(config("server-a:9123"));
+        pool.giveBack(lease);
+
+        now.addAndGet(FlussConnectionPool.IDLE_TIMEOUT_NANOS - 1);
+        pool.closeIdleLongerThan(FlussConnectionPool.IDLE_TIMEOUT_NANOS);
+        Assertions.assertEquals(1, pool.idleCount(), "closed before its 
timeout");
+        Assertions.assertEquals(1, closes.get(0).getCount(), "closed before 
its timeout");
+
+        now.incrementAndGet();
+        pool.closeIdleLongerThan(FlussConnectionPool.IDLE_TIMEOUT_NANOS);
+        Assertions.assertEquals(0, pool.idleCount());
+        Assertions.assertTrue(closes.get(0).await(PATIENCE_SECONDS, 
TimeUnit.SECONDS), "never closed");
+        Assertions.assertTrue(pool.borrow(config("server-a:9123")).opened(), 
"a closed connection was lent");
+    }
+
+    /** Idle means idle since the last range gave it back, not since it was 
opened. */
+    @Test
+    public void connectionInUseIsNotIdle() {
+        FlussConnectionPool.Lease lease = pool.borrow(config("server-a:9123"));
+        pool.giveBack(lease);
+        now.addAndGet(50 * SECOND);
+        pool.giveBack(pool.borrow(config("server-a:9123")));
+
+        now.addAndGet(FlussConnectionPool.IDLE_TIMEOUT_NANOS - 10 * SECOND);
+        pool.closeIdleLongerThan(FlussConnectionPool.IDLE_TIMEOUT_NANOS);
+        Assertions.assertEquals(1, pool.idleCount(), "closed while it had been 
idle for less than the timeout");
+        Assertions.assertEquals(1, closes.get(0).getCount());
+    }
+
+    /**
+     * The reaper runs each sweep through {@link FlussConnectionPool#sweep} on 
a scheduled executor, which
+     * never runs a periodic task again once it has thrown. A scan that fills 
BE's JVM heap makes
+     * whatever allocates throw {@code OutOfMemoryError}, a sweep included: 
that must cost the one sweep,
+     * not every sweep after it.
+     */
+    @Test
+    public void sweepThatRunsIntoAnErrorLeavesTheNextSweepsComing() throws 
Exception {
+        AtomicInteger sweeps = new AtomicInteger();
+        FlussConnectionPool failing = new FlussConnectionPool(config -> 
open(), () -> {
+            sweeps.incrementAndGet();
+            throw new OutOfMemoryError("simulated: Java heap space");
+        }, FlussConnectionCloser::close, NO_REAPER);
+        ScheduledExecutorService reaper = 
Executors.newSingleThreadScheduledExecutor();
+        try {
+            reaper.scheduleWithFixedDelay(() -> 
FlussConnectionPool.sweep(failing), 0, 1, TimeUnit.MILLISECONDS);
+            long deadline = System.nanoTime() + 
TimeUnit.SECONDS.toNanos(PATIENCE_SECONDS);
+            while (sweeps.get() < 3 && System.nanoTime() < deadline) {
+                Thread.sleep(1);
+            }
+        } finally {
+            reaper.shutdownNow();
+        }
+        Assertions.assertTrue(sweeps.get() >= 3, "sweeps after the first 
error: " + (sweeps.get() - 1));
+    }
+
+    /**
+     * The same error, struck while the sweep is handing its connections to 
the closer: it may cost the
+     * one being handed over, never the ones the sweep had not reached. Those 
stay in the pool, and the
+     * next sweep closes them; out of it, nothing would ever close them.
+     */
+    @Test
+    public void sweepThatFailsMidwayKeepsTheConnectionsItHadNotHandedOver() 
throws Exception {
+        AtomicInteger handedOver = new AtomicInteger();
+        FlussConnectionPool failing = new FlussConnectionPool(config -> 
open(), now::get, connection -> {
+            if (handedOver.incrementAndGet() == 2) {
+                throw new OutOfMemoryError("simulated: Java heap space");
+            }
+            FlussConnectionCloser.close(connection);
+        }, NO_REAPER);
+        FlussConnectionPool.Lease first = 
failing.borrow(config("server-a:9123"));
+        FlussConnectionPool.Lease second = 
failing.borrow(config("server-a:9123"));
+        FlussConnectionPool.Lease third = 
failing.borrow(config("server-a:9123"));
+        failing.giveBack(first);
+        failing.giveBack(second);
+        failing.giveBack(third);
+        now.addAndGet(FlussConnectionPool.IDLE_TIMEOUT_NANOS);
+
+        FlussConnectionPool.sweep(failing);
+        Assertions.assertTrue(closes.get(0).await(PATIENCE_SECONDS, 
TimeUnit.SECONDS), "the first was never closed");
+        Assertions.assertEquals(1, failing.idleCount(), "connections the 
failed sweep had not reached");
+
+        FlussConnectionPool.sweep(failing);
+        Assertions.assertEquals(0, failing.idleCount());
+        Assertions.assertTrue(closes.get(2).await(PATIENCE_SECONDS, 
TimeUnit.SECONDS), "the third was never closed");
+    }
+
+    /**
+     * The reaper is started with the first connection given back, not while 
the class initializes, where
+     * an {@code OutOfMemoryError} - no native thread to be had, or a full 
heap - would leave the class
+     * unusable and every fluss read of the BE failing until it restarted. A 
start that fails costs the
+     * sweeps until a later connection given back starts the reaper; ranges go 
on borrowing and giving back
+     * meanwhile.
+     */
+    @Test
+    public void reaperThatFailedToStartIsStartedByALaterGiveBack() throws 
Exception {
+        AtomicInteger starts = new AtomicInteger();
+        List<Runnable> reapers = new ArrayList<>();
+        FlussConnectionPool starting = new FlussConnectionPool(config -> 
open(), now::get,
+                FlussConnectionCloser::close, sweep -> {
+                    if (starts.incrementAndGet() == 1) {
+                        throw new OutOfMemoryError("simulated: unable to 
create native thread");
+                    }
+                    reapers.add(sweep);
+                });
+
+        FlussConnectionPool.Lease first = 
starting.borrow(config("server-a:9123"));
+        starting.giveBack(first);
+        Assertions.assertEquals(1, starts.get(), "the first connection given 
back did not try to start the reaper");
+        Assertions.assertEquals(1, starting.idleCount(), "the failed start 
lost the connection given back");
+
+        FlussConnectionPool.Lease second = 
starting.borrow(config("server-a:9123"));
+        Assertions.assertSame(first.connection(), second.connection(), "the 
failed start stopped the lending");
+        starting.giveBack(second);
+        Assertions.assertEquals(1, reapers.size(), "the next connection given 
back did not start the reaper");
+
+        starting.giveBack(starting.borrow(config("server-a:9123")));
+        Assertions.assertEquals(2, starts.get(), "a reaper already running was 
started again");
+
+        now.addAndGet(FlussConnectionPool.IDLE_TIMEOUT_NANOS);
+        reapers.get(0).run();
+        Assertions.assertEquals(0, starting.idleCount());
+        Assertions.assertTrue(closes.get(0).await(PATIENCE_SECONDS, 
TimeUnit.SECONDS), "never closed");
+    }
+
+    private Connection open() {
+        CountDownLatch latch = new CountDownLatch(1);
+        closes.add(latch);
+        return new StubConnection(latch::countDown);
+    }
+
+    private static Configuration config(String bootstrapServers) {
+        Configuration config = new Configuration();
+        config.setString("bootstrap.servers", bootstrapServers);
+        return config;
+    }
+}
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerLogTest.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerLogTest.java
index e5c652d0646..eea0d3b2d50 100644
--- 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerLogTest.java
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerLogTest.java
@@ -53,6 +53,7 @@ import java.time.Instant;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.LinkedHashMap;
@@ -430,20 +431,37 @@ public class FlussJniScannerLogTest {
     // ---------------------------------------------------------------- 
lifecycle and fail-loud
 
     /**
-     * A fluss connection owns netty and metadata-updater threads that outlive 
the query if the scanner
-     * leaks it — and a BE runs scanners for the life of the process.
+     * A fluss connection owns netty and metadata-updater threads, and a BE 
runs scanners for the life of
+     * the process. Ranges read one after another borrow one connection from 
the pool instead of opening
+     * one each; the pool keeps it, threads and all, while it is idle, and 
once the pool closes it the
+     * threads go.
      */
     @Test
-    public void closingTheScannerReleasesItsClientThreads() throws Exception {
+    public void 
rangesReadOneAfterAnotherShareAConnectionWhoseThreadsGoWhenItIsClosed() throws 
Exception {
         TablePath tablePath = TablePath.of(db, "threads");
         createIntTable(tablePath);
         appendInts(tablePath, 0, 3);
-        int before = countClientThreads();
-
-        for (int i = 0; i < 3; i++) {
-            scanAll(tablePath, columns("id", "int"), 0, 3, 1024);
+        // Start from an empty pool - earlier tests leave their connections 
idle in it - so that the
+        // first range below has to open the connection the others borrow.
+        FlussConnectionPool.INSTANCE.closeIdleLongerThan(0);
+        int before = settledClientThreads();
+
+        List<String> opened = new ArrayList<>();
+        for (int i = 0; i < 4; i++) {
+            FlussJniScanner scanner = new FlussJniScanner(1024, 
params(tablePath, columns("id", "int"), 0, 3));
+            scanner.open();
+            while (scanner.getNextBatchMeta() != 0) {
+                scanner.resetTable();
+            }
+            scanner.releaseTable();
+            
opened.add(scanner.getStatistics().get("counter:FlussJniConnectionsOpened"));
+            scanner.close();
         }
+        Assertions.assertEquals(Arrays.asList("1", "0", "0", "0"), opened, 
"connections opened, range by range");
+        Assertions.assertEquals(1, FlussConnectionPool.INSTANCE.idleCount());
+        Assertions.assertTrue(countClientThreads() > before, "the idle 
connection has no threads of its own");
 
+        FlussConnectionPool.INSTANCE.closeIdleLongerThan(0);
         long deadline = System.currentTimeMillis() + 30_000;
         while (countClientThreads() > before && System.currentTimeMillis() < 
deadline) {
             Thread.sleep(100);
@@ -452,6 +470,85 @@ public class FlussJniScannerLogTest {
                 "fluss client threads leaked: " + before + " before, " + 
countClientThreads() + " after");
     }
 
+    /**
+     * A connection serves one range at a time, so the scanner opens it with 
one network thread rather
+     * than fluss's default four; a catalog that sets the number keeps it. 
Netty starts a thread for each
+     * server a connection reaches, up to that number, and a range here 
reaches two - the coordinator and
+     * the one tablet server - so one thread carries both, or two carry one 
each.
+     */
+    @Test
+    public void connectionRunsOneNetworkThreadUnlessTheCatalogSetsTheNumber() 
throws Exception {
+        TablePath tablePath = TablePath.of(db, "network_threads");
+        createIntTable(tablePath);
+        appendInts(tablePath, 0, 3);
+        FlussConnectionPool.INSTANCE.closeIdleLongerThan(0);
+        int before = settledClientThreads();
+
+        scanAll(tablePath, columns("id", "int"), 0, 3, 1024);
+        Assertions.assertEquals(1, FlussConnectionPool.INSTANCE.idleCount());
+        Assertions.assertEquals(before + 1, settledClientThreads(), "network 
threads of the default connection");
+
+        Map<String, String> params = params(tablePath, columns("id", "int"), 
0, 3);
+        params.put("fluss.client.netty.client.num-network-threads", "2");
+        runScanner(params, 1024);
+        Assertions.assertEquals(2, FlussConnectionPool.INSTANCE.idleCount(),
+                "the setting opened a connection of its own");
+        Assertions.assertEquals(before + 3, settledClientThreads(), "network 
threads of both connections");
+
+        FlussConnectionPool.INSTANCE.closeIdleLongerThan(0);
+    }
+
+    /**
+     * Only a range read to its end gives its connection back. One closed 
early - a LIMIT, a cancel - may
+     * still have work in flight on it, and one that failed may have failed 
because of it; either way the
+     * next range opens a connection of its own.
+     */
+    @Test
+    public void rangeClosedBeforeItsEndOrFailedDoesNotGiveItsConnectionBack() 
throws Exception {
+        TablePath tablePath = TablePath.of(db, "early_close");
+        createIntTable(tablePath);
+        appendInts(tablePath, 0, 3);
+        FlussConnectionPool.INSTANCE.closeIdleLongerThan(0);
+
+        // One row of three, then closed.
+        FlussJniScanner early = new FlussJniScanner(1, params(tablePath, 
columns("id", "int"), 0, 3));
+        early.open();
+        Assertions.assertNotEquals(0, early.getNextBatchMeta());
+        early.releaseTable();
+        early.close();
+        Assertions.assertEquals(0, FlussConnectionPool.INSTANCE.idleCount(), 
"a range closed early gave back");
+
+        // Opened, then failed: the column is not in the table.
+        Assertions.assertThrows(Exception.class,
+                () -> runScanner(params(tablePath, columns("nope", "int"), 0, 
3), 1024));
+        Assertions.assertEquals(0, FlussConnectionPool.INSTANCE.idleCount(), 
"a range that failed gave back");
+
+        scanAll(tablePath, columns("id", "int"), 0, 3, 1024);
+        Assertions.assertEquals(1, FlussConnectionPool.INSTANCE.idleCount(), 
"a range read to its end did not");
+    }
+
+    /**
+     * A range gives its connection back to the pool rather than closing it, 
and closing one waits out
+     * netty's two-second graceful shutdown: closing the scanner must not wait 
for anything of the kind.
+     */
+    @Test
+    public void closingTheScannerDoesNotWaitForItsConnectionToShutDown() 
throws Exception {
+        TablePath tablePath = TablePath.of(db, "close_latency");
+        createIntTable(tablePath);
+        appendInts(tablePath, 0, 3);
+
+        FlussJniScanner scanner = new FlussJniScanner(1024, params(tablePath, 
columns("id", "int"), 0, 3));
+        scanner.open();
+        while (scanner.getNextBatchMeta() != 0) {
+            scanner.resetTable();
+        }
+        scanner.releaseTable();
+        long start = System.nanoTime();
+        scanner.close();
+        long closeMillis = (System.nanoTime() - start) / 1_000_000;
+        Assertions.assertTrue(closeMillis < 1000, "closing the scanner took " 
+ closeMillis + " ms");
+    }
+
     /** An unreadable range must name what is wrong, not hand fluss a null and 
fail somewhere inside. */
     @Test
     public void missingParameterIsNamed() {
@@ -501,6 +598,25 @@ public class FlussJniScannerLogTest {
         return message.toString();
     }
 
+    /**
+     * The client thread count once it has stopped changing: connections that 
are closing hold their
+     * threads through netty's two-second quiet period, so it has to stay put 
for longer than that.
+     */
+    private static int settledClientThreads() throws InterruptedException {
+        long deadline = System.currentTimeMillis() + 30_000;
+        int count = countClientThreads();
+        long stableSince = System.currentTimeMillis();
+        while (System.currentTimeMillis() - stableSince < 2_500 && 
System.currentTimeMillis() < deadline) {
+            Thread.sleep(100);
+            int now = countClientThreads();
+            if (now != count) {
+                count = now;
+                stableSince = System.currentTimeMillis();
+            }
+        }
+        return count;
+    }
+
     private static int countClientThreads() {
         int count = 0;
         for (Thread thread : Thread.getAllStackTraces().keySet()) {
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerPkTest.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerPkTest.java
index ba17118180a..4d0ab10e103 100644
--- 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerPkTest.java
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/FlussJniScannerPkTest.java
@@ -27,6 +27,7 @@ import org.apache.fluss.client.admin.OffsetSpec;
 import org.apache.fluss.client.metadata.KvSnapshots;
 import org.apache.fluss.client.table.Table;
 import org.apache.fluss.client.table.writer.UpsertWriter;
+import org.apache.fluss.fs.FsPathAndFileName;
 import org.apache.fluss.metadata.DatabaseDescriptor;
 import org.apache.fluss.metadata.PartitionSpec;
 import org.apache.fluss.metadata.Schema;
@@ -49,7 +50,13 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.RegisterExtension;
 
+import java.io.FileInputStream;
+import java.io.IOException;
 import java.math.BigDecimal;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.time.Duration;
 import java.time.Instant;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
@@ -61,6 +68,9 @@ import java.util.HashMap;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.Callable;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Stream;
 
 /**
  * Reading a fluss primary-key table through the scanner, against a real 
cluster in this JVM.
@@ -451,6 +461,54 @@ public class FlussJniScannerPkTest {
         Assertions.assertThrows(Exception.class, () -> read(params, 1024));
     }
 
+    // ---------------------------------------------------------------- 
lifecycle
+
+    /**
+     * BE can close a range while fluss is still copying its kv snapshot - the 
range of a query that was
+     * cancelled or failed elsewhere - and fluss's reader can only be closed 
once that copy is over. The
+     * copy runs on the download threads of the range's connection, so the 
connection has to outlive it:
+     * shut down under the copy, its download pool drops the files still 
queued, and the reader waits for
+     * them forever, with the pool thread it runs on, the thread waiting to 
close it and the half-copied
+     * snapshot directory.
+     *
+     * <p>The test holds the copy up for as long as it needs: the snapshot 
file fluss copies first is
+     * swapped for a FIFO, on which the copy blocks until the file's bytes are 
written into it, and a
+     * single download thread leaves every other file queued behind it.
+     */
+    @Test
+    public void 
rangeClosedWhileItsSnapshotIsCopiedLeavesNothingWaitingForTheCopy() throws 
Exception {
+        TablePath tablePath = TablePath.of(db, "pk_closed_while_copying");
+        createPkTable(tablePath);
+        upsert(tablePath, row(1, "one"), row(2, "two"));
+        long snapshotId = snapshot(tablePath);
+        long logOffset = snapshotLogOffset(tablePath);
+        List<FsPathAndFileName> files = admin.getKvSnapshotMetadata(
+                new TableBucket(tableId(tablePath), 0), 
snapshotId).get().getSnapshotFiles();
+        Assertions.assertTrue(files.size() > 1, "no file would wait behind the 
first: " + files);
+        // Copied first: fluss hands the files to the download threads in this 
order.
+        Path first = Paths.get(files.get(0).getPath().toUri());
+        byte[] firstBytes = Files.readAllBytes(first);
+        Files.delete(first);
+        Assertions.assertEquals(0, new ProcessBuilder("mkfifo", 
first.toString()).start().waitFor());
+
+        Path copyDir = Files.createTempDirectory("doris-fluss-snapshot-copy");
+        Map<String, String> params = pkParams(tablePath, columns("id", "int", 
"name", "string"),
+                snapshotId, logOffset, logOffset);
+        params.put("fluss.client.client.remote-file.download-thread-num", "1");
+        params.put("fluss.client.client.scanner.io.tmpdir", 
copyDir.toString());
+        FlussJniScanner scanner = new FlussJniScanner(1024, params);
+        scanner.open();
+        await(FlussJniScannerPkTest::copyBlockedOnItsFirstFile, "the snapshot 
copy never reached its first file");
+        scanner.close();
+        // Long enough for a connection closed under the copy to have shut its 
download pool down.
+        Thread.sleep(1000);
+
+        // The copy is the FIFO's reader, so this returns once the bytes are 
handed over.
+        Assertions.assertTimeoutPreemptively(Duration.ofSeconds(60), () -> 
Files.write(first, firstBytes));
+        await(() -> isEmpty(copyDir), "the snapshot reader is still waiting 
for files nobody will copy");
+        Files.delete(copyDir);
+    }
+
     // ---------------------------------------------------------------- helpers
 
     private void createPkTable(TablePath tablePath) throws Exception {
@@ -591,4 +649,35 @@ public class FlussJniScannerPkTest {
         Arrays.sort(sorted, Comparator.comparingInt(row -> (Integer) row[0]));
         return sorted;
     }
+
+    /** Whether a thread copying kv snapshot files is blocked opening one of 
them - the FIFO, above. */
+    private static boolean copyBlockedOnItsFirstFile() {
+        for (StackTraceElement[] stack : Thread.getAllStackTraces().values()) {
+            boolean opening = false;
+            boolean copying = false;
+            for (StackTraceElement frame : stack) {
+                opening |= 
frame.getClassName().equals(FileInputStream.class.getName())
+                        && frame.getMethodName().equals("open0");
+                copying |= 
frame.getClassName().equals("org.apache.fluss.fs.utils.FileDownloadUtils");
+            }
+            if (opening && copying) {
+                return true;
+            }
+        }
+        return false;
+    }
+
+    private static boolean isEmpty(Path dir) throws IOException {
+        try (Stream<Path> entries = Files.list(dir)) {
+            return !entries.findAny().isPresent();
+        }
+    }
+
+    private static void await(Callable<Boolean> condition, String failure) 
throws Exception {
+        long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(60);
+        while (!condition.call()) {
+            Assertions.assertTrue(System.nanoTime() < deadline, failure);
+            Thread.sleep(50);
+        }
+    }
 }
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScannerTest.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScannerTest.java
index bd4e099e2b8..1b41e96a5eb 100644
--- 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScannerTest.java
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/SafeKvSnapshotAndLogBatchScannerTest.java
@@ -126,6 +126,81 @@ public class SafeKvSnapshotAndLogBatchScannerTest {
                 "the SDK scanner must be closed exactly once, after 
publication");
     }
 
+    /**
+     * Until the SDK scanner is closed it is still copying its snapshot on the 
download threads of the
+     * connection it was created on, and closing the connection under that 
copy strands it. So the
+     * connection is released only once the scanner is closed, however late 
that is.
+     */
+    @Test
+    public void earlyClosedSnapshotReaderReleasesItsConnectionOnlyOnceClosed() 
throws Exception {
+        LatePublishingSnapshotScanner snapshot = new 
LatePublishingSnapshotScanner();
+        SafeKvSnapshotAndLogBatchScanner.ScannerResources resources = 
acquireSnapshotOnly(snapshot);
+        resources.snapshotScanner.close();
+
+        Assertions.assertTrue(snapshot.pollEntered.await(5, TimeUnit.SECONDS),
+                "early close did not start a publication waiter");
+        Assertions.assertFalse(resources.snapshotScanner.released().isDone(),
+                "the connection was released while the snapshot was still 
being copied on it");
+
+        snapshot.publishNativeReader();
+        resources.snapshotScanner.released().get(5, TimeUnit.SECONDS);
+        Assertions.assertEquals(1, snapshot.closeCalls.get(), "released before 
the SDK scanner was closed");
+    }
+
+    /** A snapshot that had arrived is closed at once, and with it the 
connection is released at once. */
+    @Test
+    public void 
snapshotReaderClosedAfterItsSnapshotArrivedReleasesItsConnectionAtOnce() throws 
Exception {
+        LatePublishingSnapshotScanner snapshot = new 
LatePublishingSnapshotScanner();
+        SafeKvSnapshotAndLogBatchScanner.ScannerResources resources = 
acquireSnapshotOnly(snapshot);
+        snapshot.publishNativeReader();
+        // Ready and empty: the SDK's null.
+        
Assertions.assertNull(resources.snapshotScanner.pollBatch(Duration.ofSeconds(5)));
+
+        resources.snapshotScanner.close();
+        Assertions.assertTrue(resources.snapshotScanner.released().isDone(),
+                "a scanner closed after publication must release its 
connection on the spot");
+        Assertions.assertEquals(1, snapshot.closeCalls.get());
+    }
+
+    /**
+     * A range is often closed early because its query filled BE's JVM heap, 
and the waiter's polls
+     * allocate. Ended by an OutOfMemoryError on its own thread, the waiter 
would leave the SDK scanner
+     * open and the connection waiting for released() for good; it has to keep 
waiting instead.
+     */
+    @Test
+    public void publicationWaiterOutlastsAnOutOfMemoryErrorOnItsOwnThread() 
throws Exception {
+        LatePublishingSnapshotScanner snapshot = new 
LatePublishingSnapshotScanner();
+        SafeKvSnapshotAndLogBatchScanner.ScannerResources resources =
+                acquireSnapshotOnly(new OutOfMemoryOnFirstPoll(snapshot));
+        resources.snapshotScanner.close();
+
+        Assertions.assertTrue(snapshot.pollEntered.await(5, TimeUnit.SECONDS),
+                "the waiter did not poll again after an OutOfMemoryError on 
its thread");
+        snapshot.publishNativeReader();
+        resources.snapshotScanner.released().get(5, TimeUnit.SECONDS);
+        Assertions.assertEquals(1, snapshot.closeCalls.get(),
+                "the SDK scanner must be closed exactly once, after 
publication");
+    }
+
+    /** Opens a range that reads only {@code snapshot}: its log range is 
empty. */
+    private static SafeKvSnapshotAndLogBatchScanner.ScannerResources 
acquireSnapshotOnly(BatchScanner snapshot) {
+        SafeKvSnapshotAndLogBatchScanner.ScannerFactory factory =
+                new SafeKvSnapshotAndLogBatchScanner.ScannerFactory() {
+                    @Override
+                    public BatchScanner createSnapshotScanner(
+                            TableBucket tableBucket, long snapshotId, int[] 
projectedFields) {
+                        return snapshot;
+                    }
+
+                    @Override
+                    public LogScanner createLogScanner(int[] projectedFields) {
+                        throw new AssertionError("the staged log range is 
empty");
+                    }
+                };
+        return SafeKvSnapshotAndLogBatchScanner.acquireScanners(
+                factory, new TableBucket(1L, 0), 7L, 20L, 20L, new int[] {0});
+    }
+
     private static class RecordingLogScanner implements LogScanner {
         final AtomicBoolean closed = new AtomicBoolean();
         final AtomicBoolean subscribed = new AtomicBoolean();
@@ -212,4 +287,27 @@ public class SafeKvSnapshotAndLogBatchScannerTest {
             published.countDown();
         }
     }
+
+    /** The first poll fails as an allocation on the polling thread does while 
the heap is full. */
+    private static final class OutOfMemoryOnFirstPoll implements BatchScanner {
+        private final BatchScanner delegate;
+        private final AtomicBoolean failed = new AtomicBoolean();
+
+        private OutOfMemoryOnFirstPoll(BatchScanner delegate) {
+            this.delegate = delegate;
+        }
+
+        @Override
+        public CloseableIterator<InternalRow> pollBatch(Duration timeout) 
throws IOException {
+            if (failed.compareAndSet(false, true)) {
+                throw new OutOfMemoryError("simulated: Java heap space");
+            }
+            return delegate.pollBatch(timeout);
+        }
+
+        @Override
+        public void close() throws IOException {
+            delegate.close();
+        }
+    }
 }
diff --git 
a/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/StubConnection.java
 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/StubConnection.java
new file mode 100644
index 00000000000..223daa80197
--- /dev/null
+++ 
b/fe/be-java-extensions/fluss-scanner/src/test/java/org/apache/doris/fluss/StubConnection.java
@@ -0,0 +1,64 @@
+// 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.admin.Admin;
+import org.apache.fluss.client.table.MultiTable;
+import org.apache.fluss.client.table.Table;
+import org.apache.fluss.config.Configuration;
+import org.apache.fluss.metadata.TablePath;
+
+/** A connection nobody uses for anything but lending it out and closing it. */
+final class StubConnection implements Connection {
+
+    interface CloseAction {
+        void run() throws Exception;
+    }
+
+    private final CloseAction onClose;
+
+    StubConnection(CloseAction onClose) {
+        this.onClose = onClose;
+    }
+
+    @Override
+    public Configuration getConfiguration() {
+        throw new UnsupportedOperationException();
+    }
+
+    @Override
+    public Admin getAdmin() {
+        throw new UnsupportedOperationException();
+    }
+
+    @Override
+    public Table getTable(TablePath tablePath) {
+        throw new UnsupportedOperationException();
+    }
+
+    @Override
+    public MultiTable getMultiTable() {
+        throw new UnsupportedOperationException();
+    }
+
+    @Override
+    public void close() throws Exception {
+        onClose.run();
+    }
+}


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

Reply via email to