github-actions[bot] commented on code in PR #66227: URL: https://github.com/apache/doris/pull/66227#discussion_r4122853148
########## fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java: ########## @@ -0,0 +1,245 @@ +// 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.datasource.paimon.source; + +import org.apache.doris.analysis.SlotDescriptor; +import org.apache.doris.analysis.TupleDescriptor; + +import com.google.common.collect.ImmutableSet; +import org.apache.paimon.CoreOptions; +import org.apache.paimon.io.DataFileMeta; +import org.apache.paimon.schema.TableSchema; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.source.DataSplit; +import org.apache.paimon.types.ArrayType; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataType; +import org.apache.paimon.types.DataTypeRoot; +import org.apache.paimon.types.DecimalType; +import org.apache.paimon.types.MapType; +import org.apache.paimon.types.RowType; +import org.apache.paimon.types.VarCharType; + +import java.util.Arrays; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + +/** Compatibility checks for the pinned paimon-rust reader, beyond storage capabilities. */ +final class PaimonRustReaderCapabilities { + private static final Set<String> SUPPORTED_AGGREGATE_NAMES = ImmutableSet.of( + "sum", "product", "min", "max", "last_value", "first_value", "last_non_null_value", + "first_non_null_value", "first_not_null_value", "bool_and", "bool_or", "listagg"); + private final FileStoreTable table; + private final TableSchema schema; + private final boolean tableCompatible; + private final Map<Long, Boolean> compatibleFileSchemas = new ConcurrentHashMap<>(); + + PaimonRustReaderCapabilities(FileStoreTable table, TupleDescriptor tuple) { + this.table = table; + this.schema = table.schema(); + this.tableCompatible = schema != null && hasFullNestedProjection(tuple) + && hasCompatibleAggregates(schema); + } + + boolean canRead(DataSplit split) { + if (!tableCompatible) { + return false; + } + // The pinned merge reader retains losing input batches until an output batch fills. + // Neither zero deletes nor a small read.batch-size bounds this across multiple files. + if (!schema.primaryKeys().isEmpty() && split.dataFiles().size() > 1) { + return false; + } + for (DataFileMeta file : split.dataFiles()) { + // Sort-merge retains consumed batches until it emits enough rows. Retracts can + // produce an unbounded zero-output prefix; aggregation also rejects retracts. + // Unknown counts must stay on JNI, including old files without this statistic. + if (!schema.primaryKeys().isEmpty() && file.deleteRowCount().orElse(-1L) != 0L) { + return false; + } + if (file.schemaId() != schema.id() && !compatibleFileSchemas.computeIfAbsent( + file.schemaId(), this::hasCompatibleFileSchema)) { + return false; + } + } + return true; + } + + private static boolean hasFullNestedProjection(TupleDescriptor tuple) { + for (SlotDescriptor slot : tuple.getSlots()) { + // The ABI projects only root names, while Arrow struct SerDes bind by ordinal. + // A pruned slot must use JNI's recursive read type, even inside arrays or maps. + if (slot.getType().isComplexType() && slot.getColumn() != null + && !slot.getType().equals(slot.getColumn().getType())) { + return false; + } + } + return true; + } + + private boolean hasCompatibleFileSchema(long id) { + try { + TableSchema fileSchema = table.schemaManager().schema(id); + // Match historical types by ID: renames and added fields do not require value casts. + return fileSchema != null && !hasIncompatibleEvolution(fileSchema.fields(), schema.fields()); + } catch (RuntimeException e) { + // Failure to establish compatibility must not opt a historical file into Rust. + return false; + } + } + + private static boolean hasIncompatibleEvolution(List<DataField> oldFields, List<DataField> newFields) { + Map<Integer, DataType> oldTypes = new HashMap<>(); + for (DataField field : oldFields) { + oldTypes.put(field.id(), field.type()); + } + for (DataField field : newFields) { + DataType oldType = oldTypes.get(field.id()); + if (oldType != null && hasIncompatibleEvolution(oldType, field.type())) { + return true; + } + } + return false; + } + + private static boolean hasIncompatibleEvolution(DataType oldType, DataType newType) { + // Java rounds the decimal string of a floating value; Arrow scales the binary value. + // Even an in-range value such as DOUBLE 1.005 can therefore round differently. + if (newType instanceof DecimalType && (oldType.getTypeRoot() == DataTypeRoot.FLOAT + || oldType.getTypeRoot() == DataTypeRoot.DOUBLE)) { + return true; + } + int oldWidth = integerWidth(oldType.getTypeRoot()); + int newWidth = integerWidth(newType.getTypeRoot()); + if (newWidth > 0) { + // Floating-point and decimal sources also differ under narrowing; only integer + // identity/widening casts have the same range and value semantics in both readers. + return oldWidth == 0 || oldWidth > newWidth; + } + if (oldType instanceof RowType && newType instanceof RowType) { + return hasIncompatibleEvolution(((RowType) oldType).getFields(), ((RowType) newType).getFields()); + } + if (oldType instanceof ArrayType && newType instanceof ArrayType) { + return hasIncompatibleEvolution(((ArrayType) oldType).getElementType(), + ((ArrayType) newType).getElementType()); + } + if (oldType instanceof MapType && newType instanceof MapType) { + MapType oldMap = (MapType) oldType; + MapType newMap = (MapType) newType; + return hasIncompatibleEvolution(oldMap.getKeyType(), newMap.getKeyType()) + || hasIncompatibleEvolution(oldMap.getValueType(), newMap.getValueType()); + } + return false; Review Comment: [P1] Keep integer-to-TIMESTAMP schema evolution on JNI. The gate's fallthrough at this line treats an old INTEGER/BIGINT field as compatible because TIMESTAMP has no integer width. In the pinned paimon-rust v0.4.0-rc1, nested evolution reaches Arrow's direct integer-to-timestamp cast, which reuses the raw value as timestamp ticks; Java Paimon's NumericPrimitiveToTimestamp instead interprets the value as epoch seconds and scales it (for example, 1700000000 should be in 2023, not near 1970). A historical split admitted by canRead() can therefore return a different timestamp without an error. Reject integer-to-TIMESTAMP (including nested occurrences) until the dependency with the scaling fix is pinned, and add a persisted INTEGER/BIGINT-to-TIMESTAMP differential at precisions 0/3/6. -- 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]
