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


##########
be/src/util/jni_scan_heap_gate.cpp:
##########
@@ -0,0 +1,144 @@
+// 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.
+
+#include "util/jni_scan_heap_gate.h"
+
+#include <algorithm>
+#include <chrono>
+#include <limits>
+#include <utility>
+
+#include "common/config.h"
+#include "common/logging.h"
+#include "util/jni-util.h"
+#include "util/time.h"
+
+namespace doris {
+
+namespace {
+
+// How often a waiting reader looks again although nobody told it to: to see 
its query cancelled and
+// its wait run out.
+constexpr auto POLL_INTERVAL = std::chrono::milliseconds(100);
+constexpr int64_t MB = 1024 * 1024;
+
+// A share of the -Xmx the JVM was started with, read from the same options 
the JVM was created from
+// - not measured, so nothing here calls into the JVM.
+int64_t jvm_heap_budget() {
+    const double budget = 
static_cast<double>(Jni::Util::get_max_jni_heap_memory_size()) *
+                          config::jni_scanner_heap_budget_ratio;
+    // A BE_TEST build reports an unlimited heap (SIZE_MAX), which stays 
unlimited here.
+    constexpr auto UNLIMITED = std::numeric_limits<int64_t>::max();
+    return budget >= static_cast<double>(UNLIMITED) ? UNLIMITED : 
static_cast<int64_t>(budget);
+}
+
+} // namespace
+
+void JniScanHeapGate::Permit::release() {
+    if (_gate != nullptr) {
+        _gate->_release(_bytes);
+        _gate = nullptr;
+        _bytes = 0;
+    }
+}
+
+JniScanHeapGate::JniScanHeapGate(std::function<int64_t()> budget) : 
_budget(std::move(budget)) {}
+
+JniScanHeapGate* JniScanHeapGate::instance() {
+    // Never destroyed: BE's exit runs static destructors while scanner 
threads are still alive, and
+    // a reader closing then releases its permit into this gate.
+    static auto* gate = new JniScanHeapGate(jvm_heap_budget);
+    return gate;
+}
+
+void JniScanHeapGate::acquire(int64_t bytes, const std::function<bool()>& 
stop_waiting,
+                              Permit* permit, int64_t* wait_ns) {
+    DORIS_CHECK(bytes > 0);
+    DORIS_CHECK(permit != nullptr);
+    DORIS_CHECK(!permit->held());
+    DORIS_CHECK(wait_ns != nullptr);
+    const int64_t start = MonotonicNanos();
+    std::unique_lock lock(_lock);
+    const uint64_t ticket = _next_ticket++;
+    _waiting.push_back(ticket);
+    while (true) {
+        // Asked without the lock: it is the caller's code.
+        lock.unlock();
+        const bool stop = stop_waiting();
+        lock.lock();
+        const int64_t waited = MonotonicNanos() - start;
+        bool admit = stop || _fits(ticket, bytes);
+        if (!admit && waited >= config::jni_scanner_heap_max_wait_ms * 1000 * 
1000) {
+            LOG_EVERY_T(WARNING, 10) << "A JNI scanner opens after waiting " 
<< waited / 1000 / 1000
+                                     << " ms for its share of the JVM heap, 
longer than "
+                                        "jni_scanner_heap_max_wait_ms: it 
declared "
+                                     << bytes / MB << " MB, while " << 
_holders << " scanners hold "
+                                     << _admitted_bytes / MB << " MB of a " << 
_budget() / MB
+                                     << " MB budget and " << _waiting.size() - 
1 << " others wait";
+            admit = true;
+        }
+        if (admit) {
+            _waiting.erase(std::find(_waiting.begin(), _waiting.end(), 
ticket));
+            _admitted_bytes += bytes;
+            ++_holders;
+            permit->_gate = this;
+            permit->_bytes = bytes;
+            *wait_ns = waited;
+            // Whoever is first in line now may fit.
+            _cv.notify_all();
+            return;
+        }
+        _cv.wait_for(lock, POLL_INTERVAL);
+    }

Review Comment:
   [P1] Keep heap waiters off the scan worker pool. An admitted scanner returns 
after one block with its Java reader and permit still open, and 
`ScannerContext` needs another worker turn to reach EOF/close. If other queries 
fill the finite shared scan pool in this synchronous wait, that holder cannot 
resume to release its permit; after the default 60 seconds the waiters all take 
the over-budget timeout path and may OOM the JVM. Yield and reschedule a scan 
on permit availability instead of blocking a worker here.



##########
be/src/util/jni_scan_heap_gate.cpp:
##########
@@ -0,0 +1,144 @@
+// 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.
+
+#include "util/jni_scan_heap_gate.h"
+
+#include <algorithm>
+#include <chrono>
+#include <limits>
+#include <utility>
+
+#include "common/config.h"
+#include "common/logging.h"
+#include "util/jni-util.h"
+#include "util/time.h"
+
+namespace doris {
+
+namespace {
+
+// How often a waiting reader looks again although nobody told it to: to see 
its query cancelled and
+// its wait run out.
+constexpr auto POLL_INTERVAL = std::chrono::milliseconds(100);
+constexpr int64_t MB = 1024 * 1024;
+
+// A share of the -Xmx the JVM was started with, read from the same options 
the JVM was created from
+// - not measured, so nothing here calls into the JVM.
+int64_t jvm_heap_budget() {
+    const double budget = 
static_cast<double>(Jni::Util::get_max_jni_heap_memory_size()) *
+                          config::jni_scanner_heap_budget_ratio;
+    // A BE_TEST build reports an unlimited heap (SIZE_MAX), which stays 
unlimited here.
+    constexpr auto UNLIMITED = std::numeric_limits<int64_t>::max();
+    return budget >= static_cast<double>(UNLIMITED) ? UNLIMITED : 
static_cast<int64_t>(budget);
+}
+
+} // namespace
+
+void JniScanHeapGate::Permit::release() {
+    if (_gate != nullptr) {
+        _gate->_release(_bytes);
+        _gate = nullptr;
+        _bytes = 0;
+    }
+}
+
+JniScanHeapGate::JniScanHeapGate(std::function<int64_t()> budget) : 
_budget(std::move(budget)) {}
+
+JniScanHeapGate* JniScanHeapGate::instance() {
+    // Never destroyed: BE's exit runs static destructors while scanner 
threads are still alive, and
+    // a reader closing then releases its permit into this gate.
+    static auto* gate = new JniScanHeapGate(jvm_heap_budget);
+    return gate;
+}
+
+void JniScanHeapGate::acquire(int64_t bytes, const std::function<bool()>& 
stop_waiting,
+                              Permit* permit, int64_t* wait_ns) {
+    DORIS_CHECK(bytes > 0);
+    DORIS_CHECK(permit != nullptr);
+    DORIS_CHECK(!permit->held());
+    DORIS_CHECK(wait_ns != nullptr);
+    const int64_t start = MonotonicNanos();
+    std::unique_lock lock(_lock);
+    const uint64_t ticket = _next_ticket++;
+    _waiting.push_back(ticket);
+    while (true) {
+        // Asked without the lock: it is the caller's code.
+        lock.unlock();
+        const bool stop = stop_waiting();
+        lock.lock();
+        const int64_t waited = MonotonicNanos() - start;
+        bool admit = stop || _fits(ticket, bytes);
+        if (!admit && waited >= config::jni_scanner_heap_max_wait_ms * 1000 * 
1000) {

Review Comment:
   [P2] Return cancellation instead of admitting a canceled reader. When 
`stop_waiting()` becomes true during a wait, this branch grants a real 
over-budget permit, and both `JniReader::open` and 
`JniTableReader::_open_jni_scanner` then unconditionally create and open the 
Java scanner. Thus canceling a query while other scanners fill the budget can 
launch another resource-owning scanner above it, even though the query has no 
work left. Have acquire report the stop and let the callers return before Java 
construction.



##########
fe/fe-connector/fe-connector-fluss/src/main/java/org/apache/doris/connector/fluss/FlussJniHeapEstimate.java:
##########
@@ -0,0 +1,127 @@
+// 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.connector.fluss;
+
+import org.apache.doris.connector.spi.scan.ConnectorScanRange;
+
+import org.apache.fluss.types.DataType;
+import org.apache.fluss.types.RowType;
+
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.List;
+
+/**
+ * The JVM heap BE's JNI reader of a fluss primary-key range holds, declared 
to BE's JNI heap gate
+ * ({@code TFileRangeDesc.jni_heap_bytes}) for a statement that sets {@code 
enable_jni_heap_admission}.
+ *
+ * <p>Two kinds of range keep rows in the heap until they close, and both get 
there before their first
+ * batch. A PK_FULL range replays the change log after its kv snapshot into a 
map ordered by key, which it
+ * then merges with the snapshot (SafeKvSnapshotAndLogBatchScanner); a key the 
log changed more than once
+ * keeps its first row besides its last, because the map keeps the first key 
object and that points at the
+ * first row. A PK_TAIL range keeps the last row of every key in its slice of 
the log (PkTailBatchScanner).
+ * Every row is the deep copy fluss makes of a fetched record: a GenericRow of 
boxed fields. So a PK_FULL
+ * range holds at most N x (R + 112) bytes and a PK_TAIL range N x (R + 96), N 
being the records between
+ * its offsets - more than its keys - and R a row. A LOG range streams; a lake 
range is paimon's to declare.
+ *
+ * <p>R follows from the column types, but for strings and bytes, whose length 
nothing in fluss's metadata
+ * gives: they are taken at {@link #DEFAULT_VARLEN_BYTES}.
+ */
+final class FlussJniHeapEstimate {
+
+    // A key's entry: in PK_FULL's TreeMap with the ProjectedRow standing for 
the key, in PK_TAIL's
+    // LinkedHashMap with the key encoded into a byte[].
+    static final long PK_FULL_ENTRY_BYTES = 112;
+    static final long PK_TAIL_ENTRY_BYTES = 96;
+    static final long DEFAULT_VARLEN_BYTES = 64;
+
+    private FlussJniHeapEstimate() {
+    }
+
+    /** A row of the fields at {@code fieldIndexes}: a GenericRow and its 
Object[], then every field. */
+    static long rowBytes(RowType rowType, Collection<Integer> fieldIndexes) {
+        long bytes = 16 + align8(16 + 4L * fieldIndexes.size());
+        for (int index : fieldIndexes) {
+            bytes += fieldBytes(rowType.getTypeAt(index));
+        }
+        return bytes;
+    }
+
+    static long fieldBytes(DataType type) {
+        switch (type.getTypeRoot()) {
+            case BOOLEAN:
+            case TINYINT:
+                // Boolean and Byte hand out cached instances.
+                return 0;
+            case SMALLINT:
+            case INTEGER:
+            case FLOAT:
+            case DATE:
+            case TIME_WITHOUT_TIME_ZONE:
+                return 16;
+            case BIGINT:
+            case DOUBLE:
+            case TIMESTAMP_WITHOUT_TIME_ZONE:
+            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+                return 24;
+            case DECIMAL:
+                // fluss's Decimal keeps the BigDecimal it was made from, and 
that its BigInteger.
+                return 136;
+            case BINARY:
+            case BYTES:
+                return align8(16 + DEFAULT_VARLEN_BYTES);
+            default:
+                // A string is a BinaryString over a MemorySegment over a 
byte[]; the nested types are

Review Comment:
   [P1] Bound the retained Fluss values and encoded keys. `fieldBytes` prices 
every STRING/BYTES or nested value at a fixed 64-byte payload, while PK_FULL 
and PK_TAIL retain complete rows and PK_TAIL also retains an encoded copy of 
each key. For 200,000 distinct 4 KiB STRING keys, this range declares about 61 
MiB although the encoded keys alone take about 819 MiB; two such tails pass a 1 
GiB gate while their rows and keys can exhaust a 2 GiB JVM. Use a conservative 
bound or a size ceiling for variable and nested values before declaring this 
amount.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonJniHeapEstimate.java:
##########
@@ -0,0 +1,104 @@
+// 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.connector.paimon;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.options.MemorySize;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.table.source.DataSplit;
+
+import java.util.Locale;
+import java.util.Map;
+
+/**
+ * The JVM heap BE's JNI reader of a paimon {@link DataSplit} holds, declared 
to BE's JNI heap gate
+ * ({@code TFileRangeDesc.jni_heap_bytes}) for a statement that sets {@code 
enable_jni_heap_admission}.
+ *
+ * <p>A JNI read of a split whose files overlap - a primary-key table written 
to since it was last
+ * compacted - merges them: paimon opens one file of every sorted run of a 
section at once, and an open
+ * parquet file keeps the compressed column chunks of its current row group in 
the heap until it moves
+ * on to the next, when the old and the new are there together for a moment. 
ORC keeps a stripe the
+ * same way. Sixteen such splits read at once are what runs a 2 GB heap out. 
So a split that has to be
+ * merged holds at most the sum over its files of one row group - two for a 
file that has more than one
+ * - plus a dictionary page for every column of every file. That takes all the 
files as one section,
+ * which is exact for the uncompacted tables that need the gate and over the 
mark for files that do not
+ * overlap, which are read one section after another. A split that needs no 
merging reads its files one
+ * after another and holds one file's share.
+ *
+ * <p>Every column is taken as read. A narrow projection reads less, but a 
merge also reads every key
+ * column, the sequence number and the row kind whatever is projected, and an 
estimate below what the
+ * reader holds is the one mistake the gate cannot absorb.
+ */
+final class PaimonJniHeapEstimate {
+
+    // What paimon's writers fall back to when neither file.block-size nor the 
format's own option is
+    // set (CoreOptions.FILE_BLOCK_SIZE).
+    static final long DEFAULT_PARQUET_ROW_GROUP_BYTES = 128L * 1024 * 1024;
+    static final long DEFAULT_ORC_STRIPE_BYTES = 64L * 1024 * 1024;
+    // The largest dictionary page paimon's parquet writer keeps for a column.
+    static final long DICTIONARY_BYTES_PER_COLUMN = 1024L * 1024;
+
+    private final long parquetRowGroupBytes;
+    private final long orcStripeBytes;
+    private final long dictionaryBytesPerFile;
+
+    PaimonJniHeapEstimate(long parquetRowGroupBytes, long orcStripeBytes, int 
columnsPerFile) {
+        this.parquetRowGroupBytes = parquetRowGroupBytes;
+        this.orcStripeBytes = orcStripeBytes;
+        this.dictionaryBytesPerFile = columnsPerFile * 
DICTIONARY_BYTES_PER_COLUMN;
+    }
+
+    static PaimonJniHeapEstimate of(Table table) {
+        Map<String, String> options = table.options();
+        // The precedence paimon's writers apply: file.block-size for every 
format, else the format's
+        // own option, else its default.
+        String blockSize = options.get(CoreOptions.FILE_BLOCK_SIZE.key());
+        long parquet = bytes(blockSize != null ? blockSize : 
options.get("parquet.block.size"),

Review Comment:
   [P1] Account for files written before a block-size option change. 
`of(table)` reads the current `file.block-size`/format option and 
`rowGroupsHeld` applies it to every existing file. Paimon permits changing 
table options without rewriting old files: after reducing a 128 MiB row-group 
setting to 16 MiB, a 900 MiB old file is charged at most 32 MiB plus 
dictionaries even though it can still contain 128 MiB groups. Concurrent JNI 
merge splits then pass the gate on an understated declaration and can OOM. 
Derive the bound from physical file metadata or a conservative historical 
maximum.



##########
be/src/util/jni_scan_heap_gate.cpp:
##########
@@ -0,0 +1,144 @@
+// 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.
+
+#include "util/jni_scan_heap_gate.h"
+
+#include <algorithm>
+#include <chrono>
+#include <limits>
+#include <utility>
+
+#include "common/config.h"
+#include "common/logging.h"
+#include "util/jni-util.h"
+#include "util/time.h"
+
+namespace doris {
+
+namespace {
+
+// How often a waiting reader looks again although nobody told it to: to see 
its query cancelled and
+// its wait run out.
+constexpr auto POLL_INTERVAL = std::chrono::milliseconds(100);
+constexpr int64_t MB = 1024 * 1024;
+
+// A share of the -Xmx the JVM was started with, read from the same options 
the JVM was created from
+// - not measured, so nothing here calls into the JVM.
+int64_t jvm_heap_budget() {
+    const double budget = 
static_cast<double>(Jni::Util::get_max_jni_heap_memory_size()) *
+                          config::jni_scanner_heap_budget_ratio;

Review Comment:
   [P2] Size the gate from the JVM's actual maximum heap. A nonempty 
`JAVA_OPTS` without `-Xmx` is passed to the JVM unchanged, but 
`get_max_jni_heap_memory_size()` assumes 1 GiB in that case. For example, 
`-XX:MaxRAMPercentage=10` under a 2 GiB container gives roughly a 205 MiB heap 
while this gate calculates a 512 MiB budget and can admit two 200 MiB readers. 
Query the running JVM's effective maximum instead of using the parser's 
HDFS-write fallback.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonJniHeapEstimate.java:
##########
@@ -0,0 +1,104 @@
+// 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.connector.paimon;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.options.MemorySize;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.table.source.DataSplit;
+
+import java.util.Locale;
+import java.util.Map;
+
+/**
+ * The JVM heap BE's JNI reader of a paimon {@link DataSplit} holds, declared 
to BE's JNI heap gate
+ * ({@code TFileRangeDesc.jni_heap_bytes}) for a statement that sets {@code 
enable_jni_heap_admission}.
+ *
+ * <p>A JNI read of a split whose files overlap - a primary-key table written 
to since it was last
+ * compacted - merges them: paimon opens one file of every sorted run of a 
section at once, and an open
+ * parquet file keeps the compressed column chunks of its current row group in 
the heap until it moves
+ * on to the next, when the old and the new are there together for a moment. 
ORC keeps a stripe the
+ * same way. Sixteen such splits read at once are what runs a 2 GB heap out. 
So a split that has to be
+ * merged holds at most the sum over its files of one row group - two for a 
file that has more than one
+ * - plus a dictionary page for every column of every file. That takes all the 
files as one section,
+ * which is exact for the uncompacted tables that need the gate and over the 
mark for files that do not
+ * overlap, which are read one section after another. A split that needs no 
merging reads its files one
+ * after another and holds one file's share.
+ *
+ * <p>Every column is taken as read. A narrow projection reads less, but a 
merge also reads every key
+ * column, the sequence number and the row kind whatever is projected, and an 
estimate below what the
+ * reader holds is the one mistake the gate cannot absorb.
+ */
+final class PaimonJniHeapEstimate {
+
+    // What paimon's writers fall back to when neither file.block-size nor the 
format's own option is
+    // set (CoreOptions.FILE_BLOCK_SIZE).
+    static final long DEFAULT_PARQUET_ROW_GROUP_BYTES = 128L * 1024 * 1024;
+    static final long DEFAULT_ORC_STRIPE_BYTES = 64L * 1024 * 1024;
+    // The largest dictionary page paimon's parquet writer keeps for a column.
+    static final long DICTIONARY_BYTES_PER_COLUMN = 1024L * 1024;
+
+    private final long parquetRowGroupBytes;
+    private final long orcStripeBytes;
+    private final long dictionaryBytesPerFile;
+
+    PaimonJniHeapEstimate(long parquetRowGroupBytes, long orcStripeBytes, int 
columnsPerFile) {
+        this.parquetRowGroupBytes = parquetRowGroupBytes;
+        this.orcStripeBytes = orcStripeBytes;
+        this.dictionaryBytesPerFile = columnsPerFile * 
DICTIONARY_BYTES_PER_COLUMN;
+    }
+
+    static PaimonJniHeapEstimate of(Table table) {
+        Map<String, String> options = table.options();
+        // The precedence paimon's writers apply: file.block-size for every 
format, else the format's
+        // own option, else its default.
+        String blockSize = options.get(CoreOptions.FILE_BLOCK_SIZE.key());
+        long parquet = bytes(blockSize != null ? blockSize : 
options.get("parquet.block.size"),
+                DEFAULT_PARQUET_ROW_GROUP_BYTES);
+        long orc = bytes(blockSize != null ? blockSize : 
options.get("orc.stripe.size"),
+                DEFAULT_ORC_STRIPE_BYTES);
+        // A primary-key table's data file stores the keys a second time, as 
_KEY_ columns, beside the
+        // sequence number and the row kind.
+        int keys = table.primaryKeys().size();
+        int columns = table.rowType().getFieldCount() + (keys > 0 ? keys + 2 : 
0);
+        return new PaimonJniHeapEstimate(parquet, orc, columns);
+    }
+
+    long bytesOf(DataSplit split) {
+        long total = 0;
+        long largest = 0;
+        for (DataFileMeta file : split.dataFiles()) {
+            long share = rowGroupsHeld(file) + dictionaryBytesPerFile;

Review Comment:
   [P1] Include the transient decoded batches in a Paimon permit. `bytesOf` 
sums row-group/file and dictionary allowances only, but this PR's Release-BE 
measurements report merge splits declaring 160-180 MiB and briefly holding 
474-560 MiB while decoded batches are alive. With the default 1 GiB budget, six 
170 MiB declarations fit although six measured 474 MiB peaks require about 2.8 
GiB, above the default 2 GiB JVM. The gate cannot prevent concurrent-split OOM 
on these measured cases without bounding that additional live heap.



##########
regression-test/suites/external_table_p0/fluss/test_fluss_jni_heap_admission.groovy:
##########
@@ -0,0 +1,94 @@
+// 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.
+
+// enable_jni_heap_admission has the connectors declare, on every range whose 
JNI reader holds much of
+// BE's JVM heap, how much it will hold, and BE opens those readers only while 
what it admitted fits
+// its budget. It is off by default. Turned on it may make a reader wait for 
room, and nothing else: the
+// rows must come back exactly as they do with it off.
+//
+// The reads that declare are here: fluss primary-key buckets read whole 
(PK_FULL), and a union read of a
+// primary-key table, whose tail is a PK_TAIL range and whose lake half the 
paimon connector plans. A log
+// read declares nothing and is here as the case that must not change either. 
Each query is recorded
+// with the variable on, and compared with the same query with it off; the 
comparison stays in the code
+// because what it asserts is the agreement.
+//
+// Fixtures come from docker/thirdparties/docker-compose/fluss/sql/init.sql 
and init-lake-tail.sql, and
+// are static: this suite never writes.
+suite("test_fluss_jni_heap_admission", "p0,external") {
+    String enabled = context.config.otherConfigs.get("enableFlussTest")
+    if (enabled == null || !enabled.equalsIgnoreCase("true")) {
+        return
+    }
+
+    String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+    String coordinatorPort = 
context.config.otherConfigs.get("fluss_coordinator_port")
+    String minioPort = context.config.otherConfigs.get("fluss_minio_port")
+    String bootstrapServers = "${externalEnvIp}:${coordinatorPort}"
+    String catalogName = "test_fluss_jni_heap_admission"
+
+    // Off unless a statement asks for it.
+    qt_default_off """show variables like 'enable_jni_heap_admission'"""
+
+    sql """drop catalog if exists ${catalogName}"""
+    // required: a union read that quietly fell back to fluss alone would 
never plan a lake split, and
+    // the paimon half of the declaration would go untested.
+    sql """
+        create catalog ${catalogName} properties (
+            "type" = "fluss",
+            "fluss.bootstrap.servers" = "${bootstrapServers}",
+            "fluss.lake.paimon.s3.endpoint" = 
"http://${externalEnvIp}:${minioPort}";,
+            "fluss.lake.paimon.s3.access-key" = "minioadmin",
+            "fluss.lake.paimon.s3.secret-key" = "minioadmin",
+            "fluss.union_read.mode" = "required"
+        );
+    """
+    sql """switch ${catalogName}"""
+    sql """use fluss_test"""
+    // The C++ glue exists only for the v2 file scanner, and the session 
variable that picks between
+    // them is randomised by the fuzzy mode this pipeline runs.
+    sql """set enable_file_scanner_v2 = true"""
+
+    def rowsOf = { String query -> sql(query).collect { row -> row.collect { 
it.toString() } } }
+    def sameWithAdmissionOff = { String query ->
+        sql """set enable_jni_heap_admission = false"""
+        def off = rowsOf(query)

Review Comment:
   [P2] Assert admission in the end-to-end regression. This helper compares 
rows with the option off and on, which must agree even if the FE never sends 
`jni_heap_bytes` or the BE never waits. The suite's fixtures do not create 
contention, and the expected output contains only rows/counts, so the new 
feature can be disconnected while this test stays green. Check positive 
`JvmHeapDeclaredBytes` for the PK and JNI lake cases and exercise a 
waiter/release path with a constrained budget.



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