morningman commented on code in PR #68713:
URL: https://github.com/apache/doris/pull/68713#discussion_r4179272813


##########
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:
   Confirmed, and kept out of this PR. A scan turn reads one block and gives 
its worker back (`ScannerScheduler::_scanner_scan`), the permit stays with the 
open Java scanner, and the holder is resubmitted only after its block is 
consumed, onto the same scheduler -- so waiters that fill the pool keep holders 
from closing until `jni_scanner_heap_max_wait_ms` lets everyone in.
   
   How far that reaches with default settings: file scans run on the workload 
group's remote scan scheduler, which grows its threads on demand up to 
`max(512, 10 x cores)` (`ScannerScheduler::default_remote_scan_thread_num()`), 
and a scan instance keeps at most `max_file_scanners_concurrency` (16) scanners 
in flight. A holder therefore finds a thread unless that many declaring readers 
are blocked on one BE at once -- several heavy queries together, or one query 
under a workload group with a small `max_remote_scan_thread_num`. The same 
parking has a cheaper trigger, for the record: in a join whose two sides both 
declare, the probe side's readers take permits and then wait for the build to 
finish before their blocks are consumed, so the build side's later readers can 
wait out the full 60 s. Either way the end state, after the stall, is the one 
the option being off has from the start: readers opened above the budget.
   
   Not blocking a worker means giving a scan a parked state the gate can end: 
`ScannerSplitRunner::process_for` runs one synchronous turn and returns a ready 
future today, and the thread-pool path has no state for a task waiting outside 
a worker. That changes how every scan is scheduled, not only JNI readers that 
declare, so I'd rather do it as a follow-up than grow this opt-in PR, which 
keeps the synchronous wait bounded by `jni_scanner_heap_max_wait_ms`.
   



##########
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:
   The arithmetic holds in magnitude: every STRING, BYTES or nested value is 
priced at a 64-byte payload, and a `PK_TAIL` entry at a fixed 96 bytes whatever 
the length of the key it encodes, so 200,000 distinct 4 KiB keys declare about 
61 MB against roughly 1.6 GB of strings and encoded keys.
   
   There is no conservative bound to switch to, though. Fluss's `STRING` and 
`BYTES` carry no length, and planning knows a range's offsets, not its bytes. A 
per-value ceiling large enough to be safe for any table would declare every 
primary-key range as if its values were huge and admit them one at a time -- 
the reads this option exists to let through; one small enough to keep them 
running is again not a bound. A bound that follows the data would have to come 
from BE as it reads, which is a different design.
   
   So the claim now matches the code instead (99e33c47bd1). The class comment 
no longer says a range holds "at most" its declaration: it says the declaration 
is an upper bound only for tables whose variable-length values and keys are 
short, and that a table with longer ones -- JSON in a string column, long 
string keys -- is declared below what its readers hold and is not protected by 
admission. A read of such a table that needs more heap than the JVM has fails 
as it does with the option off, the default, with an error that names the heap 
and how to raise it. The PR description says the same.
   



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