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


##########
gensrc/thrift/PlanNodes.thrift:
##########
@@ -398,6 +399,11 @@ struct TPaimonFileDesc {
     // Original Paimon RawFile.path() before Doris storage path normalization. 
Native readers use this
     // to materialize the public file-location metadata column.
     18: optional string original_file_path;
+    // paimon-rust: TableSchema serialized via 
JsonSerdeUtil.toJson(table.schema()).
+    // Consumed by paimon_table_from_schema_json on BE; skips 
catalog/warehouse resolution.
+    19: optional string paimon_table_schema_json;

Review Comment:
   [P1] Keep Paimon Rust payload IDs compatible across release branches. 
Branches 4.1 and 4.2 use field 18 for schema JSON and field 19 for branch, 
whereas this new field 19 is schema JSON and field 20 is branch. Those branches 
already advertise Rust capability, so a mixed-version FE/BE pair can select 
Rust and decode the schema as a branch, or find no schema at all, causing scans 
to fail. Master already used field 18 for `original_file_path` before this PR, 
so a compatible versioned payload/capability or equivalent migration is needed; 
merely renumbering that existing field would break master-side compatibility.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -991,16 +1009,24 @@ && hasTrustworthyLimitAccounting(table)
                             dataSplit.bucket()));
                 }
             } else {
-                // JNI reader path
+                // Logical reader path: use paimon-rust only for the 
compatibility-checked subset.
                 if (ignoreJni) {
                     // FIX-L14: ignore_split_type=IGNORE_JNI drops JNI splits 
(legacy getSplits:483).
                     continue;
                 }
                 if (requiresMetadataColumns) {
                     validateMetadataColumnReader(true, false);
                 }
-                ranges.add(buildJniScanRange(dataSplit, defaultFileFormat,
-                        partitionValues, true, weightDenominator));
+                if (rustReaderSelector != null && 
rustReaderSelector.canRead(dataSplit)) {

Review Comment:
   [P2] Preserve the `force_jni_scanner` escape hatch when Rust is enabled. 
That session variable is documented as forcing JNI for external tables. It 
disables the native arm, but this logical arm still selects Rust when 
`enable_paimon_rust_reader=true` and the split is eligible; the added 
differential test uses exactly that combination. A user setting the force flag 
to bypass a failing Rust read will still run Rust. Make the override reach the 
Rust selector, and use a separate logical-split test mechanism if needed.



##########
be/src/format_v2/table/paimon_rust_table_reader.cpp:
##########
@@ -0,0 +1,924 @@
+// 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);
+        // Open may dominate the scan or fail before get_block; include it in 
the parent timer.
+        SCOPED_TIMER(_rust_total_time);
+        SCOPED_TIMER(_rust_open_split_time);
+        RETURN_IF_ERROR(_open_split_reader(options.current_range));
+    }
+    return Status::OK();
+}
+
+Status PaimonRustTableReader::get_block(Block* block, bool* eos) {
+    SCOPED_TIMER(_profile.total_timer);
+    SCOPED_TIMER(_profile.exec_timer);
+    SCOPED_TIMER(_rust_total_time);
+    DORIS_CHECK(block != nullptr);
+    DORIS_CHECK(eos != nullptr);
+    DORIS_CHECK(block->columns() == _projected_columns.size());
+    block->clear_column_data(_projected_columns.size());
+    *eos = false;
+
+    if (_is_table_level_count_active()) {
+        return _read_table_level_count(block, eos);
+    }
+
+    // num_splits == 0 yields an empty (but valid) stream: report EOF.
+    if (_split_eof) {
+        *eos = true;
+        return Status::OK();
+    }
+    if (!_handles || !_handles->reader) {
+        return Status::InternalError("paimon-rust reader is not initialized");
+    }
+
+    while (true) {
+        // Mirror the base TableReader cancellation contract so a cancelled 
query does not
+        // drain the whole split.
+        if (_io_ctx != nullptr && _io_ctx->should_stop) {
+            _split_eof = true;
+            _close_split_reader();
+            *eos = true;
+            return Status::OK();
+        }
+
+        paimon_result_next_batch next;
+        {
+            SCOPED_TIMER(_rust_read_batch_time);
+            next = paimon_record_batch_reader_next(_handles->reader.get());
+        }
+        if (next.error != nullptr) {
+            return Status::InternalError("paimon-rust read batch failed: {}",
+                                         consume_error(next.error));
+        }
+        // End of stream: both pointers are null.
+        if (next.batch.array == nullptr && next.batch.schema == nullptr) {
+            _split_eof = true;
+            _close_split_reader();
+            *eos = true;
+            return Status::OK();
+        }
+
+        // RAII: the batch's Arrow release callbacks + container free run when
+        // `batch` leaves this scope, including on any early return.
+        ArrowBatch batch(next.batch);
+
+        auto* c_array = batch.array();
+        auto* c_schema = batch.schema();
+        arrow::Result<std::shared_ptr<arrow::RecordBatch>> import_result =
+                arrow::ImportRecordBatch(c_array, c_schema);
+        if (!import_result.ok()) {
+            return Status::InternalError("failed to import paimon-rust arrow 
batch: {}",
+                                         import_result.status().message());
+        }
+
+        auto record_batch = std::move(import_result).ValueUnsafe();
+        const auto rows = static_cast<size_t>(record_batch->num_rows());
+        if (rows == 0) {
+            // Skip empty batches and keep draining the stream.
+            continue;
+        }
+        RETURN_IF_ERROR(_fill_block_from_record_batch(record_batch, block, 
rows));
+        _record_scan_rows(rows);
+        *eos = false;
+        return Status::OK();
+    }
+}
+
+Status PaimonRustTableReader::abort_split() {
+    {
+        SCOPED_TIMER(_profile.total_timer);
+        SCOPED_TIMER(_profile.close_timer);
+        _close_split_reader();
+        _split_eof = false;
+    }
+    return format::TableReader::abort_split();
+}
+
+#ifdef BE_TEST
+std::string PaimonRustTableReader::TEST_format_options(
+        const std::map<std::string, std::string>& options) {
+    return format_options(options);
+}
+
+std::map<std::string, std::string> PaimonRustTableReader::TEST_build_options(
+        TFileScanRangeParams* scan_params, const TFileRangeDesc& range) {
+    TFileScanRangeParams* previous_params = _scan_params;
+    TFileRangeDesc previous_range = _current_range;
+    _scan_params = scan_params;
+    _current_range = range;
+    std::map<std::string, std::string> options = _build_options();
+    _scan_params = previous_params;
+    _current_range = std::move(previous_range);
+    return options;
+}
+#endif
+
+Status PaimonRustTableReader::close() {
+    {
+        SCOPED_TIMER(_profile.total_timer);
+        SCOPED_TIMER(_profile.close_timer);
+        _close_split_reader();
+        _close_table();
+    }
+    return format::TableReader::close();
+}
+
+Status PaimonRustTableReader::_validate_rust_split(const TFileRangeDesc& 
range) const {
+    if (!range.__isset.table_format_params || 
!range.table_format_params.__isset.paimon_params) {
+        return Status::InternalError(
+                "missing paimon_params for paimon rust reader, possibly caused 
by FE/BE protocol "
+                "mismatch");
+    }
+    const auto& params = range.table_format_params.paimon_params;
+    if (!params.__isset.paimon_split || params.paimon_split.empty()) {
+        return Status::InternalError(
+                "missing paimon_split for paimon rust reader, possibly caused 
by FE/BE protocol "
+                "mismatch");
+    }
+    if (params.__isset.reader_type && params.reader_type != 
TPaimonReaderType::PAIMON_RUST) {
+        return Status::InternalError(
+                "invalid reader_type for paimon rust reader, possibly caused 
by FE/BE protocol "
+                "mismatch");
+    }
+    if (!_resolve_table_path(range).has_value()) {
+        return Status::InternalError(
+                "paimon-rust missing paimon_table; cannot resolve paimon table 
location");
+    }
+    if (!_resolve_db_name(range).has_value()) {
+        return Status::InternalError(
+                "paimon-rust missing db_name; cannot open paimon table via 
schema json");
+    }
+    if (!_resolve_table_name(range).has_value()) {
+        return Status::InternalError(
+                "paimon-rust missing table_name; cannot open paimon table via 
schema json");
+    }
+    if (!_resolve_table_schema_json(range).has_value()) {
+        return Status::InternalError(
+                "paimon-rust missing paimon_table_schema_json; cannot open 
paimon table via "
+                "schema json");
+    }
+    return Status::OK();
+}
+
+Status PaimonRustTableReader::_open_split_reader(const TFileRangeDesc& range) {
+    // 1. Decode the FE-planned split first so we fail fast (and without any
+    // filesystem IO) when it is missing or malformed.
+    std::string split_bytes;
+    RETURN_IF_ERROR(_decode_split_bytes(&split_bytes));
+
+    // 2. Resolve identifier + table_path + FE-supplied TableSchema JSON.
+    auto table_path = _resolve_table_path(range).value();
+    auto db_name = _resolve_db_name(range).value();
+    auto table_name = _resolve_table_name(range).value();
+    auto schema_json = _resolve_table_schema_json(range).value();
+    auto branch_opt = _resolve_branch(range);
+
+    // 3. Assemble storage options: FE-supplied paimon options + hadoop_conf +
+    // OSS/S3 → AWS_* translations. These feed FileIO only (per
+    // paimon_table_from_schema_json contract); they are NOT merged into the
+    // supplied table schema.
+    auto options = _build_options();
+
+    auto opened_table_key =
+            std::make_tuple(table_path, schema_json, db_name, table_name, 
branch_opt, options);
+    if (!_handles || !_handles->table || _opened_table_key != 
opened_table_key) {
+        // A paimon scan reads one table, so the handle is opened at most once 
per
+        // distinct identity (e.g. re-created after a close); splits of the 
same
+        // table reuse it and only rebuild the read pipeline below.
+        _close_table();
+        _handles = std::make_unique<PaimonHandles>();
+
+        std::vector<paimon_option> c_options;
+        c_options.reserve(options.size());
+        for (const auto& kv : options) {
+            c_options.push_back(paimon_option {kv.first.c_str(), 
kv.second.c_str()});
+        }
+
+        LOG(INFO) << "paimon-rust opening table via schema json: db=" << 
db_name
+                  << " table=" << table_name << " path=" << table_path
+                  << " branch=" << (branch_opt.has_value() ? 
branch_opt.value() : "main")
+                  << " storage_options=[" << format_options(options) << "]";
+
+        // Build the table directly from the FE-supplied schema JSON. The Rust
+        // side rejects null / empty branch, so we default to paimon's 
canonical
+        // "main" sentinel when FE did not set paimon_branch (i.e. the table is
+        // on the main branch — matches upstream 
Identifier.DEFAULT_MAIN_BRANCH).
+        const std::string& branch_str = branch_opt.has_value() ? 
branch_opt.value() : "main";
+        paimon_result_get_table tbl_res = paimon_table_from_schema_json(
+                table_path.c_str(), schema_json.c_str(), db_name.c_str(), 
table_name.c_str(),
+                branch_str.c_str(), c_options.empty() ? nullptr : 
c_options.data(),
+                c_options.size());
+        if (tbl_res.error != nullptr) {
+            return Status::InternalError(
+                    "paimon-rust table_from_schema_json failed: db={} table={} 
err={}", db_name,
+                    table_name, consume_error(tbl_res.error));
+        }
+        _handles->table.reset(tbl_res.table);
+        _opened_table_key = std::move(opened_table_key);
+    }
+
+    // 4. Build the read pipeline: read_builder -> case-insensitive -> 
projection.
+    // Bound the Rust Arrow allocation itself: wide rows must respect the 
scanner's
+    // adaptive probe size before they are materialized into a Doris block.
+    const auto batch_size = std::to_string(
+            _batch_size > 0 ? _batch_size : std::max(1, 
_runtime_state->batch_size()));
+    const paimon_option batch_option {"read.batch-size", batch_size.c_str()};

Review Comment:
   [P2] Apply adaptive row caps to the active Rust stream. FileScannerV2 seeds 
32 rows before `prepare_split`, then learns bytes per row and calls 
`set_batch_size` before later `get_block` calls. Here `read.batch-size` is 
fixed when the Rust builder opens, while the inherited setter only changes 
`_batch_size`. A split with 1 MiB rows therefore keeps emitting roughly 32 MiB 
Arrow batches after the scanner asks for about eight rows to meet its default 8 
MiB target. The added test changes the cap only between splits; please cover 
and honor a cap change after the first batch.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonRustReaderCapabilities.java:
##########
@@ -0,0 +1,281 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.connector.paimon;
+
+import org.apache.doris.connector.spi.handle.ConnectorColumnHandle;
+
+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.DataTypeChecks;
+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, 
List<ConnectorColumnHandle> columns) {
+        this.table = table;
+        this.schema = table.schema();
+        this.tableCompatible = schema != null && 
hasCompatibleSchemaVersion(schema)
+                && hasFullNestedProjection(columns) && 
hasCompatibleMergeEngine(schema)
+                && 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(

Review Comment:
   [P1] Guard current-schema nanosecond timestamps before Rust selection. 
Paimon permits `TIMESTAMP(9)` year-3000 values and stores precision >6 as 
Parquet INT96 ([type 
range](https://paimon.apache.org/docs/master/concepts/data-types/); [Parquet 
mapping](https://paimon.apache.org/docs/master/concepts/spec/fileformat/)). The 
pinned Rust/Parquet reader maps INT96 to signed i64 Arrow nanoseconds, and 
[`Int96::to_nanos` wraps on 
overflow](https://github.com/apache/arrow-rs/blob/58.4.0/parquet/src/data_type.rs#L110-L119)
 beyond year 2262. This `schemaId() != schema.id()` condition skips the only 
timestamp precision check for current-schema files, so an eligible logical 
split can return a different date from JNI without error. Gate precision 7–9, 
including nested leaves, until decoding preserves the full range; add a 
year-3000 parity case.



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