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]
