github-actions[bot] commented on code in PR #66227:
URL: https://github.com/apache/doris/pull/66227#discussion_r4124160423


##########
be/src/format_v2/table/paimon_rust_predicate_converter.cpp:
##########
@@ -0,0 +1,836 @@
+// 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 "format_v2/table/paimon_rust_predicate_converter.h"
+
+#include <algorithm>
+#include <cctype>
+#include <memory>
+#include <utility>
+
+#include "common/logging.h"
+#include "core/column/column_const.h"
+#include "core/column/column_nullable.h"
+#include "core/data_type/data_type.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/field.h"
+#include "core/types.h"
+#include "core/value/decimalv2_value.h"
+#include "core/value/timestamptz_value.h"
+#include "core/value/vdatetime_value.h"
+#include "exprs/runtime_filter_expr.h"
+#include "exprs/vcompound_pred.h"
+#include "exprs/vdirect_in_predicate.h"
+#include "exprs/vectorized_fn_call.h"
+#include "exprs/vexpr.h"
+#include "exprs/vin_predicate.h"
+#include "exprs/vliteral.h"
+#include "exprs/vslot_ref.h"
+
+namespace doris {
+
+namespace {
+// paimon_datum tags (see paimon.h / bindings/c/src/table.rs::datum_from_c).
+constexpr int32_t kTagBool = 0;
+constexpr int32_t kTagTinyInt = 1;
+constexpr int32_t kTagSmallInt = 2;
+constexpr int32_t kTagInt = 3;
+constexpr int32_t kTagLong = 4;
+constexpr int32_t kTagDouble = 6;
+constexpr int32_t kTagString = 7;
+constexpr int32_t kTagDate = 8;
+constexpr int32_t kTagTimestamp = 10;
+constexpr int32_t kTagDecimal = 12;
+constexpr int32_t kTagBytes = 13;
+
+// paimon decimal precision ceiling (paimon::Decimal::MAX_PRECISION).
+constexpr int32_t kPaimonDecimalMaxPrecision = 38;
+
+// RAII for an owned paimon_predicate*. and/or/not consume their inputs, so we
+// release() before handing pointers to them.
+struct predicate_deleter {
+    void operator()(paimon_predicate* p) const {
+        if (p) {
+            paimon_predicate_free(p);
+        }
+    }
+};
+using predicate_ptr = std::unique_ptr<paimon_predicate, predicate_deleter>;
+
+// RAII for an owned paimon_error*.
+struct error_deleter {
+    void operator()(paimon_error* p) const {
+        if (p) {
+            paimon_error_free(p);
+        }
+    }
+};
+using error_ptr = std::unique_ptr<paimon_error, error_deleter>;
+
+// Render a paimon_error into a string. Takes ownership of `err` via RAII so it
+// is freed on every return path. Safe to call with nullptr.
+std::string consume_predicate_error(paimon_error* err) {
+    error_ptr owned(err);
+    if (!owned) {
+        return "unknown error";
+    }
+    std::string msg;
+    if (owned->message.data != nullptr && owned->message.len > 0) {
+        msg.assign(reinterpret_cast<const char*>(owned->message.data), 
owned->message.len);
+    }
+    return "code=" + std::to_string(owned->code) + ", msg=" + msg;
+}
+} // namespace
+
+PaimonRustPredicateConverter::PaimonRustPredicateConverter(
+        const std::vector<std::string>& column_names, const 
std::vector<DataTypePtr>& column_types,
+        const paimon_table* table)
+        : _table(table) {
+    DORIS_CHECK(column_names.size() == column_types.size());
+    _columns_by_name.reserve(column_names.size());
+    for (size_t i = 0; i < column_names.size(); ++i) {
+        _columns_by_name.emplace(_normalize_name(column_names[i]),
+                                 std::make_pair(column_names[i], 
column_types[i]));
+    }
+    // Paimon TIMESTAMP (wall clock) is stored as epoch-millis-of-the-wall-time
+    // and the DateTimeV2 serde decodes timezone-naive arrow values in UTC, so
+    // timestamp literals convert wall->epoch in UTC. utc_time_zone() needs no
+    // tzdata lookup, so the conversion cannot silently fall back to a
+    // machine-local zone.
+    _utc_tz = cctz::utc_time_zone();
+}
+
+paimon_predicate* PaimonRustPredicateConverter::build(const VExprContextSPtrs& 
conjuncts) {
+    _converted_conjuncts = 0;
+    _converted_runtime_filters = 0;
+    if (_table == nullptr) {
+        return nullptr;
+    }
+    predicate_ptr result;
+    for (const auto& conjunct : conjuncts) {
+        if (!conjunct || !conjunct->root()) {
+            continue;
+        }
+        auto root = conjunct->root();
+        const bool is_runtime_filter = root->is_rf_wrapper();
+        if (root->is_rf_wrapper()) {
+            if (auto impl = root->get_impl()) {
+                // A null-aware runtime filter (an EQ_FOR_NULL join) must stay
+                // residual: its wrapper execution restores NULL probe rows to
+                // true (RuntimeFilterExpr::change_null_to_true), while the
+                // unwrapped impl — rebuilt as an ordinary IN predicate through
+                // VDirectInPredicate::get_slot_in_expr — treats NULL as
+                // not-in-set and would prune the NULL probes permanently
+                // before the join sees them. Keep the wrapper itself: it fails
+                // every dispatch below, so the conjunct stays in the residual,
+                // and it is safe to execute on selected rows, so later
+                // conjuncts keep pushing. The lance pushdown declines
+                // is_null_aware() filters for the same reason.
+                // is_null_aware() is concrete on RuntimeFilterExpr (the only
+                // class whose is_rf_wrapper() is true), so the dynamic_cast
+                // never fails in practice; the null guard keeps the unwrap for
+                // any future wrapper shape.
+                auto* rf_wrapper = 
dynamic_cast<RuntimeFilterExpr*>(root.get());
+                if (rf_wrapper == nullptr || !rf_wrapper->is_null_aware()) {
+                    root = impl;
+                }
+            }
+        }
+        // Preserve a safe prefix of the conjunct order: a later pushed
+        // predicate (e.g. an arrived IN runtime filter) could otherwise prune
+        // rows on which an earlier error-preserving conjunct —
+        // assert_true(...), a failing cast, ... — must still raise. The v1
+        // partition-pruning path (FileScanner::_init_runtime_filter_partition_
+        // prune_ctxs) stops at is_safe_to_execute_on_selected_rows() for the
+        // same reason, so a convertible predicate after an unsafe conjunct
+        // must not be pushed. Safe conjuncts that cannot be converted keep
+        // the old skip: they cannot raise, so pruning rows before they are
+        // evaluated as the residual never loses an error.
+        if (!_is_safe_to_push(root)) {
+            break;
+        }
+        predicate_ptr pred(_convert_expr(root));
+        if (!pred) {
+            continue;
+        }
+        ++_converted_conjuncts;
+        _converted_runtime_filters += is_runtime_filter;
+        if (!result) {
+            result = std::move(pred);
+        } else {
+            // and consumes both inputs regardless of success.
+            result.reset(paimon_predicate_and(result.release(), 
pred.release()));
+            if (!result) {
+                _converted_conjuncts = 0;
+                _converted_runtime_filters = 0;
+                return nullptr;
+            }
+        }
+    }
+    return result.release();
+}
+
+bool PaimonRustPredicateConverter::_is_safe_to_push(const VExprSPtr& expr) {
+    if (expr->is_safe_to_execute_on_selected_rows()) {
+        return true;
+    }
+    // VectorizedFnCall::is_safe_to_execute_on_selected_rows() admits a fixed
+    // whitelist that does not include `like`, so a like conjunct would always
+    // stop the safe prefix here and never reach _convert_like. A like call
+    // whose children are themselves safe is just as total as the whitelisted
+    // comparisons: it only compares strings, and its regex is built by 
escaping
+    // the pattern operand (FunctionLike::convert_like_pattern), so no scanned
+    // value can make it raise — a malformed pattern or escape operand fails
+    // at open() on every execution, regardless of which rows survive. Admit
+    // it here rather than widening the generic whitelist, which gates shared
+    // BE pushdown paths outside this reader's scope.
+    auto* fn = dynamic_cast<VectorizedFnCall*>(expr.get());
+    if (fn == nullptr || _normalize_name(fn->function_name()) != "like") {
+        return false;
+    }
+    for (uint16_t i = 0; i < fn->get_num_children(); ++i) {
+        if (!fn->get_child(i)->is_safe_to_execute_on_selected_rows()) {
+            return false;
+        }
+    }
+    return true;
+}
+
+paimon_predicate* PaimonRustPredicateConverter::_convert_expr(const VExprSPtr& 
expr) {
+    if (!expr) {
+        return nullptr;
+    }
+
+    // Casts are not unwrapped anywhere (predicate root included): a cast node
+    // fails every dispatch below and the conjunct stays in the Doris residual,
+    // mirroring the FE converter, which keeps casted expressions unconverted.
+    if (auto* direct_in = dynamic_cast<VDirectInPredicate*>(expr.get())) {
+        VExprSPtr in_expr;
+        if (direct_in->get_slot_in_expr(in_expr)) {
+            return _convert_in(in_expr);
+        }
+        return nullptr;
+    }
+
+    if (dynamic_cast<VInPredicate*>(expr.get()) != nullptr) {
+        return _convert_in(expr);
+    }
+
+    switch (expr->op()) {
+    case TExprOpcode::COMPOUND_AND:
+    case TExprOpcode::COMPOUND_OR:
+        return _convert_compound(expr);
+    case TExprOpcode::COMPOUND_NOT:
+        return nullptr;
+    case TExprOpcode::EQ:
+    case TExprOpcode::EQ_FOR_NULL:
+    case TExprOpcode::NE:
+    case TExprOpcode::GE:
+    case TExprOpcode::GT:
+    case TExprOpcode::LE:
+    case TExprOpcode::LT:
+        return _convert_binary(expr);
+    default:
+        break;
+    }
+
+    if (auto* fn = dynamic_cast<VectorizedFnCall*>(expr.get())) {
+        auto fn_name = _normalize_name(fn->function_name());
+        if (fn_name == "is_null_pred" || fn_name == "is_not_null_pred") {
+            return _convert_is_null(expr, fn_name);
+        }
+        if (fn_name == "like") {
+            return _convert_like(expr);
+        }
+    }
+
+    return nullptr;
+}
+
+paimon_predicate* PaimonRustPredicateConverter::_convert_compound(const 
VExprSPtr& expr) {
+    if (!expr || expr->get_num_children() != 2) {
+        return nullptr;
+    }
+    predicate_ptr left(_convert_expr(expr->get_child(0)));
+    if (!left) {
+        return nullptr;
+    }
+    predicate_ptr right(_convert_expr(expr->get_child(1)));
+    if (!right) {
+        return nullptr;
+    }
+
+    if (expr->op() == TExprOpcode::COMPOUND_AND) {
+        return paimon_predicate_and(left.release(), right.release());
+    }
+    if (expr->op() == TExprOpcode::COMPOUND_OR) {
+        return paimon_predicate_or(left.release(), right.release());
+    }
+    return nullptr;
+}
+
+paimon_predicate* PaimonRustPredicateConverter::_convert_in(const VExprSPtr& 
expr) {
+    auto* in_pred = dynamic_cast<VInPredicate*>(expr.get());
+    if (!in_pred || expr->get_num_children() < 2) {
+        return nullptr;
+    }
+    auto field_meta = _resolve_field(expr->get_child(0));
+    if (!field_meta) {
+        return nullptr;
+    }
+
+    const auto num_values = expr->get_num_children() - 1;
+    // Reserve up front so the backing strings never reallocate: each datum's
+    // str_data points into storages[i], which must stay stable.
+    std::vector<std::string> storages;
+    std::vector<paimon_datum> datums;
+    storages.reserve(num_values);
+    datums.reserve(num_values);
+    for (uint16_t i = 1; i < expr->get_num_children(); ++i) {

Review Comment:
   [P1] Reject generated IN sets before the 16-bit child count wraps. With 
runtime_filter_max_in_num=70000, an arrived IN runtime filter containing 65,537 
distinct keys is expanded by VDirectInPredicate::get_slot_in_expr into 65,538 
children. VExpr::get_num_children() truncates that to 2, so this loop pushes 
only the first key to paimon-rust; a probe row matching any other build key is 
discarded before Doris's residual/join can see it. Inspect children().size() 
and leave oversized filters residual (or use a size_t-safe path), with a 
differential at the 65,535/65,536 boundary.



##########
be/src/format_v2/table/paimon_rust_table_reader.cpp:
##########
@@ -0,0 +1,916 @@
+// 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 "format_v2/table/paimon_rust_table_reader.h"
+
+#include <algorithm>
+#include <utility>
+
+#include "arrow/c/abi.h"
+#include "arrow/c/bridge.h"
+#include "arrow/record_batch.h"
+#include "arrow/result.h"
+#include "common/logging.h"
+#include "core/assert_cast.h"
+#include "core/block/block.h"
+#include "core/block/column_with_type_and_name.h"
+#include "core/column/column_const.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_string.h"
+#include "exprs/vexpr_context.h"
+#include "exprs/vliteral.h"
+#include "format_v2/column_mapper.h"
+#include "format_v2/table/paimon_rust_predicate_converter.h"
+#include "runtime/descriptors.h"
+#include "runtime/file_scan_profile.h"
+#include "runtime/runtime_state.h"
+#include "util/string_util.h"
+#include "util/timezone_utils.h"
+#include "util/url_coding.h"
+
+extern "C" {
+#include "paimon_rust/paimon.h"
+}
+
+namespace doris::format::paimon {
+
+namespace {
+constexpr const char* VALUE_KIND_FIELD = "_VALUE_KIND";
+
+// ---------------------------------------------------------------------------
+// RAII wrappers over the paimon-rust C handles. Each handle is an opaque
+// pointer owned by Rust and released by a matching paimon_*_free function.
+// ---------------------------------------------------------------------------
+#define PAIMON_OWNED(type, freefn)                \
+    struct type##_deleter {                       \
+        void operator()(paimon_##type* p) const { \
+            if (p) {                              \
+                freefn(p);                        \
+            }                                     \
+        }                                         \
+    };                                            \
+    using type##_ptr = std::unique_ptr<paimon_##type, type##_deleter>
+
+PAIMON_OWNED(table, paimon_table_free);
+PAIMON_OWNED(read_builder, paimon_read_builder_free);
+PAIMON_OWNED(plan, paimon_plan_free);
+PAIMON_OWNED(table_read, paimon_table_read_free);
+PAIMON_OWNED(record_batch_reader, paimon_record_batch_reader_free);
+PAIMON_OWNED(error, paimon_error_free);
+
+#undef PAIMON_OWNED
+
+// One Arrow batch (schema + array containers). Owning it requires a two-step
+// teardown that the unique_ptr deleters above can't express: first invoke the
+// Arrow C Data Interface `release` callback on each struct (hands buffers back
+// to the producer), then free the container structs via 
paimon_arrow_batch_free.
+class ArrowBatch {
+public:
+    explicit ArrowBatch(paimon_arrow_batch batch) : batch_(batch) {}
+    ~ArrowBatch() {
+        auto* schema = static_cast<ArrowSchema*>(batch_.schema);
+        auto* array = static_cast<ArrowArray*>(batch_.array);
+        if (array && array->release) {
+            array->release(array);
+        }
+        if (schema && schema->release) {
+            schema->release(schema);
+        }
+        paimon_arrow_batch_free(batch_);
+    }
+
+    ArrowBatch(const ArrowBatch&) = delete;
+    ArrowBatch& operator=(const ArrowBatch&) = delete;
+
+    ArrowSchema* schema() const { return 
static_cast<ArrowSchema*>(batch_.schema); }
+    ArrowArray* array() const { return static_cast<ArrowArray*>(batch_.array); 
}
+
+private:
+    paimon_arrow_batch batch_;
+};
+
+// Render a paimon_error into a string. Takes ownership of `err` via RAII so it
+// is freed on every return path. Safe to call with nullptr.
+std::string consume_error(paimon_error* err) {
+    error_ptr owned(err);
+    if (!owned) {
+        return "unknown error";
+    }
+    std::string msg;
+    if (owned->message.data != nullptr && owned->message.len > 0) {
+        msg.assign(reinterpret_cast<const char*>(owned->message.data), 
owned->message.len);
+    }
+    return "code=" + std::to_string(owned->code) + ", msg=" + msg;
+}
+
+// Render storage option KEYS for diagnostics. Values are never rendered:
+// credential keys arrive under many spellings and cases (AWS_SECRET_KEY,
+// AWS_TOKEN, fs.oss.accessKeySecret, s3.secret-key, ...), and a key-name
+// blocklist that misses one alias leaks the value into the INFO log, so
+// only the key names are printed at all.
+std::string format_options(const std::map<std::string, std::string>& options) {
+    std::string out;
+    for (const auto& kv : options) {
+        if (!out.empty()) {
+            out += ", ";
+        }
+        out += kv.first;
+    }
+    return out;
+}
+
+} // namespace
+
+// Paimon-rust handles. Order of members matters: destruction runs in reverse
+// declaration order, and the read_builder depends on the table while the arrow
+// reader depends on the whole pipeline above it. So the table MUST be declared
+// first (destroyed last) and the record batch reader last.
+struct PaimonRustTableReader::PaimonHandles {
+    table_ptr table;
+    read_builder_ptr read_builder;
+    plan_ptr plan;
+    table_read_ptr table_read;
+    record_batch_reader_ptr reader;
+};
+
+PaimonRustTableReader::PaimonRustTableReader() = default;
+
+PaimonRustTableReader::~PaimonRustTableReader() = default;
+
+Status PaimonRustTableReader::init(format::TableReadOptions&& options) {
+    RETURN_IF_ERROR(format::TableReader::init(std::move(options)));
+    {
+        // Base and derived scopes must not overlap on the same counter: 
RuntimeProfile timers
+        // add deltas, so nested use would double-count instead of extending 
lifecycle coverage.
+        SCOPED_TIMER(_profile.total_timer);
+        SCOPED_TIMER(_profile.init_timer);
+        // Materialize TIMESTAMP_LTZ in the session timezone — the same
+        // convention as the JNI reader (PaimonJniScanner reads time_zone from
+        // its scan params) and lance_reader. Timezone-naive (paimon TIMESTAMP)
+        // arrow values are decoded in UTC by the DateTimeV2 serde regardless
+        // of _ctz, so NTZ wall-clock semantics are preserved.
+        DORIS_CHECK(_runtime_state != nullptr);
+        _ctz = _runtime_state->timezone_obj();
+        if (_scanner_profile != nullptr) {
+            file_scan_profile::ensure_hierarchy(_scanner_profile);
+            _rust_total_time = ADD_CHILD_TIMER(_scanner_profile, 
"PaimonRustReader",
+                                               
file_scan_profile::TABLE_READER);
+            _rust_open_split_time =
+                    ADD_CHILD_TIMER(_scanner_profile, "OpenSplitTime", 
"PaimonRustReader");
+            _rust_read_batch_time =
+                    ADD_CHILD_TIMER(_scanner_profile, "ReadBatchTime", 
"PaimonRustReader");
+            _rust_arrow_to_block_time =
+                    ADD_CHILD_TIMER(_scanner_profile, "ArrowToBlockTime", 
"PaimonRustReader");
+            _rust_predicates_input = ADD_CHILD_COUNTER(_scanner_profile, 
"RustPredicatesInput",
+                                                       TUnit::UNIT, 
"PaimonRustReader");
+            _rust_predicates_converted = ADD_CHILD_COUNTER(
+                    _scanner_profile, "RustPredicatesConverted", TUnit::UNIT, 
"PaimonRustReader");
+            _rust_predicates_applied = ADD_CHILD_COUNTER(_scanner_profile, 
"RustPredicatesApplied",
+                                                         TUnit::UNIT, 
"PaimonRustReader");
+            _rust_runtime_filters_input = ADD_CHILD_COUNTER(
+                    _scanner_profile, "RustRuntimeFiltersInput", TUnit::UNIT, 
"PaimonRustReader");
+            _rust_runtime_filters_applied = ADD_CHILD_COUNTER(
+                    _scanner_profile, "RustRuntimeFiltersApplied", 
TUnit::UNIT, "PaimonRustReader");
+        }
+        // Projected column name -> fixed output position, registered with 
both the exact and
+        // the lower-case spelling so mixed-case Rust schema output still 
resolves (v1
+        // semantics: exact match first, lower-case fallback on lookup).
+        _output_name_to_idx.reserve(_projected_columns.size() * 2);
+        for (size_t idx = 0; idx < _projected_columns.size(); ++idx) {
+            _output_name_to_idx.emplace(_projected_columns[idx].name, idx);
+            
_output_name_to_idx.emplace(to_lower(_projected_columns[idx].name), idx);
+        }
+    }
+    return Status::OK();
+}
+
+Status PaimonRustTableReader::prepare_split(const format::SplitReadOptions& 
options) {
+    // EOF belongs to the previous split. Keep it set after closing that split 
so repeated reads
+    // are idempotent, and clear it only when a new split is explicitly 
prepared.
+    _close_split_reader();
+    _split_eof = false;
+    _current_range = options.current_range;
+    RETURN_IF_ERROR(format::TableReader::prepare_split(options));
+    if (current_split_pruned()) {
+        return Status::OK();
+    }
+    if (_is_table_level_count_active()) {
+        // No rust pipeline is opened; get_block emits the synthetic count 
rows.
+        return Status::OK();
+    }
+    RETURN_IF_ERROR(_validate_rust_split(options.current_range));
+    {
+        SCOPED_TIMER(_profile.total_timer);
+        SCOPED_TIMER(_profile.prepare_split_timer);
+        SCOPED_TIMER(_rust_open_split_time);

Review Comment:
   [P2] Include split opening in the PaimonRustReader total timer. 
OpenSplitTime is registered under that total, but this scope times only 
OpenSplitTime while _rust_total_time is scoped only in get_block(). A slow or 
failed _open_split_reader can therefore report OpenSplitTime greater than its 
parent (or a zero parent), obscuring where scan time went. Scope the Rust total 
around the open as well and cover an open-heavy profile.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,257 @@
+// 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) {

Review Comment:
   [P1] Bound the single-file primary-key merge before allowing this split. 
Zero deletes and one data file do not imply unique keys: pinned paimon-rust 
explicitly routes single files through sort-merge because older files can 
contain duplicate keys. Its merge loop accumulates every row in one same-key 
run and retains each source RecordBatch until the key changes, so a large run 
can exhaust BE memory before producing one row; reducing read.batch-size only 
creates more retained batches. JNI releases consumed batches. Keep this shape 
on JNI or make the Rust merge incremental/spillable, and cover a one-file, 
zero-delete repeated-key fixture.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,257 @@
+// 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());
+        DataTypeRoot oldRoot = oldType.getTypeRoot();
+        DataTypeRoot newRoot = newType.getTypeRoot();
+        boolean targetTimestamp = newRoot == 
DataTypeRoot.TIMESTAMP_WITHOUT_TIME_ZONE
+                || newRoot == DataTypeRoot.TIMESTAMP_WITH_LOCAL_TIME_ZONE;
+        // Java accepts epoch-day/millisecond strings and interprets numeric 
timestamps as seconds.
+        // Arrow instead parses calendar strings or reuses numeric timestamp 
ticks, including in
+        // nested values. Keep these historical casts on JNI until the pinned 
reader matches Java.
+        boolean sourceString = oldRoot == DataTypeRoot.CHAR || oldRoot == 
DataTypeRoot.VARCHAR;
+        if ((sourceString && (newRoot == DataTypeRoot.DATE || targetTimestamp))
+                || (oldWidth > 0 && targetTimestamp)) {
+            return true;
+        }
+        int newWidth = integerWidth(newRoot);
+        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 historical TIMESTAMP-to-DATE files on JNI until this cast matches 
Paimon. The gate admits an old TIMESTAMP field changed to DATE. Paimon 1.4.2's 
Java cast computes `(int)(epochMillis / 86400000)`, so a stored 1969-12-31 
23:59:59.999 (-1 ms) becomes 1970-01-01; the pinned Rust reader delegates to 
Arrow's calendar-date cast and returns 1969-12-31. This changes projected 
values even without a filter. Add a persisted pre-epoch differential for a root 
column (and nested form if supported).



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,257 @@
+// 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());
+        DataTypeRoot oldRoot = oldType.getTypeRoot();
+        DataTypeRoot newRoot = newType.getTypeRoot();
+        boolean targetTimestamp = newRoot == 
DataTypeRoot.TIMESTAMP_WITHOUT_TIME_ZONE

Review Comment:
   [P1] Gate historical NTZ-to-LTZ conversion until Rust uses the Java cast 
timezone. An old TIMESTAMP_WITHOUT_TIME_ZONE column changed to 
TIMESTAMP_WITH_LOCAL_TIME_ZONE passes this check. Java Paimon interprets the 
old civil timestamp in `TimeZone.getDefault()`; the pinned Rust reader casts to 
Arrow Timestamp(UTC). On an Asia/Shanghai host, stored 2024-01-01 00:00:00 
becomes 2023-12-31 16:00:00Z through JNI but 2024-01-01 00:00:00Z through Rust. 
This applies to ordinary Parquet data, apart from the existing ORC LTZ thread. 
Add a non-UTC persisted differential.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,257 @@
+// 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());
+        DataTypeRoot oldRoot = oldType.getTypeRoot();
+        DataTypeRoot newRoot = newType.getTypeRoot();
+        boolean targetTimestamp = newRoot == 
DataTypeRoot.TIMESTAMP_WITHOUT_TIME_ZONE
+                || newRoot == DataTypeRoot.TIMESTAMP_WITH_LOCAL_TIME_ZONE;
+        // Java accepts epoch-day/millisecond strings and interprets numeric 
timestamps as seconds.
+        // Arrow instead parses calendar strings or reuses numeric timestamp 
ticks, including in
+        // nested values. Keep these historical casts on JNI until the pinned 
reader matches Java.
+        boolean sourceString = oldRoot == DataTypeRoot.CHAR || oldRoot == 
DataTypeRoot.VARCHAR;
+        if ((sourceString && (newRoot == DataTypeRoot.DATE || targetTimestamp))
+                || (oldWidth > 0 && targetTimestamp)) {
+            return true;
+        }
+        int newWidth = integerWidth(newRoot);
+        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] Gate historical FLOAT/TIMESTAMP/TIME-to-STRING files until Rust matches 
Paimon's text casts. Java renders FLOAT 10000000f as `1.0E7`, but Arrow/ryu 
renders `10000000.0`; Java renders TIMESTAMP with a space (`2024-01-01 
00:00:00`), but Arrow/chrono inserts `T`; Java renders TIME(3) zero 
milliseconds as `00:00:00.0`, but Arrow omits `.0`. All pass this gate. Plain 
projections return different text, and Rust's exact residual casts old values 
to Utf8 before a current STRING equality, so matching JNI rows can be 
discarded. Add persisted projection and equality differentials for these source 
types.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java:
##########
@@ -411,10 +727,330 @@ private void setPaimonParams(TFileRangeDesc rangeDesc, 
PaimonSplit paimonSplit)
 
         String fileFormat = getFileFormat(paimonSplit.getPathString());
         if (split != null) {
+            // use jni reader / paimon-cpp reader / paimon-rust reader
             rangeDesc.setFormatType(TFileFormatType.FORMAT_JNI);
-            // A logical DataSplit may span multiple files, so keep it intact 
for the JNI reader.
-            fileDesc.setReaderType(TPaimonReaderType.PAIMON_JNI);
-            fileDesc.setPaimonSplit(PaimonUtil.encodeObjectToString(split));
+            // paimon-cpp and paimon-rust both consume Paimon native binary 
serialization,
+            // which only supports DataSplit. Any other split type falls back 
to JNI.
+            boolean nativeSplit = split instanceof DataSplit;
+            // Fallback-read splits stay on JNI: FallbackDataSplit extends
+            // DataSplit, so the instanceof above passes, but its serializer
+            // appends an isFallback byte after the ordinary split that the
+            // pinned rust decoder rejects outright ("trailing bytes after
+            // DataSplit" — it requires full-buffer consumption), and even a
+            // permissive decode would still lack the second table identity
+            // needed to honor the fallback-side discriminator. Both sides of a
+            // FallbackReadFileStoreTable wrap their splits, so the table
+            // wrapper is gated as a whole (any split from it routes to JNI)
+            // until the rust ABI represents both sides; the FallbackSplit
+            // interface also catches a wrapper split regardless of how the
+            // table was resolved here.
+            boolean fallbackRead = split instanceof 
FallbackReadFileStoreTable.FallbackSplit
+                    || processedTable instanceof FallbackReadFileStoreTable;
+            // Serialize the same effective table that planning and the JNI 
reader use.
+            // Relation options such as t@options('read.batch-size'='1') are 
applied by
+            // getProcessedTable() (doInitialize caches it in processedTable), 
and the
+            // rust reader derives its read batch size from the schema options 
— the raw
+            // cached table would silently drop the override. Copies, 
delegates and
+            // fallback wrappers of getProcessedTable() are still 
FileStoreTable, so the
+            // instanceof gate keeps its semantics.
+            Table paimonTable = processedTable;
+            FileStoreTable paimonFileStoreTable =
+                    paimonTable instanceof FileStoreTable ? (FileStoreTable) 
paimonTable : null;
+            // query-auth.enabled tables stay on JNI: when catalog 
authorization
+            // succeeds with no row filter or column mask, Paimon still leaves 
an
+            // ordinary DataSplit (restricted results use QueryAuthSplit and 
are
+            // already handled by the nativeSplit gate above), so this table 
shape
+            // passes the compound gate — but the shipped schema keeps
+            // query-auth.enabled=true and the pinned rust ReadBuilder rejects
+            // every such table (its CoreOptions::ensure_read_authorized fails
+            // closed because the client cannot enforce the row filter / column
+            // masking), turning a valid authorized scan into a BE-open 
failure.
+            // Until the authorization result can be transported and enforced 
by
+            // the rust ABI, these tables route to JNI.
+            boolean queryAuthTable = false;
+            // REST-token tables stay on JNI: doInitialize snapshots
+            // RESTTokenFileIO.validToken().token() into the backend storage
+            // properties, discarding expireAtMillis and the REST refresh
+            // context, so the shipped credentials look static — but the
+            // pinned rust table reuses one option map with no refresh
+            // callback, while paimon 1.4.2's JNI RESTTokenFileIO checks
+            // expiry before each file operation and obtains a replacement
+            // token. A queued or long scan that crosses the token TTL would
+            // start on rust and later fail authentication. Gate until the
+            // rust ABI can refresh and atomically update credentials.
+            boolean restTokenTable = false;
+            // Partial-update / aggregation tables with deletion vectors only 
pass
+            // the rust reader in the fully materialized shape: the pinned rust
+            // read_pk rejects merge-engine=partial-update/aggregation with
+            // deletion-vectors.merge-on-read=true outright, and otherwise 
requires
+            // every split to be compacted and known free of retract rows
+            // (DataSplit::is_fully_materialized_pk_dv). Their ordinary 
DataSplits
+            // sail through the compound gate above, so without this check a 
valid
+            // Java/JNI scan reaches BE and the rust open fails. Deduplicate 
stays
+            // rust-eligible: its read_pk routes uncompacted splits to the KV
+            // reader, which applies the attached per-file DVs. 
merge-on-read=true
+            // is a table option, so the whole table routes to JNI;
+            // non-materialized splits are gated per split below.
+            boolean puAggDeletionVectors = false;
+            boolean dvMergeOnRead = false;
+            boolean deduplicateIgnoreDelete = false;
+            boolean rustUnsupportedMergeOption = false;
+            if (paimonFileStoreTable != null) {
+                // A renewable REST token reached the shipped properties as a
+                // plain value; only the table's FileIO type reveals it 
expires.
+                // Null-safe: a table handle whose FileIO is not resolved stays
+                // rust-eligible, mirroring the CoreOptions null-safety below.
+                restTokenTable = paimonFileStoreTable.fileIO() instanceof 
RESTTokenFileIO;
+                CoreOptions resolvedCoreOptions = 
paimonFileStoreTable.coreOptions();
+                // Null-safe: a table handle whose CoreOptions is not resolved
+                // (e.g. some wrapper shapes) stays rust-eligible rather than
+                // failing the scan here — the rust open itself rejects such a
+                // table if the option is really set.
+                if (resolvedCoreOptions != null) {
+                    queryAuthTable = resolvedCoreOptions.queryAuthEnabled();
+                    CoreOptions.MergeEngine mergeEngine = 
resolvedCoreOptions.mergeEngine();
+                    if (resolvedCoreOptions.deletionVectorsEnabled()
+                            && (mergeEngine == 
CoreOptions.MergeEngine.PARTIAL_UPDATE
+                                    || mergeEngine == 
CoreOptions.MergeEngine.AGGREGATE)) {
+                        puAggDeletionVectors = true;
+                        // The merge-engine and deletion-vectors.enabled checks
+                        // above resolve through the Java CoreOptions 
accessors,
+                        // which the table builds from this same schema options
+                        // map — the one the BE rust reader deserializes from
+                        // the shipped schema JSON — so they cannot diverge 
from
+                        // what BE sees. merge-on-read has no Java accessor in
+                        // paimon 1.4, so it is read raw from the map, with the
+                        // rust parsing semantics (any case-insensitive "true"
+                        // is on, default false).
+                        TableSchema dvSchema = paimonFileStoreTable.schema();
+                        Map<String, String> dvOptions = dvSchema == null ? 
null : dvSchema.options();
+                        String mergeOnRead = dvOptions == null
+                                ? null : 
dvOptions.get(DELETION_VECTORS_MERGE_ON_READ);
+                        dvMergeOnRead = "true".equalsIgnoreCase(mergeOnRead);
+                    }
+                    // deduplicate.ignore-delete=true tables stay on JNI:
+                    // Java's DeduplicateMergeFunction skips retract records
+                    // when the option is set — including old, uncompacted
+                    // files that still contain them — but the pinned rust
+                    // read_pk does not pass table options into its
+                    // deduplicate merge: it picks the latest row and omits
+                    // the key when that row is DELETE/UPDATE_BEFORE. An
+                    // uncompacted insert followed by a delete therefore
+                    // returns the insert through JNI but silently disappears
+                    // through rust. Gate the option until the rust merge
+                    // implements it.
+                    if (mergeEngine == CoreOptions.MergeEngine.DEDUPLICATE
+                            && resolvedCoreOptions.ignoreDelete()) {
+                        deduplicateIgnoreDelete = true;
+                    }
+                    // Non-DV merge options the pinned rust read rejects: Java
+                    // supports partial-update.remove-record-on-delete /
+                    // aggregation.remove-record-on-delete and the wider
+                    // per-field retract matrix, but the rust
+                    // PartialUpdateConfig / AggregationConfig validations
+                    // return Unsupported for them — and the DV-derived gates
+                    // above only cover deletion-vector tables, so an ordinary
+                    // non-DV DataSplit with one of these options would pass 
the
+                    // compound gate and fail during the rust merge
+                    // construction. Mirror the exact rust key matrix 
(presence,
+                    // not values) against the same schema options map BE
+                    // deserializes.
+                    if (mergeEngine == CoreOptions.MergeEngine.PARTIAL_UPDATE
+                            || mergeEngine == 
CoreOptions.MergeEngine.AGGREGATE) {
+                        TableSchema mergeSchema = 
paimonFileStoreTable.schema();
+                        Map<String, String> mergeOptions =
+                                mergeSchema == null ? null : 
mergeSchema.options();
+                        rustUnsupportedMergeOption = mergeOptions != null
+                                && hasRustUnsupportedMergeOption(mergeOptions, 
mergeEngine);
+                    }
+                }
+            }
+            // paimon-rust additionally requires (a) FileScannerV2: the V1 
FileScanner
+            // explicitly rejects PAIMON_RUST, so with enable_file_scanner_v2 
disabled
+            // the split falls back to JNI instead of encoding a rust request 
that the
+            // selected scanner cannot consume, and (b) a FileStoreTable: BE 
opens the
+            // table via paimon_table_from_schema_json, which needs the 
resolved
+            // TableSchema that only FileStoreTable exposes via schema(). If 
the table
+            // is not a FileStoreTable (e.g. a sys table backed by DataSplit), 
we cannot
+            // ship a schema JSON, so fall back to CPP / JNI rather than 
sending an
+            // incomplete PAIMON_RUST request that BE would reject.
+            //
+            // The paimon-rust S3 bridge maps static credentials, anonymous
+            // access (AWS_CREDENTIALS_PROVIDER_TYPE=ANONYMOUS -> s3.anonymous)
+            // and assume-role (AWS_ROLE_ARN / AWS_EXTERNAL_ID ->
+            // s3.assumed.role.*), but the remaining credential-provider modes
+            // are ambient JVM provider chains (ENV, SYSTEM_PROPERTIES,
+            // WEB_IDENTITY, CONTAINER, INSTANCE_PROFILE) with no paimon-rust
+            // equivalent — rust would silently sign with whatever the ambient
+            // chain resolves to. Gate those modes away from the rust reader
+            // here so the configured provider is honored via the JNI path.
+            boolean providerModeTranslatable = true;
+            String providerType = backendStorageProperties == null
+                    ? null : 
backendStorageProperties.get("AWS_CREDENTIALS_PROVIDER_TYPE");
+            String mode = providerType == null ? "DEFAULT" : 
providerType.trim().toUpperCase(Locale.ROOT);
+            if (mode.isEmpty()) {
+                mode = "DEFAULT";
+            }
+            String location = source.getTableLocation();
+            if (location != null && (location.startsWith("s3://") || 
location.startsWith("s3a://"))) {
+                // Java DEFAULT may resolve JVM properties or anonymous 
credentials;
+                // Rust's ambient chain is different. Only explicit keys or 
explicit
+                // anonymous access can cross this boundary without changing 
identity.
+                String accessKey = backendStorageProperties == null
+                        ? null : 
backendStorageProperties.get("AWS_ACCESS_KEY");
+                String secretKey = backendStorageProperties == null
+                        ? null : 
backendStorageProperties.get("AWS_SECRET_KEY");
+                boolean staticKeys = accessKey != null && 
!accessKey.trim().isEmpty()
+                        && secretKey != null && !secretKey.trim().isEmpty();
+                providerModeTranslatable = mode.equals("ANONYMOUS") || 
(mode.equals("DEFAULT") && staticKeys);
+            } else if (providerType != null) {
+                providerModeTranslatable = mode.equals("DEFAULT") || 
mode.equals("ANONYMOUS");
+                // OSS has no anonymous FileIO mode in the pinned Rust 
dependency.
+                if (mode.equals("ANONYMOUS") && location != null && 
location.startsWith("oss://")) {
+                    providerModeTranslatable = false;
+                }
+            }
+            // Incremental scans (binlog / changelog / delta / diff) must stay
+            // on the JNI path: this wire format carries only an ordinary
+            // DataSplit and the rust reader invokes TableRead::to_arrow, but
+            // paimon 1.4 marks incremental splits as streaming (which the
+            // pinned rust deserializer rejects), diff requires a separate
+            // IncrementalPlan instead of an ordinary plan, and ordinary
+            // primary-key reads can merge versions rather than return the
+            // changes — until the C ABI transports the mode and plan, the
+            // rust reader cannot express any of these.
+            TableScanParams incrementalParams = getScanParams();
+            boolean isIncremental = incrementalParams != null && 
incrementalParams.incrementalRead();
+            // ORC TIMESTAMP_WITH_LOCAL_TIME_ZONE schemas stay on JNI: the 
pinned
+            // paimon-rust ORC decoder materializes LTZ instants shifted by the
+            // writer timezone (an upstream crate limitation), so a logical ORC
+            // DataSplit that selects rust (e.g. with force_jni_scanner=true or
+            // when raw conversion is unavailable) returns a different instant
+            // than JNI — applying the session timezone in BE cannot repair an
+            // epoch already shifted during decode. Two bypasses are covered:
+            // (a) the format must come from EVERY member file — paimon allows
+            // per-level file.format, so one DataSplit can mix Parquet and ORC
+            // files and the split path's suffix (the first file) would hide
+            // the ORC members; (b) the LTZ search must recurse into nested
+            // types — an LTZ under MAP/ARRAY/ROW reaches the same shifted ORC
+            // decode through the container's field materialization. Parquet
+            // files with any LTZ, and ORC without any recursive LTZ, stay
+            // rust-eligible. nativeSplit only guards the cast — non-DataSplit
+            // splits already route to JNI.
+            boolean orcLtzSchema = paimonFileStoreTable != null
+                    && nativeSplit
+                    && splitHasOrcFile((DataSplit) split)
+                    && paimonFileStoreTable.schema().fields().stream()
+                            .anyMatch(field -> 
containsTimestampLtz(field.type()));
+            // data-file.external-paths splits stay on JNI (see
+            // splitHasExternalFiles): the rust table's single FileIO cannot
+            // serve an external file's backend. nativeSplit only guards the
+            // cast — non-DataSplit splits already route to JNI.
+            boolean externalFileSplit = nativeSplit && 
splitHasExternalFiles((DataSplit) split);
+            // Projected VARIANT columns stay on JNI: the rust leaf feeds its
+            // Arrow arrays to the slot serdes, and DataTypeVariantV2SerDe::
+            // read_column_from_arrow unconditionally returns
+            // NOT_IMPLEMENTED_ERROR — a nested Variant (ARRAY / MAP / STRUCT
+            // containing one) reaches the same decoder through the container
+            // serdes. desc carries only the slots this query projects, so a
+            // table whose VARIANT column is not projected still scans on
+            // rust. Gate until the rust leaf has a Variant Arrow decoder.
+            boolean projectedVariant = desc.getSlots().stream()
+                    .anyMatch(slot -> 
PaimonUtil.containsVariant(slot.getType()));
+            // Scheme capability gate: the pinned paimon-rust storage
+            // dispatcher (io/storage.rs) selects the FileIO parser from the
+            // table location's URI scheme, and libpaimon_c.a compiles in
+            // separate COS, OBS, GCS and Azdls parsers besides the OSS and S3
+            // ones. Doris normalizes every object store's credentials into
+            // the AWS_* / use_path_style aliases (see the *Properties storage
+            // classes), which the BE rust bridge translates only into the
+            // fs.oss.* and s3.* key families — a cosn:// / obs:// / gs:// /
+            // abfs:// warehouse would reach its scheme's parser without the
+            // key family it reads (fs.cosn.userinfo.*, fs.obs.*, gcs.*,
+            // azure.*) and fail the open instead of using JNI. Only the
+            // schemes whose property translation is implemented and
+            // open-tested (s3 / s3a / oss, via RUST_VERIFIED_LOCATION_SCHEMES)
+            // plus the credential-free hdfs and local-filesystem parsers stay
+            // rust-eligible; every other scheme falls back to JNI. A null
+            // location also routes to JNI: the rust path needs the
+            // paimon_table that only a real location can provide (BE rejects
+            // a split without it).
+            boolean schemeCapabilityVerified = 
isRustVerifiedLocationScheme(source.getTableLocation());
+            // An hdfs:// location is scheme-verified only together with the 
credential-free
+            // backend shape: the backend storage properties that ship to BE 
also carry an
+            // HDFS catalog's authentication (kerberos principal / keytab, 
proxy user, HA
+            // nameservice config), none of which the pinned rust HDFS parser 
reads — the
+            // scan would open as the BE process's ambient identity instead of 
the
+            // catalog's configured one and fail the access JNI honors. See
+            // isRustVerifiedHdfsBackend.
+            boolean hdfsBackendVerified = 
!isHdfsLocationScheme(source.getTableLocation())
+                    || isRustVerifiedHdfsBackend(backendStorageProperties);
+            // With merge-on-read=true the whole table already routes to JNI 
(dvMergeOnRead);
+            // for the remaining partial-update/aggregation DV tables, a split 
that is
+            // not fully materialized (uncompacted level-0 data, or 
retractions not
+            // known to be applied — even a split with no deletion file 
attached yet)
+            // fails the rust is_fully_materialized_pk_dv guard, so it falls 
back per
+            // split instead of turning into a BE-open failure. nativeSplit and
+            // !fallbackRead only guard the cast — those splits already route 
to JNI.
+            boolean splitDvNotMaterialized = puAggDeletionVectors && 
!dvMergeOnRead
+                    && nativeSplit && !fallbackRead
+                    && !isFullyMaterializedPkDvSplit((DataSplit) split);
+            boolean canUseRust = sessionVariable.isEnablePaimonRustReader()
+                    && sessionVariable.enableFileScannerV2 && nativeSplit && 
!fallbackRead
+                    && !isIncremental && providerModeTranslatable && 
!queryAuthTable
+                    && !restTokenTable
+                    && !dvMergeOnRead && !splitDvNotMaterialized
+                    && !orcLtzSchema && !projectedVariant && !externalFileSplit
+                    && !deduplicateIgnoreDelete && !rustUnsupportedMergeOption
+                    && schemeCapabilityVerified && hdfsBackendVerified
+                    && paimonFileStoreTable != null;
+            if (canUseRust) {
+                if (rustReaderCapabilities == null) {
+                    rustReaderCapabilities = new 
PaimonRustReaderCapabilities(paimonFileStoreTable, desc);
+                }
+                canUseRust = rustReaderCapabilities.canRead((DataSplit) split);
+            }
+            if (canUseRust) {
+                fileDesc.setReaderType(TPaimonReaderType.PAIMON_RUST);

Review Comment:
   [P2] Fence this opt-in reader type during mixed-BE upgrades. The FE selects 
`PAIMON_RUST=3` from session/table eligibility without checking which BE 
receives the range. A BE at the PR base accepts explicit Paimon ranges only for 
`PAIMON_JNI` in FileScannerV2 and fails validation on value 3; its V1 JNI 
reader also cannot consume the native-binary DataSplit sent here. Thus enabling 
the session switch while even one older BE remains available makes otherwise 
valid queries fail depending on scheduling. Route these ranges only to capable 
BEs or hold Rust emission until all eligible BEs support it, with a 
mixed-version routing check.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,257 @@
+// 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());
+        DataTypeRoot oldRoot = oldType.getTypeRoot();
+        DataTypeRoot newRoot = newType.getTypeRoot();
+        boolean targetTimestamp = newRoot == 
DataTypeRoot.TIMESTAMP_WITHOUT_TIME_ZONE
+                || newRoot == DataTypeRoot.TIMESTAMP_WITH_LOCAL_TIME_ZONE;
+        // Java accepts epoch-day/millisecond strings and interprets numeric 
timestamps as seconds.
+        // Arrow instead parses calendar strings or reuses numeric timestamp 
ticks, including in
+        // nested values. Keep these historical casts on JNI until the pinned 
reader matches Java.
+        boolean sourceString = oldRoot == DataTypeRoot.CHAR || oldRoot == 
DataTypeRoot.VARCHAR;

Review Comment:
   [P1] Fall back to JNI for historical TIME-to-TIMESTAMP files. This gate 
admits TIME_WITHOUT_TIME_ZONE as a timestamp source, and Paimon 1.4.2 allows 
the explicit type change by default. Java converts the old millisecond-of-day 
value to a Timestamp; pinned Rust maps it to Arrow Time32(Millisecond), then 
asks arrow-cast to convert to Timestamp, which has no such cast and fails the 
scan. Add a persisted TIME-to-TIMESTAMP root-column differential.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,257 @@
+// 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());
+        DataTypeRoot oldRoot = oldType.getTypeRoot();
+        DataTypeRoot newRoot = newType.getTypeRoot();
+        boolean targetTimestamp = newRoot == 
DataTypeRoot.TIMESTAMP_WITHOUT_TIME_ZONE
+                || newRoot == DataTypeRoot.TIMESTAMP_WITH_LOCAL_TIME_ZONE;
+        // Java accepts epoch-day/millisecond strings and interprets numeric 
timestamps as seconds.
+        // Arrow instead parses calendar strings or reuses numeric timestamp 
ticks, including in
+        // nested values. Keep these historical casts on JNI until the pinned 
reader matches Java.
+        boolean sourceString = oldRoot == DataTypeRoot.CHAR || oldRoot == 
DataTypeRoot.VARCHAR;
+        if ((sourceString && (newRoot == DataTypeRoot.DATE || targetTimestamp))
+                || (oldWidth > 0 && targetTimestamp)) {
+            return true;
+        }
+        int newWidth = integerWidth(newRoot);
+        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) {

Review Comment:
   [P1] Keep historical ROW-to-STRING files on JNI until Rust implements 
Paimon's cast. Paimon 1.4.2 permits this explicit type change and Java 
RowToStringCastRule returns text such as `{1, x}`. This check recurses only 
when both versions are ROW, so an old ROW/new STRING reaches the compatible 
fallthrough; the current scalar slot also passes the projection gate. The 
pinned Rust reader then asks Arrow to cast the old Struct array to Utf8, which 
Arrow explicitly rejects, failing an otherwise valid scan. Add a persisted root 
ROW-to-STRING differential.



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