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]

Reply via email to