github-actions[bot] commented on code in PR #68825:
URL: https://github.com/apache/doris/pull/68825#discussion_r4230741120
##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -575,7 +581,8 @@ public List<ConnectorScanRange> planScan(ConnectorSession
session, ConnectorScan
return scanReuse.computeIfAbsent(reuseKey,
key -> Collections.unmodifiableList(planScanInternal(session,
request.getTableHandle(), request.getColumns(),
request.getFilter(),
Review Comment:
[P1] Include Rust eligibility in the scan reuse key. The key records only
root column names, so with nested pruning enabled two aliases reading the same
struct root, one full and one `s.b`, collide even though only the full
projection may use Rust. If it plans first, this cache hands its Rust ranges to
the nested scan, bypassing the selector's required JNI fallback. The common
backend capability set is also absent, so independently selected aliases in a
mixed-version BE pool can reuse ranges assigned to an incapable backend. Key
the projection shape and backend capabilities before reusing these ranges.
##########
fe/fe-connector/fe-connector-spi/src/main/java/org/apache/doris/connector/spi/scan/ConnectorScanRequest.java:
##########
@@ -147,10 +151,20 @@ public boolean isExplainOnly() {
return explainOnly;
}
+ /** Capabilities shared by every backend eligible to execute this scan. */
+ public Set<String> getBackendCapabilities() {
+ return backendCapabilities;
Review Comment:
[P1] Bump the connector SPI major for the new request method. The Paimon
plugin now invokes `ConnectorScanRequest.allBackendsSupport`, but the plugin
API stays at 12.0. An updated directory plugin can pass the old FE's
major-version gate and fail its first scan with `NoSuchMethodError`. Increment
the API major and update the pinned surface/version baselines so that older FEs
reject the incompatible plugin at load time.
##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonScanPlanProvider.java:
##########
@@ -1650,6 +1668,25 @@ private PaimonScanRange buildJniScanRange(Split split,
String defaultFileFormat,
return builder.build();
}
+ private PaimonScanRange buildRustScanRange(DataSplit split, FileStoreTable
table,
+ PaimonTableHandle handle, String defaultFileFormat,
Review Comment:
[P2] Serialize the table schema once per scan. `buildRustScanRange` runs for
every logical split and calls `table.schema()`, strips options, and encodes a
fresh JSON string each time; every range retains that identical string. A
split-heavy table can consume hundreds of MB or more of FE heap and repeat
substantial JSON work during planning. Compute the schema JSON and branch once
outside the split loop and share them among range builders.
##########
be/src/format_v2/table/paimon_rust_predicate_converter.cpp:
##########
@@ -0,0 +1,840 @@
+// 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 <limits>
+#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());
+ // Runtime-generated IN sets can exceed VExpr's 16-bit child count. A
wrapped
+ // count would push only a prefix and permanently discard valid probe rows.
+ if (!in_pred || expr->children().size() < 2 ||
+ expr->children().size() > std::numeric_limits<uint16_t>::max()) {
+ 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) {
+ // Casted list values are rejected by _convert_literal (the same rule
+ // as the binary RHS): with debug_skip_fold_constant the cast reaches
+ // the BE un-folded, and unwrapping it would filter on the pre-cast
+ // value — in `amount IN (CAST(1.24 AS DECIMAL(10,1)))` Doris keeps
+ // the 1.2 rows while the unwrapped 1.24 push removes them. Rejecting
+ // the value rejects the whole predicate; the residual applies the
+ // cast correctly.
+ auto holder = _convert_literal(expr->get_child(i), field_meta->type);
+ if (!holder) {
+ return nullptr;
+ }
+ storages.emplace_back(std::move(holder->storage));
+ paimon_datum datum = holder->datum;
+ _bind_datum_storage(&datum, storages.back());
+ datums.emplace_back(datum);
+ }
+
+ if (datums.empty()) {
+ return nullptr;
+ }
+ if (in_pred->is_not_in()) {
+ return _take(paimon_predicate_is_not_in(_table,
field_meta->column.c_str(), datums.data(),
+ datums.size()));
+ }
+ return _take(paimon_predicate_is_in(_table, field_meta->column.c_str(),
datums.data(),
+ datums.size()));
+}
+
+paimon_predicate* PaimonRustPredicateConverter::_convert_binary(const
VExprSPtr& expr) {
+ if (!expr || expr->get_num_children() != 2) {
+ return nullptr;
+ }
+ auto field_meta = _resolve_field(expr->get_child(0));
+ if (!field_meta) {
+ return nullptr;
+ }
+ const char* column = field_meta->column.c_str();
+
+ // Convert the RHS first so EQ_FOR_NULL (<=>) only converts when the RHS is
+ // a convertible literal, mirroring the FE converter, which rejects a
+ // non-literal RHS. A column-to-column `a <=> b` must therefore stay in the
+ // Doris residual: it has no single-column rust predicate, and pushing
+ // `a IS NULL` would wrongly discard rows like (1, 1) — rows dropped by the
+ // pushed filter cannot be recovered by the residual conjunct.
+ auto holder = _convert_literal(expr->get_child(1), field_meta->type);
+ if (!holder) {
+ return nullptr;
+ }
+
+ if (expr->op() == TExprOpcode::EQ_FOR_NULL) {
+ return _take(paimon_predicate_is_null(_table, column));
Review Comment:
[P1] Preserve non-null null-safe equality when pushing this predicate. With
`disable_nereids_expression_rules='NULL_SAFE_EQUAL_TO_EQUAL'`, `WHERE id <=> 1`
reaches this converter with a non-null literal. It currently pushes `id IS
NULL`, so Rust removes every matching `id=1` row before Doris applies the
residual and the query returns no rows. Push equality for a non-null literal,
reserving `IS NULL` for a null literal, or leave the expression as a residual.
##########
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()};
+ paimon_result_read_builder rb_res =
+ paimon_table_new_read_builder_with_options(_handles->table.get(),
&batch_option, 1);
+ if (rb_res.error != nullptr) {
+ return Status::InternalError("paimon-rust new read builder failed: {}",
+ consume_error(rb_res.error));
+ }
+ _handles->read_builder.reset(rb_res.read_builder);
+
+ // Fold column casing on the Rust side so FE-normalized lowercase names
+ // resolve against tables with mixed-case column definitions.
+ if (paimon_error* case_err =
+
paimon_read_builder_with_case_sensitive(_handles->read_builder.get(), false)) {
+ return Status::InternalError("paimon-rust set case_sensitive failed:
{}",
+ consume_error(case_err));
+ }
+
+ // Partition keys are excluded: they are materialized from split metadata
+ // (see _fill_non_arrow_columns), and paimon-rust does not emit them.
+ auto read_columns = _build_read_columns();
+ std::vector<const char*> projection;
+ projection.reserve(read_columns.size() + 1);
+ for (const auto& col : read_columns) {
+ projection.push_back(col.c_str());
+ }
+ projection.push_back(nullptr);
+ if (paimon_error* proj_err =
paimon_read_builder_with_projection(_handles->read_builder.get(),
+
projection.data())) {
+ return Status::InternalError("paimon-rust set projection failed: {}",
+ consume_error(proj_err));
+ }
+
+ // Convert the scanner conjuncts into a paimon-rust filter and apply it.
+ RETURN_IF_ERROR(_apply_predicate());
+
+ // 5. Deserialize the FE-planned split into a one-split plan, so this
+ // scanner reads exactly the split it was assigned rather than replanning
+ // the whole table. The wire form is identical to what paimon-cpp consumes
+ // (`paimon::table::DataSplit::serialize`).
+ paimon_result_plan plan_res = paimon_plan_from_split_bytes(
+ reinterpret_cast<const uint8_t*>(split_bytes.data()),
split_bytes.size());
+ if (plan_res.error != nullptr) {
+ return Status::InternalError("paimon-rust build plan failed: {}",
+ consume_error(plan_res.error));
+ }
+ _handles->plan.reset(plan_res.plan);
+
+ size_t num_splits = paimon_plan_num_splits(_handles->plan.get());
+ if (num_splits == 0) {
+ _split_eof = true;
+ return Status::OK();
+ }
+
+ // 6. Open the arrow stream over the plan.
+ paimon_result_new_read read_res =
paimon_read_builder_new_read(_handles->read_builder.get());
+ if (read_res.error != nullptr) {
+ return Status::InternalError("paimon-rust new read failed: {}",
+ consume_error(read_res.error));
+ }
+ _handles->table_read.reset(read_res.read);
+
+ paimon_result_record_batch_reader rdr_res = paimon_table_read_to_arrow(
+ _handles->table_read.get(), _handles->plan.get(), /*offset=*/0,
/*length=*/num_splits);
+ if (rdr_res.error != nullptr) {
+ return Status::InternalError("paimon-rust open arrow reader failed:
{}",
+ consume_error(rdr_res.error));
+ }
+ _handles->reader.reset(rdr_res.reader);
+ return Status::OK();
+}
+
+void PaimonRustTableReader::_close_split_reader() {
+ if (!_handles) {
+ return;
+ }
+ // Reverse of the declaration order in PaimonHandles.
+ _handles->reader.reset();
+ _handles->table_read.reset();
+ _handles->plan.reset();
+ _handles->read_builder.reset();
+}
+
+void PaimonRustTableReader::_close_table() {
+ if (!_handles) {
+ return;
+ }
+ _close_split_reader();
+ _handles->table.reset();
+ _opened_table_key.reset();
+}
+
+Status PaimonRustTableReader::_apply_predicate() {
+ if (_conjuncts.empty() || !_handles || !_handles->table ||
!_handles->read_builder) {
+ return Status::OK();
+ }
+ if (_scanner_profile != nullptr) {
+ COUNTER_UPDATE(_rust_predicates_input, _conjuncts.size());
+ for (const auto& conjunct : _conjuncts) {
+ if (conjunct && conjunct->root() &&
conjunct->root()->is_rf_wrapper()) {
+ COUNTER_UPDATE(_rust_runtime_filters_input, 1);
+ }
+ }
+ }
+ LOG(INFO) << "paimon-rust predicate pushdown: " << _conjuncts.size() << "
conjunct(s) input";
Review Comment:
[P2] Avoid INFO logging for every filtered split. `_apply_predicate` runs
once per Rust logical split and emits this line plus an applied/none line
whenever conjuncts exist. A 50,000-split scan therefore writes about 100,000
INFO records, adding avoidable log I/O and rotation pressure to a hot scan
path. Put per-split details behind verbose logging or use the existing
predicate profile counters for the aggregate result.
--
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]