morningman opened a new pull request, #68711:
URL: https://github.com/apache/doris/pull/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 a 3-row log table: 2069 ms, 2044 ms of it in close;
      - `COUNT(*)` on a 30M-row log table: 4480 ms, about 32 s of close summed 
over 16 ranges;
      - union read of a log table with a 1% log tail: `COUNT(*)` 8320 ms, of 
which 32 ms reading the tail, against 130 ms with no tail and 4610 ms with 
union read disabled. 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:
   
   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.
   
      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 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.
   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:
   
   - 64-bucket log table, `COUNT(*)` 20 times back to back: every fourth query 
took 2.2 s; now the slowest takes 128 ms (median 133 -> 103 ms).
   - Partitioned primary-key lake table (30 partitions x 4 buckets, 1% tail), 
union `COUNT(*)`, 240 JNI reads per query: 2090 -> 1667 ms; four at once 11.3 
-> 4.9 s each; BE threads 3326 -> 2067.
   - 16 buckets read in parallel do not slow down (log `SELECT *` 10.8 -> 10.5 
s, primary-key 6.67 -> 6.46 s).
   
   One network thread per connection (four A/B rounds, only the plugin 
swapped): three selectors fewer per connection; with 4 and 8 concurrent `SELECT 
*`, peak file descriptors 1238-2415 -> 811-1167 and peak threads 2186-2361 -> 
2085-2245, latency within 5% either way.
   
   Snapshot copy (commit 2): six out-of-memory `SELECT *` queries in a row over 
a 16-bucket primary-key table used to leave two copy threads parked in 
`FileDownloadUtils.transferAllDataToDirectory`, two waiter threads blocked 
behind them and two snapshot directories (17 MB and 4.3 MB); now nothing is 
left.
   
   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.
   - `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 #PR3: 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, 8): 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. `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` (2 new). BE 
`FileScanOperatorFlussTest.EveryInstanceOfTh
 eNodeReadsThroughTheNodesCache` (new). On this branch: fluss-scanner 77 tests 
pass. 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 -->
   


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

To unsubscribe, e-mail: [email protected]

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


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

Reply via email to