morningman opened a new pull request, #68713:
URL: https://github.com/apache/doris/pull/68713
### What problem does this PR solve?
Issue Number: None
Related PR: #PR3 (scans run up to 16 scanners per instance again), #PR1 (an
out-of-memory error no longer exits BE)
Problem Summary:
**In short.** A few kinds of JNI reader keep hundreds of megabytes of BE's
JVM heap each, from before their first batch until they close; sixteen of them
opened at once, one scan instance's worth, run the default 2 GB heap out and
the query fails with `OutOfMemoryError`. This PR adds an opt-in admission: a
statement that sets `enable_jni_heap_admission` has the paimon and fluss
connectors declare the heap each such reader will hold, and BE opens those
readers only while their declarations fit in half the JVM's heap. It also makes
the out-of-memory error say which heap ran out and how to give it more. Off by
default; without the variable only the error text changes.
**Background**
- Every JNI reader in BE (paimon, fluss, hudi, max-compute, trino, JDBC,
iceberg system tables) runs its Java scanner in the one JVM BE embeds,
`-Xmx2048m` by default.
- Nothing bounds what those scanners keep in the Java heap: BE's memory
trackers do not see it, and a scan's concurrency (up to
`max_file_scanners_concurrency`, 16, per instance; see #PR3) counts blocks, not
Java objects.
- Most readers stream and keep a few megabytes. Three kinds hold a lot, from
before their first batch until they close:
- a paimon split read through JNI that has to be merged (a primary-key
table written since its last compaction) keeps a row group of every file it
merges in the heap, about 130 MB per split on a 30M-row table;
- a fluss primary-key bucket read (`PK_FULL`) replays the change log after
the bucket's kv snapshot into a sorted map;
- the log tail of a fluss union read of a primary-key table (`PK_TAIL`)
keeps the last row of every key in its slice.
**The problem, and what it cost**
On a Release BE with the default 2 GB heap, each of these fails every time
with `OutOfMemoryError: Java heap space`:
- `SELECT *` over a 30M-row paimon primary-key table that is not compacted
(32 JNI merge splits), through the paimon catalog or fluss's `$lake`, with 16
scanners;
- `SELECT *` over a fluss primary-key table of 2M keys of which 1M were
updated since its last kv snapshot, with 16 scanners, and four such queries at
once;
- eight concurrent union reads of a primary-key table with 100K updated keys
per bucket in its log tail.
And the error did not help the user out:
```
errCode = 2, detailMessage = (127.0.0.1)[JNI_ERROR]OutOfMemoryError: Java
heap space
```
It names neither the JVM nor where its size is set, and the statements that
came right after one that ran out failed with `[JNI_ERROR]Failed to attach the
current thread to the JVM, code=-1`, which says nothing about heap at all:
until the failed statement's scanners have exited, its heap is not freed and
the JVM cannot even allocate the `java.lang.Thread` an attaching thread needs.
Measuring the heap does not make a usable gate. `Runtime.totalMemory() -
freeMemory()` counts garbage G1 has not collected yet (a gigabyte of it can sit
in the old generation for minutes while young collections keep running) and
everybody else in the JVM (libhdfs reads and writes, Java UDFs, JDBC), so a
gate built on it held readers back for heap they would never need: in a trial,
eight concurrent `SELECT *` over a fluss log table waited a median 62 s instead
of 28 s. Asking the JVM to collect before measuring stops every JNI reader, UDF
and HDFS stream for each collection.
**How this PR fixes it**
The connector that plans such a range knows what the reader will hold (file
sizes and table options for paimon, log offsets and the schema for fluss), so
the range says it, and BE keeps an account. Five commits:
1. **BE admission gate.** `TFileRangeDesc.jni_heap_bytes` (new, optional)
carries the declaration; the session variable `enable_jni_heap_admission`
(default `false`) decides whether connectors declare at all. `JniScanHeapGate`,
one per process, admits a declaring reader before it creates its Java scanner,
on both JNI reader paths (V1 `JniReader`, V2 `JniTableReader`), while what the
readers it admitted declared, plus this one, fits in
`jni_scanner_heap_budget_ratio` (new BE config, default 0.5) of the `-Xmx` BE
started the JVM with. The share comes back when the scanner closes, not at its
first batch, because these readers keep what they hold until they close.
Readers that do not fit wait in arrival order, so a large one is not passed
over by small ones. A reader that is alone is always admitted, a cancelled
query stops waiting, and any wait ends after `jni_scanner_heap_max_wait_ms`
(new BE config, 60 s). A range that declares nothing - every reader of a
statement that did not ask
, and always JDBC, trino, max-compute, iceberg system tables, fluss log ranges
and the catalog connection test - opens at once, as it does now. The profile
gains `JvmHeapDeclaredBytes` and `JvmHeapWaitTime`.
2. **paimon declares.** For every JNI `DataSplit` of such a statement: if
the split has to be merged, the sum over its files of one row group (two for a
file holding more than one, never more than the file) plus a 1 MB dictionary
page per column per file; otherwise the largest file's share. The row group
size is what paimon's writers use: `file.block-size` if the table sets it, else
`parquet.block.size` (128 MB) or, for ORC, `orc.stripe.size` (64 MB). Every
column counts, because a merge reads the key, sequence and row-kind columns
whatever is projected. Fluss's lake splits hand their descriptor to the paimon
range they wrap and declare through it.
3. **fluss declares.** For `PK_FULL` and `PK_TAIL` ranges: records x (row +
map entry), the records being the range's offset span (an upper bound on its
keys), the row a `GenericRow` of boxed fields over the projected and key
columns (strings and bytes at 64 bytes, since fluss's metadata gives no
length), the entry 112 bytes for `PK_FULL`'s `TreeMap` and 96 for `PK_TAIL`'s
`LinkedHashMap`. `LOG` ranges stream and declare nothing, nor does a bucket
whose snapshot holds everything.
4. **The out-of-memory error says what to do.** An error whose Java
exception, or a cause of it, is `OutOfMemoryError: Java heap space`, the
"likely out of memory" errors raised when rendering the exception itself fails,
and an attach that fails with `JNI_ERR` now end with: BE's JVM is out of heap;
its maximum is N MB, the `-Xmx` in `JAVA_OPTS_FOR_JDK_17` of be.conf; raise it
and restart the BE, or have the statement hold less at once - set
`max_file_scanners_concurrency` below its default of 16, or
`enable_jni_heap_admission = true`; a statement that ran out frees the heap
only once its scanners have exited. Direct-memory, metaspace and native-thread
errors do not get it, since `-Xmx` does nothing for them.
5. **Regression test.**
`external_table_p0/fluss/test_fluss_jni_heap_admission` reads `PK_FULL` ranges,
a primary-key union read (paimon lake splits and a `PK_TAIL` tail, both
declaring) and a log read (declaring nothing) with the variable on, and checks
every result against the same query with it off: admission may only make a
reader wait, never change what it returns.
Why off by default: an estimate can only be an upper bound to be safe, so it
can serialize reads that would have fit (a `PK_TAIL` whose tail is all updates
declares twice its need). Without the variable, a scan that needs more heap
than the JVM has fails as before, now with an error that says how to fix it.
**Results**
A/B on one Release BE with the default 2 GB heap and #PR3 applied (16
scanners per instance), same process, only the session variable switched, three
rounds each:
| Read | variable off (default) | variable on |
|---|---|---|
| fluss primary-key table, 1M of 2M keys updated after the snapshot, `SELECT
*`, 16 scanners | OOM every round | 1.68-1.72 s |
| the same, 4 concurrent | 4 of 4 OOM every round | 5.8-6.2 s |
| paimon uncompacted primary-key table (32 merge splits), `SELECT *` | OOM
every round | 11.2-11.7 s |
| fluss union read, 100K updated keys per bucket in the tail, 8 concurrent,
16 scanners each | 8 of 8 fail every round | no OOM; 81.6-84.1 s (on macOS some
of these queries hit libcurl's 1024-descriptor limit, unrelated to the heap) |
| fluss primary-key table without a change log after its snapshot, `SELECT
*` (declares 0) | 2.76-2.78 s | 2.71-2.77 s |
| fluss log table, `SELECT *`, 8 concurrent (declares nothing) | 21.3-22.0 s
| 21.2-21.8 s |
- Reads that declare nothing are as fast with the variable on as off.
- When the variable is on and the budget is exceeded, readers queue: 16
readers declaring 2.15 GB in all against a 1 GB budget waited 7.9 s summed over
the scanners and finished in 1.7 s.
- How close the declarations are: a fluss `PK_FULL` bucket declared 144 MB
and held 116 MB at its peak (1.24x). A paimon merge split declared 160-180 MB
against 150-170 MB held steadily, but a merge briefly holds up to about 3x the
declaration (several hundred decoded column batches at once); that did not
cause an OOM in the runs above.
- With the variable off the failing statements now carry the hint, e.g.
`OutOfMemoryError: Java heap space. BE's JVM is out of heap. Its maximum is
2048 MB, the -Xmx in JAVA_OPTS_FOR_JDK_17 of be.conf: ...`.
**Classes, and how they call each other**
- `SessionVariable.enableJniHeapAdmission` (FE, new):
`enable_jni_heap_admission`, forwarded to the connectors in the session
properties.
- `PaimonScanPlanProvider` -> `PaimonJniHeapEstimate` (new) ->
`PaimonScanRange` (paimon connector): computes and carries the declaration of a
JNI split.
- `FlussScanPlanProvider` -> `FlussJniHeapEstimate` (new) ->
`FlussScanRange` (fluss connector): the same for `PK_FULL` / `PK_TAIL`.
- `TFileRangeDesc.jni_heap_bytes` (thrift, new field 17): unset or 0 means
the reader opens without waiting.
- `JniScanHeapGate` (BE, new, `util/jni_scan_heap_gate.{h,cpp}`):
`acquire(bytes, stop_waiting, &permit, &wait_ns)`; the `Permit` gives the share
back on `release()` or destruction.
- `JniReader::open` (V1, paimon) and `JniTableReader::_open_jni_scanner`
(V2, fluss and the others) (BE): acquire before creating the Java scanner,
release when it closes or if it never opened.
- `FlussUnionLakeReader::tail_scan_range` (BE, untouched apart from a
comment): the tail-key read it builds declares nothing, so it never waits while
the lake split it serves may hold a permit on the same thread.
- `Jni::Util::jvm_heap_exhausted_hint()` (BE, new) used by
`Env::GetJniExceptionMsg` and `JvmLauncher::_attach_current_thread`.
```
FE SET enable_jni_heap_admission = true
|-- paimon: PaimonScanPlanProvider -- JNI DataSplit -->
PaimonJniHeapEstimate.estimate(split, table options)
|-- fluss: FlussScanPlanProvider -- PK_FULL / PK_TAIL -->
FlussJniHeapEstimate.estimate(offset span, columns)
'--> TFileRangeDesc.jni_heap_bytes
(LOG ranges, other connectors: unset)
BE scanner thread
JniReader::open (V1) / JniTableReader::_open_jni_scanner (V2)
jni_heap_bytes > 0 ?
JniScanHeapGate::instance()->acquire(bytes) admitted while
sum(declared by admitted readers) + bytes
<=
jni_scanner_heap_budget_ratio x -Xmx;
FIFO; alone ->
admitted; cancelled -> stops waiting;
jni_scanner_heap_max_wait_ms -> admitted anyway
create + open the Java scanner ... read ... close -->
permit.release()
a Java exception "OutOfMemoryError: Java heap space" / attach JNI_ERR
--> error + jvm_heap_exhausted_hint()
```
Not touched: readers that declare nothing behave exactly as now (no lock, no
wait), and nothing calls into the JVM to measure or collect its heap.
**Merge order.** No code dependency on the other PRs. Its problem statement
and measurements assume #PR3 (16 scanners per instance by default), so it is
best merged right after it. #PR1 keeps an out-of-memory error from exiting BE
through fluss or paimon threads.
### Release note
New session variable `enable_jni_heap_admission` (default false). When true,
the JNI readers that hold much of BE's JVM heap (paimon splits read through
JNI, fluss primary-key bucket reads and union-read log tails) declare how much,
and BE opens them only while the declarations fit in
`jni_scanner_heap_budget_ratio` (new BE config, default 0.5) of the JVM's
maximum heap; a reader waits at most `jni_scanner_heap_max_wait_ms` (new BE
config, default 60000). The error a query gets when BE's JVM runs out of heap
now names the heap size and how to raise it.
### 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.** BE: `JniScanHeapGateTest` (new, 7: budget, arrival order,
a reader larger than the budget, a budget change applied to waiting readers,
stop, the longest wait, permit release); `JniUtilHeapSizeTest` (2 new: only an
out-of-heap error asks for more heap; the hint names the heap and the ways
out). 232 BE tests in 17 suites pass, among them the JNI reader suites (V1, V2,
paimon, fluss, JDBC, hudi), the file scanner and the scanner scheduling suites.
FE: `PaimonJniHeapEstimateTest`, `PaimonScanRangeJniHeapTest`,
`FlussJniHeapEstimateTest`, `FlussScanRangeTest`, `FlussSplitPlanTest`; on this
branch fe-connector-paimon 659 tests (1 skipped, as on master) and
fe-connector-fluss 295 tests pass.
**Regression.** `external_table_p0/fluss`: all 18 suites, including the
new `test_fluss_jni_heap_admission`, pass against a local fluss 1.0.0 docker
environment.
**Manual.** The A/B above.
- Behavior changed:
- [x] No. <!-- Explain the behavior change --> Not unless
`enable_jni_heap_admission` is set; the out-of-memory error text changes.
- [ ] Yes.
- Does this need documentation?
- [ ] No.
- [x] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 --> The session variable and
the two BE configs; doc PR to follow.
### 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]