Gabriel39 commented on code in PR #66227:
URL: https://github.com/apache/doris/pull/66227#discussion_r4128779658


##########
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:
   Fixed by checking children().size() before using the 16-bit expression 
count. Oversized IN sets remain residual. Added BE coverage at 
65,534/65,535/65,536/65,537 values and JNI/Rust SQL regressions at 
65,535/65,536/65,537 keys with complete COUNT/MAX checks and runtime-filter 
profile assertions. The new BE test reproduced the incorrect partial pushdown 
before the fix and now passes.



##########
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:
   Fixed by including split open in the PaimonRustReader parent timer. Added a 
failed-open test that never calls get_block and asserts the parent includes 
OpenSplitTime. It reproduced the zero parent timer before the fix and now 
passes.



##########
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:
   Fixed with a recursive allowlist of established schema-cast equivalence. 
Historical TIMESTAMP to DATE now uses JNI, including nested fields and map 
keys. Added a real persisted pre-epoch fixture proving the Java truncation 
result, recursive routing tests, and safe-cast controls.



##########
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:
   Fixed by keeping NTZ/LTZ schema transitions on JNI in either direction. 
Added persisted Parquet fixtures in Asia/Shanghai that verify the Java timezone 
conversion, plus recursive routing coverage and same-kind timestamp widening 
controls.



##########
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:
   Fixed: historical TIME to TIMESTAMP now stays on JNI. Added a persisted TIME 
fixture verifying the Java epoch-millisecond result and recursive routing tests 
for both NTZ and LTZ targets.



##########
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:
   Confirmed against released Rust producers: v0.2.0 partial-update and v0.3.0 
aggregation can retain repeated INSERT keys in one file with zero delete count. 
These primary-key merge engines now always use JNI because file metadata cannot 
establish producer normalization. Added one-file/zero-delete negative tests and 
retained deduplicate/first-row positive controls; SQL regressions verify the 
same routing.



##########
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:
   Fixed by keeping floating-point and temporal conversions to text on JNI, 
including nested fields. Persisted fixtures verify scientific FLOAT/DOUBLE 
formatting, TIMESTAMP separators/fractions, TIME formatting, and NULLs. 
Recursive routing tests cover STRING, VARCHAR, and CHAR targets.



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