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


##########
be/src/format_v2/jni/jni_table_reader.cpp:
##########
@@ -459,6 +472,55 @@ Status JniTableReader::_open_jni_scanner() {
     }
     _apply_common_scanner_params();
 
+    // Only the connector of a statement that set enable_jni_heap_admission 
declares the heap a
+    // reader will hold, and only for the readers that hold much; every other 
reader opens at once.
+    if (_current_range.__isset.jni_heap_bytes && _current_range.jni_heap_bytes 
> 0) {
+        // The scanner's try_stop() marks its IOContext: a cancelled query, a 
satisfied limit, a
+        // closing scan. The gate asks on a thread of its own, so the check 
owns what it reads.

Review Comment:
   Fixed in 730a221e300. Both `try_stop()` writes (`FileScanner`, 
`FileScannerV2`) and both callbacks the gate polls (`JniReader` on V1, 
`JniTableReader` on V2) now access `IOContext::should_stop` through 
`std::atomic_ref<bool>`, and `JniScanHeapGate::request()` says a `stop_waiting` 
callback must read what it reads atomically. The field stays a plain `bool`, 
since `IOContext` is copied by value in the IO layer; its other readers, on the 
scan workers, predate the gate and are not changed here. The ordering is 
relaxed because the flag publishes nothing else: the gate only uses it to 
decide that a waiting reader leaves without a share.
   



##########
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:
   Update: this went into the PR after all. 55f51167914 (cddb40bec0b before the 
rebase) parks a reader that has to wait instead of blocking a scan worker. Its 
scanner ends the turn without a block, `ScannerContext::park_scan_task()` takes 
the task off the workers and out of the in-flight count, and the gate's own 
thread completes the admission's future, whose callback schedules the task 
again. Both schedulers go through the context. The details and the end-to-end 
numbers are in the replies to the scheduler comment: 
https://github.com/apache/doris/pull/68713#issuecomment-6080975144 and 
https://github.com/apache/doris/pull/68713#issuecomment-6082410250.
   



##########
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:
   Fixed in this PR after all, in b389299bf20. Every file of a JNI split is now 
declared by the schema version it was written under 
(`DataFileMeta.schemaId()`): its row group size (`file.block-size`, else 
`parquet.block.size` or `orc.stripe.size`, else the default) and its columns 
come from that version's options and fields. The table's own version is at 
hand, and an older one is read once per planning from the table's 
`SchemaManager`, under the authenticator planning reads the manifests with. 
`PaimonJniHeapEstimateTest` covers a split of files written before and after an 
ALTER TABLE, and `PaimonScanPlanProviderTest` plans an altered table with the 
variable on and fails if the table's current options are used instead. What is 
left is a row group size that a write job passed as a dynamic option, which no 
schema version records; the class comment says so.
   



##########
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:
   After the rebase this is 1d8955875ee (1ce1680237d before it). Since 
55f51167914 there is no `acquire()` any more: `JniScanHeapGate::request()` 
returns an `Admission` at once, and a reader whose scan stops while it waits is 
settled as stopped by the gate's own thread. It takes no share and opens no 
Java scanner, and its parked task is dropped, the context being done. The tests 
named above are still there, unchanged in what they check: 
`JniScanHeapGateTest.AReaderWhoseScanStopsLeavesWithoutAShare`, 
`AReaderWhoseScanHasStoppedTakesNothingEvenWhenItFits`, 
`AStoppedReaderFirstInLineLetsTheNextOneIn`, and 
`JniTableReaderTest.ASplitWhoseScanStoppedOpensNoJavaScannerForTheHeapItDeclared`.
   



##########
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:
   After the rebase this is 710ae556384 (6de09492e88 before it); the checks are 
as described above.
   



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