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


##########
be/src/exec/scan/scanner_context.cpp:
##########
@@ -362,6 +363,61 @@ void 
ScannerContext::push_completed_scan_task(std::shared_ptr<ScanTask> scan_tas
     _dependency->set_ready();
 }
 
+void ScannerContext::park_scan_task(std::shared_ptr<ScanTask> scan_task,
+                                    SharedListenableFuture<Void> waiting_for) {
+    {
+        std::unique_lock<std::mutex> l(_transfer_lock);
+        scan_task->set_state(ScanTask::State::PARKED);
+        _in_flight_tasks_num--;
+        ++_parked_tasks_num;
+        // The slot it held may now admit a pending scanner - one whose block 
the operator consumed
+        // while this context had no slot to spare, holding what this task 
waits for, among them.
+        // Nothing else would look: a parked task never completes a scan 
attempt to make the operator
+        // schedule again.
+        if (!done()) {
+            Status status = 
_scanner_scheduler->schedule_scan_task(shared_from_this(), nullptr, l);
+            if (!status.ok()) {
+                set_context_failure(status, l);
+            }
+        }
+    }
+    std::weak_ptr<ScannerContext> weak_ctx = shared_from_this();
+    // Runs on the thread that completes the future, or right here if it is 
done already.
+    waiting_for.add_callback([weak_ctx, scan_task](const Void&, const Status&) 
{
+        if (auto ctx = weak_ctx.lock()) {
+            ctx->_resume_parked_task(scan_task);
+        }
+    });
+}
+
+void ScannerContext::_resume_parked_task(const std::shared_ptr<ScanTask>& 
scan_task) {
+    auto task_execution_lock = task_exec_ctx();
+    if (task_execution_lock == nullptr) {
+        // The query has finished; nothing will read this task again.
+        return;
+    }
+#ifndef BE_TEST
+    // Scheduling allocates for the query, not for the gate's thread that 
completed the future. A
+    // future done before park_scan_task() added this callback runs it on the 
worker that parked
+    // the task, which is attached to this query already and must stay so.
+    std::optional<AttachTask> attach;
+    if (!thread_context()->is_attach_task()) {
+        attach.emplace(_state);
+    }
+#endif
+    std::unique_lock<std::mutex> l(_transfer_lock);
+    --_parked_tasks_num;
+    if (done()) {
+        return;
+    }
+    // Like a task whose block the operator just consumed: it may run again.
+    scan_task->set_state(ScanTask::State::PENDING);
+    Status status = _scanner_scheduler->schedule_scan_task(shared_from_this(), 
scan_task, l);

Review Comment:
   Not changing this one in this PR. Resubmitting a task before the worker that 
handed it over has returned is the contract the scan scheduler already has, and 
the one unsynchronized access in that tail belongs to `execute_scan_task()`, 
where the push path has it too.
   
   **The handoff.** `_scanner_scan()` ends with `push_completed_scan_task()`, 
and the operator resubmits a non-EOS task from `get_block_from_queue()` as soon 
as it takes `_transfer_lock` 
([scanner_context.cpp#L481-L482](https://github.com/apache/doris/blob/b389299bf207474d4aff74bdbd1eb7b0d171b781/be/src/exec/scan/scanner_context.cpp#L481-L482)).
 On TaskExecutor that is `re_enqueue_split` of the same `ScannerSplitRunner`, 
while the worker that pushed the task is still unwinding `_scanner_scan()`, 
`execute_scan_task()` and `process_for()`. A parked task is handed over at the 
same point. When the future is already done, the parking worker triggers the 
resubmission itself instead of the operator; what follows on that worker is the 
same tail.
   
   **The read.** `execute_scan_task()` returns `scan_task->is_eos()` after the 
handoff 
([scanner_scheduler.cpp#L111](https://github.com/apache/doris/blob/b389299bf207474d4aff74bdbd1eb7b0d171b781/be/src/exec/scan/scanner_scheduler.cpp#L111)),
 and `ScanTask::_state` is a plain enum. That is a data race by the letter of 
the standard, after `push_completed_scan_task()` as much as after 
`park_scan_task()`, and the comment at 
[scanner_scheduler.cpp#L341-L342](https://github.com/apache/doris/blob/b389299bf207474d4aff74bdbd1eb7b0d171b781/be/src/exec/scan/scanner_scheduler.cpp#L341-L342)
 overstates it in both cases. What it can do:
   - ThreadPool discards the value 
([simplified_scan_scheduler.cpp#L139](https://github.com/apache/doris/blob/b389299bf207474d4aff74bdbd1eb7b0d171b781/be/src/exec/scan/simplified_scan_scheduler.cpp#L139)).
   - On TaskExecutor it reads true only if the resumed turn has already ended 
the task, with EOS or an error, by the time the parking worker gets there, a 
few destructors after the handoff. A resumed reader is an admitted one: 
`try_stop()` comes only from `stop_scanners()`, which makes the context 
`done()` under `_transfer_lock`, and `_resume_parked_task()` checks `done()` 
under that lock before it schedules. So unless the query stops in that same 
instant, the resumed turn first attaches to the JVM and opens its Java scanner.
   - If it does read true, both turns report the same completion. The 
completion future ignores the second `set_value`, 
`PrioritizedSplitRunner::close()` runs once (`_closed.exchange`), and the extra 
`_split_finished()` counts the split twice in the per-level statistics and may 
move the adaptive concurrency target by one. The parking worker touches no 
scanner state after the handoff (`update_scanner_profile()` runs before it), 
and `PrioritizedSplitRunner::process()` keeps its own state in atomics or under 
`_priority_mutex`.
   
   So I don't see a P1 here, nor anything specific to parking. The fix, 
`_scanner_scan()` returning what it decided before the handoff and 
`execute_scan_task()` returning that instead of reading the task again, covers 
both handoffs and belongs in `execute_scan_task()` as a change of its own. 
Deferring the resubmission until the original invocation has exited would mean 
TaskExecutor not re-enqueueing a split that is still inside `process_for()`, 
which changes how every scan is scheduled, not only parked ones.
   



##########
be/src/util/jni_scan_heap_gate.cpp:
##########
@@ -0,0 +1,256 @@
+// 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/thread.h"
+#include "util/time.h"
+
+namespace doris {
+
+namespace {
+
+// How often the gate asks the waiting readers whether their scans stopped, 
and looks for those that
+// waited too long, although nobody gave a share back.
+constexpr int64_t POLL_INTERVAL_NS = 100L * 1000 * 1000;
+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:
   These two are read the way every mutable BE config is read. `DEFINE_m*` 
values are plain globals that BE reads without a lock everywhere: `config.cpp` 
defines 474 of them, and `_scanner_scan()` reads `doris_scanner_row_bytes` the 
same way. `set_config()` holds `mutable_string_config_lock` to keep the value, 
`full_conf_map` and the update callbacks in step ("Keep the value, config map, 
and callback in the same update order"); the only readers that take it are the 
ones that list the configs, the HTTP config pages and 
`information_schema.backend_configuration`.
   
   The closest neighbour does exactly this: 
`HdfsWriteMemUsageRecorder::max_usage()` scales the same JVM max heap by 
`max_hdfs_wirter_jni_heap_usage_ratio`, an mDouble, inside a `cv.wait_for` 
predicate without that lock 
([hdfs_file_writer.cpp#L105-L114](https://github.com/apache/doris/blob/b389299bf207474d4aff74bdbd1eb7b0d171b781/be/src/io/fs/hdfs_file_writer.cpp#L105-L114)),
 and its comment says the limit "may move freely, since nothing subtracts from 
it". The same holds here. The account (`_admitted_bytes`, `_holders`) moves 
only by each admission's own `_bytes` and is never derived from the budget, and 
the wait limit is compared with each waiter's own start time. A decision taken 
with the old value or the new one is what changing it online means; nothing can 
underflow or be left behind.
   
   Making only these two atomic would make them the only atomic mutable configs 
in BE, and reading them under `mutable_string_config_lock` would put a 
process-wide lock on every admission decision. If mutable configs should become 
race-free, that is a change to the config storage for all of them.
   



##########
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:
   Still as described above at b389299bf20. After the rebase, the class comment 
change I referred to is 568c7e2b221 (99e33c47bd1 before it).
   
   On a conservative bound, for the record: fluss 1.0 gives the planner no byte 
figure to bound with. `Admin.getTableStats()` returns a `TableStats` that holds 
only a row count, `getKvSnapshotMetadata()` only the snapshot's file names and 
its log offset, and `STRING` / `BYTES` carry no length. So the choice is still 
between a ceiling that serializes every primary-key read and one that is not a 
bound. The class comment and the description say which tables admission does 
not protect.
   



##########
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:
   Still as described above at b389299bf20. After the rebase, the comment 
change is 568c7e2b221 (99e33c47bd1 before it).
   
   The measured case is the one in the description's Results: 32 JNI merge 
splits read with 16 scanners on the default 2 GB heap ran out of heap in every 
round with the variable off, and in none of the rounds with it on. Six merge 
splits peaking at once is the worst case the class comment names. The headroom 
for it is a BE setting already: lowering `jni_scanner_heap_budget_ratio` leaves 
more of the heap outside the budget, without tripling every merge split's 
declaration for every reader.
   



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