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]