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]