github-actions[bot] commented on code in PR #66227: URL: https://github.com/apache/doris/pull/66227#discussion_r4121052244
########## regression-test/suites/external_table_p0/paimon/test_paimon_rust_reader_cast_predicates.groovy: ########## @@ -0,0 +1,186 @@ +// 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. + +// CASTed operands through the paimon rust reader. The rust predicate converter +// must NOT strip a cast from a column before pushing: for a DECIMAL(10,2) +// column holding 1.24, `CAST(amount AS DECIMAL(10,1)) = 1.2` must retain the +// row, while the unwrapped `amount = 1.2` prunes it — and rows pruned by the +// pushed filter cannot be recovered by the Doris residual. The FE converter +// already rejects casted operands (PaimonPredicateConverterTest's +// rejectDecimalScaleCast / rejectStringToIntegerCast), so JNI never pushes +// them; the rust reader must match. Cast-induced NULL semantics +// (`CAST(s AS INT) IS NULL`) are covered the same way. +// +// Both differential legs run with force_jni_scanner=true: these parquet +// append tables convert to raw native splits, which getSplits() would +// otherwise prefer — both legs would silently use the native reader and +// never reach the JNI / rust converters. The actual reader path is verified +// per leg through the query profile (the rust reader's PaimonRustReader +// timer group). +suite("test_paimon_rust_reader_cast_predicates", "p0,external,paimon") { + String enabled = context.config.otherConfigs.get("enablePaimonTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disabled paimon test") + return + } + + String catalogName = "test_paimon_rust_cast" + String dbName = "test_paimon_rust_cast_db" + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + + // Table is created via Spark because Doris does not support Paimon DDL. + // t_decimal_cast mirrors the FE rejectDecimalScaleCast shape: a + // DECIMAL(10,2) column whose scale-reducing cast must not be pushed. + // t_cast_null covers cast-related NULL semantics on a string column. + spark_paimon_multi """ + CREATE DATABASE IF NOT EXISTS paimon.${dbName}; + DROP TABLE IF EXISTS paimon.${dbName}.t_decimal_cast; + CREATE TABLE paimon.${dbName}.t_decimal_cast ( + id INT, amount DECIMAL(10, 2) + ) USING paimon; + INSERT INTO paimon.${dbName}.t_decimal_cast VALUES + (1, 1.24), (2, 1.20), (3, 1.30), (4, NULL); + + DROP TABLE IF EXISTS paimon.${dbName}.t_cast_null; + CREATE TABLE paimon.${dbName}.t_cast_null ( + id INT, s VARCHAR(10) + ) USING paimon; + INSERT INTO paimon.${dbName}.t_cast_null VALUES (1, '5'), (2, NULL); + """ + + // The s3.region property is required: paimon-rust's S3 client rejects a + // missing region, while the JNI reader falls back to the SDK default. + sql """drop catalog if exists ${catalogName}""" + sql """ + CREATE CATALOG ${catalogName} PROPERTIES ( + 'type' = 'paimon', + 'paimon.catalog.type' = 'filesystem', + 'warehouse' = 's3://warehouse/wh', + 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}', + 's3.access_key' = 'admin', + 's3.secret_key' = 'password', + 's3.region' = 'us-east-1', + 'use_path_style' = 'true' + ); + """ + + // Capture the settings this suite overrides so finally can restore them. + def originalForceJni = sql("select @@force_jni_scanner")[0][0] + def originalEnableProfile = sql("select @@enable_profile")[0][0] + + try { + sql """switch ${catalogName}""" + sql """use ${dbName}""" + sql """set enable_file_scanner_v2=true""" + // These are parquet append tables whose DataSplits convert to raw + // native splits; force the logical (JNI / rust) reader path so the + // differential actually exercises both converters. + sql """set force_jni_scanner=true""" + sql """set enable_profile=true""" + + // Runs one query and returns its profile text via the FE REST API + // (`show query profile "/<id>"` only lists profiles in this version). + // The profile is finalized asynchronously after the query returns, so + // retry briefly until the endpoint serves the finished body. + def profileTextOf = { String query -> + sql(query) + def queryId = sql("select last_query_id()")[0][0] + for (int i = 0; i < 10; i++) { + def (code, out, err) = curl("GET", + "http://${context.config.feHttpAddress}/rest/v1/query_profile/text/${queryId}", + null, 30, "root", "") Review Comment: [P2] Use ProfileAction with the configured FE HTTP credentials for this differential. The hardcoded root/empty-password curl times out when FE HTTP authentication is configured, even though the query succeeds. It also returns as soon as FileScannerV2 appears, so the Rust timer can still be absent and the path assertion becomes flaky. In the Rust leg, wait for both FileScannerV2 and PaimonRustReader through ProfileAction; its readiness check also waits for COMPLETE when the profile has a completion-state marker. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonRustReaderCapabilities.java: ########## @@ -0,0 +1,226 @@ +// 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 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 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); + // Java numeric-to-integer casts can wrap; Arrow may return NULL for the same value. + // Compare IDs recursively: renames and newly added fields are not narrowing. + return fileSchema != null && !hasIntegerNarrowing(fileSchema.fields(), schema.fields()); Review Comment: [P1] Keep historical character-string-to-DATE and STRING-to-TIMESTAMP(3) files on JNI until Rust matches Paimon's casts. This gate rejects integer narrowing only, so it admits a file written with STRING v and read after v changes to nullable DATE. Paimon 1.4.2 allows that evolution by default and its Java cast interprets the stored string "42" as epoch day 42 (1970-02-12). The pinned Rust reader uses arrow-cast 58.4.0, whose Date32 parser rejects "42" and produces NULL under its safe cast. STRING-to-TIMESTAMP(3) has the same gap: Java interprets "42" as 42 epoch milliseconds, while Arrow rejects the short timestamp string. Enabling Rust can silently replace valid values with NULL; please add persisted JNI/Rust differentials for both casts. ########## be/src/format_v2/table/paimon_rust_table_reader.cpp: ########## @@ -0,0 +1,912 @@ +// 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); + 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. + paimon_result_read_builder rb_res = paimon_table_new_read_builder(_handles->table.get()); + 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"; + // The conjunct VSlotRefs carry table global indices (positions), so the v2 + // converter mode resolves fields by the projected column names; partition + // keys are excluded because the rust reader does not read them. + std::vector<std::string> names; + std::vector<DataTypePtr> types; + names.reserve(_projected_columns.size()); + types.reserve(_projected_columns.size()); + for (const auto& col : _projected_columns) { + if (col.is_partition_key) { + continue; + } + names.push_back(col.name); + types.push_back(col.type); + } + PaimonRustPredicateConverter converter(names, types, _handles->table.get()); + paimon_predicate* predicate = converter.build(_conjuncts); + if (_scanner_profile != nullptr) { + COUNTER_UPDATE(_rust_predicates_converted, converter.converted_conjuncts()); + } + if (predicate == nullptr) { + LOG(INFO) << "paimon-rust predicate pushdown: nothing convertible, no filter applied"; + return Status::OK(); + } + // paimon_read_builder_with_filter consumes the predicate (ownership moves to + // the builder) on every path, so we must not free it here. + if (paimon_error* err = + paimon_read_builder_with_filter(_handles->read_builder.get(), predicate)) { + return Status::InternalError("paimon-rust apply filter failed: {}", consume_error(err)); + } + // Count application only after the C API accepts the filter; reader timers + // alone cannot distinguish pushdown from Doris residual-only execution. + if (_scanner_profile != nullptr) { + COUNTER_UPDATE(_rust_predicates_applied, converter.converted_conjuncts()); + COUNTER_UPDATE(_rust_runtime_filters_applied, converter.converted_runtime_filters()); + } + LOG(INFO) << "paimon-rust predicate pushdown: applied"; + return Status::OK(); +} + +Status PaimonRustTableReader::_fill_block_from_record_batch( + const std::shared_ptr<arrow::RecordBatch>& batch, Block* block, size_t rows) { + SCOPED_TIMER(_rust_arrow_to_block_time); + DORIS_CHECK(batch != nullptr); + DORIS_CHECK(block != nullptr); + std::unordered_set<size_t> materialized_indices; + materialized_indices.reserve(_projected_columns.size()); + { + auto columns_guard = block->mutate_columns_scoped(); + auto& columns = columns_guard.mutable_columns(); + for (int c = 0; c < batch->num_columns(); ++c) { + const auto& field = batch->schema()->field(c); + if (field->name() == VALUE_KIND_FIELD) { + continue; + } + // Projected column names are FE-normalized to lowercase. + // paimon-rust's case_sensitive=false setting also case-folds column + // names in the schema output, so exact match works — but tolerate + // mixed-case Rust output by folding here as well. + auto it = _output_name_to_idx.find(field->name()); + if (it == _output_name_to_idx.end()) { + it = _output_name_to_idx.find(to_lower(field->name())); + } + if (it == _output_name_to_idx.end()) { + // Skip columns that are not in the block (e.g. columns dropped by + // slot pruning). + continue; + } + const auto output_idx = it->second; + if (!materialized_indices.emplace(output_idx).second) { + return Status::InternalError("paimon-rust returned duplicate column '{}'", + field->name()); + } + try { + RETURN_IF_ERROR(columns_guard.get_datatype_by_position(output_idx) + ->get_serde() + ->read_column_from_arrow(*columns[output_idx], + batch->column(c).get(), 0, rows, + _ctz)); + } catch (Exception& e) { + return Status::InternalError("Failed to convert from arrow to block: {}", e.what()); + } + } + } + // Partition columns and other projected columns absent from the arrow batch + // are back-filled from split metadata / defaults. + RETURN_IF_ERROR(_fill_non_arrow_columns(block, rows, materialized_indices)); + // This direct Arrow path bypasses TableReader::finalize_chunk, whose last + // step enforces truncate_char_or_varchar_columns — without it, a column + // narrowed by schema evolution returns untruncated historical values. + RETURN_IF_ERROR(_truncate_char_or_varchar_columns(block)); + return Status::OK(); +} + +Status PaimonRustTableReader::_truncate_char_or_varchar_columns(Block* block) { + if (_runtime_state == nullptr || + !_runtime_state->query_options().truncate_char_or_varchar_columns) { + return Status::OK(); + } + for (size_t idx = 0; idx < block->columns(); ++idx) { + const auto& column_type = block->get_by_position(idx).type; + if (column_type == nullptr) { + continue; + } + const auto type = remove_nullable(column_type); + const auto primitive = type->get_primitive_type(); + if (primitive != TYPE_VARCHAR && primitive != TYPE_CHAR) { + continue; + } + const auto target_len = assert_cast<const DataTypeString*>(type.get())->len(); + if (target_len <= 0) { + continue; + } + // Reuses TableReader's vectorized truncation (substring(column, 1, + // len)); the base variant maps through column_mapper metadata, which + // the direct rust path does not populate, so iterate the block's own + // slot-derived types here — the arrow side is always lengthless Utf8, + // so any bounded CHAR/VARCHAR target truncates to its declared length. + _truncate_char_or_varchar_column(block, idx, target_len); Review Comment: [P1] Materialize bounded partition constants before calling this helper. With truncate_char_or_varchar_columns=true, a projected nullable VARCHAR/CHAR partition key is filled by VLiteral as ColumnConst(ColumnNullable), while _truncate_char_or_varchar_column casts the column directly to ColumnNullable and throws a bad-cast error. The existing TableReader path materializes constants in _align_column_nullability first, but this direct Rust path does not. A scan of a partitioned table therefore fails instead of returning rows; please cover this shape in the truncation test. -- 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]
