morningman commented on code in PR #68713: URL: https://github.com/apache/doris/pull/68713#discussion_r4179273569
########## 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: Confirmed: `of(table)` takes today's `file.block-size` / format option and applies it to every file of the split, while paimon leaves old files as they were written when the option changes. It takes a lowered setting, and files from before it still uncompacted in a split read through JNI; raising the setting only over-declares. The precise fix needs no footer reads: every `DataFileMeta` carries the `schemaId` it was written under, and that schema version's options are what its writer used (short of options a write job passes dynamically). The connector already resolves historical schemas (`PaimonCatalogOps.schemaAt`), so taking the row-group size per file from there, cached by schema id, is a small change; I'm leaving it for a follow-up rather than this round. Until then such a file can be under-declared, and a read that then needs more heap than the JVM has fails as it does with the option off. ########## 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: The numbers are this PR's own: a merge split declaring 160-180 MB holds 150-170 MB steadily and peaked at 474-560 MB in 3 of 24 one-second samples (`GC.class_histogram`, so live objects, not garbage). The 2.8 GiB, though, needs six of those peaks at once, and the case they come from -- 32 JNI merge splits, 16 scanners, the default 1 GiB budget on a 2 GiB heap -- ran six times with the option on without an out-of-memory error. The gate does not guarantee against it; "cannot prevent" overstates it. The budget is half the heap so that the other half takes such peaks, along with the JVM's other users. Counting the peak instead would triple each merge split's declaration and cut the splits in flight to about a third -- from roughly nine to three on that table -- slowing the very reads this option exists for. Bounding the transient at its source would mean knowing why a few hundred decoded batches are alive at once, which I haven't pinned down. What was wrong is the comment's claim, fixed in 99e33c47bd1: it no longer says a declaration below what the reader holds is the one mistake the gate cannot absorb. It says the declaration is what a merge holds steadily, that the decode peaks, about three times that, are left to the heap outside the budget, which has room for some readers peaking at once but not for all, and that a read needing more heap than the JVM has fails as it does with the option off. 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]
