github-actions[bot] commented on code in PR #66227:
URL: https://github.com/apache/doris/pull/66227#discussion_r4128963657
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java:
##########
@@ -411,10 +728,335 @@ private void setPaimonParams(TFileRangeDesc rangeDesc,
PaimonSplit paimonSplit)
String fileFormat = getFileFormat(paimonSplit.getPathString());
if (split != null) {
+ // use jni reader / paimon-cpp reader / paimon-rust reader
rangeDesc.setFormatType(TFileFormatType.FORMAT_JNI);
- // A logical DataSplit may span multiple files, so keep it intact
for the JNI reader.
- fileDesc.setReaderType(TPaimonReaderType.PAIMON_JNI);
- fileDesc.setPaimonSplit(PaimonUtil.encodeObjectToString(split));
+ // paimon-cpp and paimon-rust both consume Paimon native binary
serialization,
+ // which only supports DataSplit. Any other split type falls back
to JNI.
+ boolean nativeSplit = split instanceof DataSplit;
+ // Fallback-read splits stay on JNI: FallbackDataSplit extends
+ // DataSplit, so the instanceof above passes, but its serializer
+ // appends an isFallback byte after the ordinary split that the
+ // pinned rust decoder rejects outright ("trailing bytes after
+ // DataSplit" — it requires full-buffer consumption), and even a
+ // permissive decode would still lack the second table identity
+ // needed to honor the fallback-side discriminator. Both sides of a
+ // FallbackReadFileStoreTable wrap their splits, so the table
+ // wrapper is gated as a whole (any split from it routes to JNI)
+ // until the rust ABI represents both sides; the FallbackSplit
+ // interface also catches a wrapper split regardless of how the
+ // table was resolved here.
+ boolean fallbackRead = split instanceof
FallbackReadFileStoreTable.FallbackSplit
+ || processedTable instanceof FallbackReadFileStoreTable;
+ // Serialize the same effective table that planning and the JNI
reader use.
+ // Relation options such as t@options('read.batch-size'='1') are
applied by
+ // getProcessedTable() (doInitialize caches it in processedTable),
and the
+ // rust reader derives its read batch size from the schema options
— the raw
+ // cached table would silently drop the override. Copies,
delegates and
+ // fallback wrappers of getProcessedTable() are still
FileStoreTable, so the
+ // instanceof gate keeps its semantics.
+ Table paimonTable = processedTable;
+ FileStoreTable paimonFileStoreTable =
+ paimonTable instanceof FileStoreTable ? (FileStoreTable)
paimonTable : null;
+ // query-auth.enabled tables stay on JNI: when catalog
authorization
+ // succeeds with no row filter or column mask, Paimon still leaves
an
+ // ordinary DataSplit (restricted results use QueryAuthSplit and
are
+ // already handled by the nativeSplit gate above), so this table
shape
+ // passes the compound gate — but the shipped schema keeps
+ // query-auth.enabled=true and the pinned rust ReadBuilder rejects
+ // every such table (its CoreOptions::ensure_read_authorized fails
+ // closed because the client cannot enforce the row filter / column
+ // masking), turning a valid authorized scan into a BE-open
failure.
+ // Until the authorization result can be transported and enforced
by
+ // the rust ABI, these tables route to JNI.
+ boolean queryAuthTable = false;
+ // REST-token tables stay on JNI: doInitialize snapshots
+ // RESTTokenFileIO.validToken().token() into the backend storage
+ // properties, discarding expireAtMillis and the REST refresh
+ // context, so the shipped credentials look static — but the
+ // pinned rust table reuses one option map with no refresh
+ // callback, while paimon 1.4.2's JNI RESTTokenFileIO checks
+ // expiry before each file operation and obtains a replacement
+ // token. A queued or long scan that crosses the token TTL would
+ // start on rust and later fail authentication. Gate until the
+ // rust ABI can refresh and atomically update credentials.
+ boolean restTokenTable = false;
+ // Partial-update / aggregation tables with deletion vectors only
pass
+ // the rust reader in the fully materialized shape: the pinned rust
+ // read_pk rejects merge-engine=partial-update/aggregation with
+ // deletion-vectors.merge-on-read=true outright, and otherwise
requires
+ // every split to be compacted and known free of retract rows
+ // (DataSplit::is_fully_materialized_pk_dv). Their ordinary
DataSplits
+ // sail through the compound gate above, so without this check a
valid
+ // Java/JNI scan reaches BE and the rust open fails. Deduplicate
stays
+ // rust-eligible: its read_pk routes uncompacted splits to the KV
+ // reader, which applies the attached per-file DVs.
merge-on-read=true
+ // is a table option, so the whole table routes to JNI;
+ // non-materialized splits are gated per split below.
+ boolean puAggDeletionVectors = false;
+ boolean dvMergeOnRead = false;
+ boolean deduplicateIgnoreDelete = false;
+ boolean rustUnsupportedMergeOption = false;
+ if (paimonFileStoreTable != null) {
+ // A renewable REST token reached the shipped properties as a
+ // plain value; only the table's FileIO type reveals it
expires.
+ // Null-safe: a table handle whose FileIO is not resolved stays
+ // rust-eligible, mirroring the CoreOptions null-safety below.
+ restTokenTable = paimonFileStoreTable.fileIO() instanceof
RESTTokenFileIO;
+ CoreOptions resolvedCoreOptions =
paimonFileStoreTable.coreOptions();
+ // Null-safe: a table handle whose CoreOptions is not resolved
+ // (e.g. some wrapper shapes) stays rust-eligible rather than
+ // failing the scan here — the rust open itself rejects such a
+ // table if the option is really set.
+ if (resolvedCoreOptions != null) {
+ queryAuthTable = resolvedCoreOptions.queryAuthEnabled();
+ CoreOptions.MergeEngine mergeEngine =
resolvedCoreOptions.mergeEngine();
+ if (resolvedCoreOptions.deletionVectorsEnabled()
+ && (mergeEngine ==
CoreOptions.MergeEngine.PARTIAL_UPDATE
+ || mergeEngine ==
CoreOptions.MergeEngine.AGGREGATE)) {
+ puAggDeletionVectors = true;
+ // The merge-engine and deletion-vectors.enabled checks
+ // above resolve through the Java CoreOptions
accessors,
+ // which the table builds from this same schema options
+ // map — the one the BE rust reader deserializes from
+ // the shipped schema JSON — so they cannot diverge
from
+ // what BE sees. merge-on-read has no Java accessor in
+ // paimon 1.4, so it is read raw from the map, with the
+ // rust parsing semantics (any case-insensitive "true"
+ // is on, default false).
+ TableSchema dvSchema = paimonFileStoreTable.schema();
+ Map<String, String> dvOptions = dvSchema == null ?
null : dvSchema.options();
+ String mergeOnRead = dvOptions == null
+ ? null :
dvOptions.get(DELETION_VECTORS_MERGE_ON_READ);
+ dvMergeOnRead = "true".equalsIgnoreCase(mergeOnRead);
+ }
+ // deduplicate.ignore-delete=true tables stay on JNI:
+ // Java's DeduplicateMergeFunction skips retract records
+ // when the option is set — including old, uncompacted
+ // files that still contain them — but the pinned rust
+ // read_pk does not pass table options into its
+ // deduplicate merge: it picks the latest row and omits
+ // the key when that row is DELETE/UPDATE_BEFORE. An
+ // uncompacted insert followed by a delete therefore
+ // returns the insert through JNI but silently disappears
+ // through rust. Gate the option until the rust merge
+ // implements it.
+ if (mergeEngine == CoreOptions.MergeEngine.DEDUPLICATE
+ && resolvedCoreOptions.ignoreDelete()) {
+ deduplicateIgnoreDelete = true;
+ }
+ // Non-DV merge options the pinned rust read rejects: Java
+ // supports partial-update.remove-record-on-delete /
+ // aggregation.remove-record-on-delete and the wider
+ // per-field retract matrix, but the rust
+ // PartialUpdateConfig / AggregationConfig validations
+ // return Unsupported for them — and the DV-derived gates
+ // above only cover deletion-vector tables, so an ordinary
+ // non-DV DataSplit with one of these options would pass
the
+ // compound gate and fail during the rust merge
+ // construction. Mirror the exact rust key matrix
(presence,
+ // not values) against the same schema options map BE
+ // deserializes.
+ if (mergeEngine == CoreOptions.MergeEngine.PARTIAL_UPDATE
+ || mergeEngine ==
CoreOptions.MergeEngine.AGGREGATE) {
+ TableSchema mergeSchema =
paimonFileStoreTable.schema();
+ Map<String, String> mergeOptions =
+ mergeSchema == null ? null :
mergeSchema.options();
+ rustUnsupportedMergeOption = mergeOptions != null
+ && hasRustUnsupportedMergeOption(mergeOptions,
mergeEngine);
+ }
+ }
+ }
+ // paimon-rust additionally requires (a) FileScannerV2: the V1
FileScanner
+ // explicitly rejects PAIMON_RUST, so with enable_file_scanner_v2
disabled
+ // the split falls back to JNI instead of encoding a rust request
that the
+ // selected scanner cannot consume, and (b) a FileStoreTable: BE
opens the
+ // table via paimon_table_from_schema_json, which needs the
resolved
+ // TableSchema that only FileStoreTable exposes via schema(). If
the table
+ // is not a FileStoreTable (e.g. a sys table backed by DataSplit),
we cannot
+ // ship a schema JSON, so fall back to CPP / JNI rather than
sending an
+ // incomplete PAIMON_RUST request that BE would reject.
+ //
+ // The paimon-rust S3 bridge maps static credentials, anonymous
+ // access (AWS_CREDENTIALS_PROVIDER_TYPE=ANONYMOUS -> s3.anonymous)
+ // and assume-role (AWS_ROLE_ARN / AWS_EXTERNAL_ID ->
+ // s3.assumed.role.*), but the remaining credential-provider modes
+ // are ambient JVM provider chains (ENV, SYSTEM_PROPERTIES,
+ // WEB_IDENTITY, CONTAINER, INSTANCE_PROFILE) with no paimon-rust
+ // equivalent — rust would silently sign with whatever the ambient
+ // chain resolves to. Gate those modes away from the rust reader
+ // here so the configured provider is honored via the JNI path.
+ boolean providerModeTranslatable = true;
+ String providerType = backendStorageProperties == null
+ ? null :
backendStorageProperties.get("AWS_CREDENTIALS_PROVIDER_TYPE");
+ String mode = providerType == null ? "DEFAULT" :
providerType.trim().toUpperCase(Locale.ROOT);
+ if (mode.isEmpty()) {
+ mode = "DEFAULT";
+ }
+ String location = source.getTableLocation();
+ if (location != null && (location.startsWith("s3://") ||
location.startsWith("s3a://"))) {
+ // Java DEFAULT may resolve JVM properties or anonymous
credentials;
+ // Rust's ambient chain is different. Only explicit keys or
explicit
+ // anonymous access can cross this boundary without changing
identity.
Review Comment:
[P1] Preserve the Java credential choice for S3 ANONYMOUS catalogs. Doris
accepts a catalog with both access/secret keys and
`s3.credentials_provider_type=ANONYMOUS`; its Java/Hadoop Paimon path installs
`SimpleAWSCredentialsProvider` whenever keys are set. This gate admits it
because the mode is ANONYMOUS, then the BE sets `s3.anonymous=true`, which
makes the pinned Rust reader [skip
signing](https://github.com/apache/paimon-rust/blob/v0.4.0-rc1/crates/paimon/src/io/storage_s3.rs#L108-L111)
even though it also received the keys. A private-bucket scan that works on JNI
fails when Rust is enabled. Set Rust anonymous mode only when Java would
actually use anonymous access, or route this combination to JNI; test the
crossed setting.
##########
thirdparty/build-thirdparty.sh:
##########
@@ -2162,11 +2185,133 @@ build_lance_c() {
mkdir -p "${TP_INSTALL_DIR}/include" "${TP_INSTALL_DIR}/lib64"
rm -rf "${TP_INSTALL_DIR}/include/lance"
cp -av include/lance "${TP_INSTALL_DIR}/include/"
- cp -v "${BUILD_DIR}/release/liblance_c.a" "${TP_INSTALL_DIR}/lib64/"
+ install_rust_archive "${BUILD_DIR}/release/liblance_c.a"
+}
+
+# paimon-rust
+build_paimon_rust() {
+ check_if_source_exist "${PAIMON_RUST_SOURCE}"
+ cd "${TP_SOURCE_DIR}/${PAIMON_RUST_SOURCE}"
+
+ rm -rf "${BUILD_DIR}"
+ mkdir -p "${BUILD_DIR}"
+
+ local cargo_bin="${PAIMON_RUST_CARGO:-${CARGO:-cargo}}"
+ if ! command -v "${cargo_bin}" >/dev/null 2>&1; then
+ echo "cargo is required to build paimon-rust. Install Rust 1.94.0 or
set PAIMON_RUST_CARGO."
+ exit 1
+ fi
+
+ local required_rust_version="1.94.0"
+ local cargo_env=(
+ "CARGO_BUILD_JOBS=${PARALLEL}"
+ "CARGO_TARGET_DIR=${PWD}/${BUILD_DIR}"
+ )
+ if command -v rustup >/dev/null 2>&1 && [[ -z "${RUSTUP_TOOLCHAIN}" ]];
then
+ # The presence check must look for the toolchain the minimum actually
+ # requires, not a literal: with only an older toolchain installed the
+ # stale check would skip the install below and then force
+ # RUSTUP_TOOLCHAIN to a version rustup cannot dispatch, failing the
+ # build before any archive is produced.
+ local required_rust_regex="${required_rust_version//./\\.}"
+ if ! rustup toolchain list | grep -Eq
"^${required_rust_regex}([[:space:]-]|$)"; then
+ rustup toolchain install "${required_rust_version}" --profile
minimal
+ fi
+ cargo_env+=("RUSTUP_TOOLCHAIN=${required_rust_version}")
+ fi
+
+ local cargo_version
+ if ! cargo_version="$(env "${cargo_env[@]}" "${cargo_bin}" --version | awk
'{print $2}')"; then
+ echo "failed to get cargo version for paimon-rust. Install Rust
${required_rust_version} or set PAIMON_RUST_CARGO/RUSTUP_TOOLCHAIN."
+ exit 1
+ fi
+ # Rust 1.94.0 is the minimum supported version. Allow newer toolchains when
+ # callers explicitly select one or rustup is unavailable on the system.
+ # NOTE: paimon_c and lance_c are both Rust staticlibs linked into the same
+ # BE binary; they must be built with the SAME rustc toolchain so the linker
+ # resolves both crates' std references against a single std copy. Mixing
+ # toolchains makes the precompiled std hashes differ and the linker pulls
+ # both std copies in, colliding on the unmangled `rust_eh_personality`
+ # (duplicate symbol). Rebuild both lance_c and paimon_rust whenever the
+ # toolchain changes, using the same RUSTUP_TOOLCHAIN for both builds.
+ if ! awk -v required="${required_rust_version}" -v
actual="${cargo_version}" 'BEGIN {
+ split(required, r, ".");
+ split(actual, a, ".");
+ for (i = 1; i <= 3; i++) {
+ if ((a[i] + 0) > (r[i] + 0)) {
+ exit 0;
+ }
+ if ((a[i] + 0) < (r[i] + 0)) {
+ exit 1;
+ }
+ }
+ exit 0;
+ }'; then
+ echo "paimon-rust requires Rust/Cargo ${required_rust_version} or
newer, but found ${cargo_version}."
+ echo "Install Rust ${required_rust_version} or set
PAIMON_RUST_CARGO/RUSTUP_TOOLCHAIN."
+ exit 1
+ fi
- if [[ "${STRIP_TP_LIB}" = "ON" && "${KERNEL}" != 'Darwin' ]]; then
- strip --strip-debug --strip-unneeded
"${TP_INSTALL_DIR}/lib64/liblance_c.a"
+ if [[ "${KERNEL}" != 'Darwin' ]]; then
+ cargo_env+=("CFLAGS=${CFLAGS:-} -std=gnu17")
+ fi
+
+ local cargo_args=(build --release --locked -p paimon-c --features
paimon/storage-hdfs)
+ # cbindgen invokes cargo metadata itself; command-line flags on the build
+ # and install calls do not propagate to that child process.
+ cargo_env+=("CARGO=${cargo_bin}")
+ if [[ "$(echo "${PAIMON_RUST_CARGO_OFFLINE}" | tr '[:lower:]'
'[:upper:]')" == "ON" ]]; then
+ cargo_args+=(--offline)
+ cargo_env+=("CARGO_NET_OFFLINE=true")
+ fi
+ env "${cargo_env[@]}" "${cargo_bin}" "${cargo_args[@]}"
+
+ # Generate the C header from the Rust extern "C" surface via cbindgen.
+ # cbindgen is a pinned, Doris-controlled input: an unpinned "current"
+ # release would regenerate paimon.h differently between builds. The
+ # pinned version installs under a Doris-controlled --root and its
+ # resolved absolute path is invoked directly (a custom CARGO_HOME does
+ # not necessarily put cargo-installed binaries on PATH). Offline builds
+ # pass --offline to the install command, exactly like the fetch/build
+ # handling above — cargo fails on a missing local crate cache instead
+ # of reaching for the network.
+ local cbindgen_version="0.29.4"
+ local cbindgen_bin="${PAIMON_RUST_CBINDGEN:-}"
+ if [[ -z "${cbindgen_bin}" ]]; then
+ local
cbindgen_root="${TP_SOURCE_DIR}/.doris-cbindgen-${cbindgen_version}"
+ cbindgen_bin="${cbindgen_root}/bin/cbindgen"
+ if [[ ! -x "${cbindgen_bin}" ]]; then
+ local cbindgen_install_args=(install cbindgen
+ --version "${cbindgen_version}" --locked --root
"${cbindgen_root}")
+ if [[ "$(echo "${PAIMON_RUST_CARGO_OFFLINE}" | tr '[:lower:]'
'[:upper:]')" == "ON" ]]; then
+ cbindgen_install_args+=(--offline)
+ fi
+ echo "cbindgen not found; installing pinned ${cbindgen_version}
via cargo install ..."
+ env "${cargo_env[@]}" "${cargo_bin}" "${cbindgen_install_args[@]}"
+ fi
+ elif [[ ! -x "${cbindgen_bin}" ]]; then
+ echo "PAIMON_RUST_CBINDGEN=${cbindgen_bin} is not an executable file."
+ exit 1
fi
+ # Write a temporary cbindgen.toml so the generated header carries our
+ # include-guard / cpp-compat settings without touching the upstream tree.
+ local cbindgen_toml="${BUILD_DIR}/cbindgen.toml"
+ mkdir -p "${BUILD_DIR}"
+ cat >"${cbindgen_toml}" <<'EOF'
+language = "C"
+include_guard = "PAIMON_C_H"
+pragma_once = true
+cpp_compat = true
+EOF
+ env "${cargo_env[@]}" "${cbindgen_bin}" bindings/c \
+ --config "${cbindgen_toml}" \
+ --output "${BUILD_DIR}/release/paimon.h"
+
Review Comment:
[P2] Publish the generated header as part of a recoverable library pair.
This direct `cp` can leave an empty or partial installed `paimon.h` after an
interrupted replacement, while the previous archive remains. If archive
publication fails after the header copy, new declarations can also remain
beside the old archive. `build.sh` only checks that both paths are regular
files, so the next build skips repair. Stage the header and validate the
published header/archive pair before reuse; cover header-copy and
pair-publication failures.
##########
be/src/format_v2/table/paimon_rust_table_reader.cpp:
##########
@@ -0,0 +1,918 @@
+// 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.
+ 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.
+ // Partition/default expressions can produce
ColumnConst(ColumnNullable). The shared
+ // truncation helper expects a materialized nullable column, as in
finalize_chunk.
+ auto& column = block->get_by_position(idx).column;
+ column = column->convert_to_full_column_if_const();
+ _truncate_char_or_varchar_column(block, idx, target_len);
+ }
+ return Status::OK();
+}
+
+Status PaimonRustTableReader::_fill_non_arrow_columns(
+ Block* block, size_t rows, const std::unordered_set<size_t>&
materialized_indices) {
+ for (size_t idx = 0; idx < _projected_columns.size(); ++idx) {
+ if (materialized_indices.count(idx) != 0) {
+ continue;
+ }
+ const auto& column = _projected_columns[idx];
+ VExprContextSPtr constant_expr;
+ if (const Field* value = find_partition_value(column,
_partition_values);
+ column.is_partition_key && value != nullptr) {
+ // Partition values are split constants (same materialization the
+ // TableColumnMapper builds for native readers).
+ constant_expr =
+
VExprContext::create_shared(VLiteral::create_shared(column.type, *value));
+ } else if (column.default_expr != nullptr) {
+ constant_expr = column.default_expr;
+ } else {
+ // The column is genuinely absent from the arrow batch. Schema
+ // evolution is handled by paimon-rust itself, so reaching here
means
+ // an unexpected schema drift: fill defaults so the scan remains
+ // well-defined instead of failing the query.
+ LOG(WARNING) << "paimon-rust did not return projected column '" <<
column.name
+ << "'; filling with defaults";
+ auto data = column.type->create_column();
+ data->insert_many_defaults(rows);
+ block->replace_by_position(idx, std::move(data));
+ continue;
+ }
+ ColumnPtr constant_column;
+ RETURN_IF_ERROR(_materialize_constant_column(constant_expr,
column.type, column.name, rows,
+ &constant_column));
+ block->replace_by_position(idx, std::move(constant_column));
+ }
+ return Status::OK();
+}
+
+Status PaimonRustTableReader::_materialize_constant_column(const
VExprContextSPtr& expr,
+ const DataTypePtr&
type,
+ const std::string&
name, size_t rows,
+ ColumnPtr* column) {
+ DORIS_CHECK(expr != nullptr);
+ DORIS_CHECK(column != nullptr);
+ RowDescriptor row_desc;
+ RETURN_IF_ERROR(expr->prepare(_runtime_state, row_desc));
+ RETURN_IF_ERROR(expr->open(_runtime_state));
+ // Constants evaluate per input row, so a rows-sized synthetic block
yields a
+ // rows-sized result for both plain literals and default expressions.
+ Block eval_block;
+ eval_block.insert({type->create_column_const_with_default_value(rows),
type, name});
+ int result_column_id = -1;
+ RETURN_IF_ERROR(expr->execute(&eval_block, &result_column_id));
+ DORIS_CHECK(result_column_id >= 0);
+ ColumnPtr result_column =
eval_block.get_by_position(result_column_id).column;
+ if (result_column->size() == 1 && rows > 1) {
+ result_column = ColumnConst::create(std::move(result_column), rows);
+ }
+ *column = std::move(result_column);
+ return Status::OK();
+}
+
+Status PaimonRustTableReader::_decode_split_bytes(std::string* out) const {
+ if (!_current_range.__isset.table_format_params ||
+ !_current_range.table_format_params.__isset.paimon_params ||
+
!_current_range.table_format_params.paimon_params.__isset.paimon_split) {
+ return Status::InternalError("paimon-rust missing paimon_split in scan
range");
+ }
+ const auto& encoded_split =
_current_range.table_format_params.paimon_params.paimon_split;
+ if (!base64_decode(encoded_split, out)) {
+ return Status::InternalError("paimon-rust base64 decode paimon_split
failed");
+ }
+ if (out->empty()) {
+ return Status::InternalError("paimon-rust decoded paimon_split is
empty");
+ }
+ return Status::OK();
+}
+
+std::optional<std::string> PaimonRustTableReader::_resolve_table_path(
+ const TFileRangeDesc& range) const {
+ if (range.__isset.table_format_params &&
range.table_format_params.__isset.paimon_params &&
+ range.table_format_params.paimon_params.__isset.paimon_table &&
+ !range.table_format_params.paimon_params.paimon_table.empty()) {
+ return range.table_format_params.paimon_params.paimon_table;
+ }
+ return std::nullopt;
+}
+
+std::optional<std::string> PaimonRustTableReader::_resolve_db_name(
+ const TFileRangeDesc& range) const {
+ if (range.__isset.table_format_params &&
range.table_format_params.__isset.paimon_params &&
+ range.table_format_params.paimon_params.__isset.db_name &&
+ !range.table_format_params.paimon_params.db_name.empty()) {
+ return range.table_format_params.paimon_params.db_name;
+ }
+ return std::nullopt;
+}
+
+std::optional<std::string> PaimonRustTableReader::_resolve_table_name(
+ const TFileRangeDesc& range) const {
+ if (range.__isset.table_format_params &&
range.table_format_params.__isset.paimon_params &&
+ range.table_format_params.paimon_params.__isset.table_name &&
+ !range.table_format_params.paimon_params.table_name.empty()) {
+ return range.table_format_params.paimon_params.table_name;
+ }
+ return std::nullopt;
+}
+
+std::optional<std::string> PaimonRustTableReader::_resolve_table_schema_json(
+ const TFileRangeDesc& range) const {
+ if (range.__isset.table_format_params &&
range.table_format_params.__isset.paimon_params &&
+
range.table_format_params.paimon_params.__isset.paimon_table_schema_json &&
+
!range.table_format_params.paimon_params.paimon_table_schema_json.empty()) {
+ return
range.table_format_params.paimon_params.paimon_table_schema_json;
+ }
+ return std::nullopt;
+}
+
+std::optional<std::string> PaimonRustTableReader::_resolve_branch(
+ const TFileRangeDesc& range) const {
+ // FE only sets paimon_branch when the branch is not `main` (matches
+ // upstream paimon commit 742da63: null-if-DEFAULT_MAIN_BRANCH). Unset here
+ // means main-branch semantics.
+ if (range.__isset.table_format_params &&
range.table_format_params.__isset.paimon_params &&
+ range.table_format_params.paimon_params.__isset.paimon_branch &&
+ !range.table_format_params.paimon_params.paimon_branch.empty()) {
+ return range.table_format_params.paimon_params.paimon_branch;
+ }
+ return std::nullopt;
+}
+
+std::vector<std::string> PaimonRustTableReader::_build_read_columns() const {
+ std::vector<std::string> columns;
+ columns.reserve(_projected_columns.size());
+ for (const auto& column : _projected_columns) {
+ if (column.is_partition_key) {
+ continue;
+ }
+ columns.emplace_back(column.name);
+ }
+ return columns;
+}
+
+std::map<std::string, std::string> PaimonRustTableReader::_build_options()
const {
+ std::map<std::string, std::string> options;
+ if (_scan_params && _scan_params->__isset.paimon_options &&
+ !_scan_params->paimon_options.empty()) {
+ options.insert(_scan_params->paimon_options.begin(),
_scan_params->paimon_options.end());
+ } else if (_current_range.__isset.table_format_params &&
+ _current_range.table_format_params.__isset.paimon_params &&
+
_current_range.table_format_params.paimon_params.__isset.paimon_options) {
+
options.insert(_current_range.table_format_params.paimon_params.paimon_options.begin(),
+
_current_range.table_format_params.paimon_params.paimon_options.end());
+ }
+
+ if (_scan_params && _scan_params->__isset.properties &&
!_scan_params->properties.empty()) {
+ for (const auto& kv : _scan_params->properties) {
+ options[kv.first] = kv.second;
+ }
+ } else if (_current_range.__isset.table_format_params &&
+ _current_range.table_format_params.__isset.paimon_params &&
+
_current_range.table_format_params.paimon_params.__isset.hadoop_conf) {
+ for (const auto& kv :
_current_range.table_format_params.paimon_params.hadoop_conf) {
+ options[kv.first] = kv.second;
+ }
+ }
+
+ auto copy_if_missing = [&](const char* from_key, const char* to_key) {
+ if (options.find(to_key) != options.end()) {
+ return;
+ }
+ auto it = options.find(from_key);
+ if (it != options.end() && !it->second.empty()) {
+ options[to_key] = it->second;
+ }
+ };
+
+ // The pinned paimon-rust storage dispatcher (io/storage.rs) selects the
+ // FileIO parser from the table path's URI scheme, and libpaimon_c.a
+ // compiles in separate COS, OBS, GCS and Azdls parsers besides the OSS
+ // and S3 ones. Doris's FE normalizes every object store's credentials
+ // into the AWS_* / use_path_style aliases, which this bridge can only
+ // translate into the two key families those two parsers read:
+ // oss:// -> fs.oss.endpoint / fs.oss.accessKeyId /
+ // fs.oss.accessKeySecret (+ fs.oss.securityToken STS)
+ // s3:// / s3a:// -> paimon-java's s3.* family (access-key, secret-key,
+ // session.token, endpoint, region, path-style-access,
+ // normalized from the fs.s3a. / s3a. / s3. prefixes)
+ // The FE therefore gates the rust reader to exactly these schemes
+ // (PaimonScanNode: cosn:// / obs:// / gs:// / abfs:// warehouses fall
+ // back to JNI before the split is encoded — their parsers read the
+ // fs.cosn.userinfo.* / fs.obs.* / gcs.* / azure.* families the AWS_*
+ // aliases cannot express). The mapping here is scoped the same way so a
+ // version-skewed FE cannot smuggle AWS_* aliases into another parser's
+ // property map: a cosn:// table reaching this bridge with synthesized
+ // s3.* keys would hit the COS parser without
+ // fs.cosn.userinfo.secretId / secretKey and fail the open with a
+ // misleading auth error instead of the JNI fallback.
+ const std::string table_path =
_resolve_table_path(_current_range).value_or("");
+ std::string scheme;
+ if (const auto sep = table_path.find("://"); sep != std::string::npos) {
+ scheme = to_lower(table_path.substr(0, sep));
+ }
+ if (scheme == "oss") {
+ // The OSS parser reads only the four fs.oss.* keys (and retry
+ // settings); native fs.oss.* options pass through untouched. It has
+ // no region / anonymous / assume-role handling, so nothing else is
+ // mapped for this scheme.
+ copy_if_missing("AWS_ENDPOINT", "fs.oss.endpoint");
+ copy_if_missing("AWS_ACCESS_KEY", "fs.oss.accessKeyId");
+ copy_if_missing("AWS_SECRET_KEY", "fs.oss.accessKeySecret");
+ copy_if_missing("AWS_TOKEN", "fs.oss.securityToken");
+ return options;
+ }
+ if (scheme == "s3" || scheme == "s3a") {
+ copy_if_missing("AWS_ACCESS_KEY", "s3.access-key");
+ copy_if_missing("AWS_SECRET_KEY", "s3.secret-key");
+ copy_if_missing("AWS_TOKEN", "s3.session.token");
+ copy_if_missing("AWS_ENDPOINT", "s3.endpoint");
+ copy_if_missing("AWS_REGION", "s3.region");
+ copy_if_missing("use_path_style", "s3.path-style-access");
+ // Authentication modes: the FE storage-properties channel marks
anonymous
+ // access with AWS_CREDENTIALS_PROVIDER_TYPE=ANONYMOUS (emitted when no
+ // static credentials are configured) and assume-role with
+ // AWS_ROLE_ARN / AWS_EXTERNAL_ID (from the s3.role_arn /
s3.external_id
+ // catalog properties). The crate reads s3.anonymous (skip_signature)
and
+ // the s3.assumed.role.* family, so map both; without these, anonymous
+ // catalogs would consult the ambient credential chain and role-only
+ // catalogs would never assume the requested role. The remaining
provider
+ // modes are ambient JVM credential chains (ENV, SYSTEM_PROPERTIES,
+ // WEB_IDENTITY, CONTAINER, INSTANCE_PROFILE) with no paimon-rust
+ // equivalent — the FE gates those away from the rust reader before the
+ // split is encoded.
+ if (options.contains("AWS_CREDENTIALS_PROVIDER_TYPE") &&
+ options.at("AWS_CREDENTIALS_PROVIDER_TYPE") == "ANONYMOUS") {
+ options["s3.anonymous"] = "true";
+ }
+ copy_if_missing("AWS_ROLE_ARN", "s3.assumed.role.arn");
Review Comment:
[P1] Preserve static-key precedence before forwarding the role. With both S3
keys and `s3.role_arn`, Java's Hadoop catalog installs
`SimpleAWSCredentialsProvider` and returns before configuring assume-role. The
FE admits the keyed DEFAULT mode, but this mapping sends `s3.assumed.role.arn`
with the keys; [OpenDAL
0.58.2](https://github.com/apache/opendal/blob/v0.58.2/core/services/s3/src/backend.rs#L884-L921)
then wraps the static provider in STS AssumeRole and uses the role identity
for S3. A key pair with bucket read access but no STS permission succeeds
through JNI and fails through Rust. Omit the role option when static keys are
present or keep this combination on JNI; add a crossed-setting test.
##########
regression-test/suites/external_table_p0/paimon/test_paimon_rust_reader_eq_for_null.groovy:
##########
@@ -0,0 +1,559 @@
+// 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.
+
+// NULL-safe equality (`<=>`, EQ_FOR_NULL) through the paimon rust reader.
+// The rust predicate converter must NOT push `a <=> b` down as `a IS NULL`:
+// with rows (NULL, NULL), (1, 1), (1, 2) that would wrongly drop (1, 1), and
+// rows dropped by the pushed filter cannot be recovered by the residual
+// conjunct. FE rewrites the literal forms (`a <=> 1` -> `a = 1`,
+// `a <=> NULL` -> `a IS NULL`) before they reach the BE, so only the
+// column-to-column form exercises EQ_FOR_NULL here; the literal forms still
+// guard the rewrite + pushdown chain end to end.
+//
+// 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), and the join leg forces an IN runtime filter
+// (runtime_filter_type=1 + runtime_filter_wait_infinitely) that must be
+// planned onto the probe scan and arrive before the split opens.
+import org.apache.doris.regression.action.ProfileAction
+
+suite("test_paimon_rust_reader_eq_for_null", "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_eq_null"
+ String dbName = "test_paimon_rust_eq_null_db"
+ String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+ String minioPort = context.config.otherConfigs.get("iceberg_minio_port")
+
+ // Use Spark to prepare the timestamp and nested-schema fixtures.
+ // Both columns stay nullable: `a <=> b` survives FE's NullSafeEqualToEqual
+ // rewrite (which only fires when one side is non-nullable / a NULL
literal)
+ // precisely when both sides are nullable.
+ //
+ // t_frac_ts uses Spark TIMESTAMP_NTZ, which maps to Paimon TIMESTAMP
+ // (wall-clock, microsecond precision). Note Spark's plain TIMESTAMP maps
+ // to Paimon TIMESTAMP_LTZ, so the NTZ semantics must be spelled out in the
+ // DDL. spark.sql.timestampType=TIMESTAMP_NTZ makes the TIMESTAMP '...'
+ // literals parse as NTZ civil times too.
+ //
+ // The predicate literals are millisecond-aligned on purpose: FE truncates
+ // plan-time pushed-down timestamp predicates to 3 fractional digits, so a
+ // 6-digit literal would reach the readers truncated to milliseconds and
+ // the exact residual conjunct would then drop every row it kept — for the
+ // JNI, Rust and native readers alike. Literal conversion in the BE retains
+ // microseconds, but timestamp predicates remain
+ // residual because a DATETIMEV2 slot cannot reveal source nanosecond
precision.
+ // ---- NaN differential on DOUBLE (total-ordering semantics) ----
+ // t_nan is created below: Doris defines NaN as equal to itself and
+ // greater than every finite value, but the pinned rust evaluator compares
+ // doubles with f64::partial_cmp (IEEE: NaN unordered, NaN != NaN), so the
+ // rust converter must NOT push DOUBLE predicates — a pushed `d > 1.0`
+ // would drop the stored NaN row Doris retains, and rows pruned by the
+ // rust filter cannot be recovered by the residual. The JNI path is
+ // unaffected: paimon-java's CompareUtils compares through
Double.compareTo,
+ // which matches Doris's total ordering.
+ spark_paimon_multi """
+ SET spark.sql.timestampType=TIMESTAMP_NTZ;
+ CREATE DATABASE IF NOT EXISTS paimon.${dbName};
+ DROP TABLE IF EXISTS paimon.${dbName}.t_eq_null;
+ CREATE TABLE paimon.${dbName}.t_eq_null (
+ a INT, b INT
+ ) USING paimon;
+ INSERT INTO paimon.${dbName}.t_eq_null VALUES (NULL, NULL), (1, 1),
(1, 2);
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_frac_ts;
+ CREATE TABLE paimon.${dbName}.t_frac_ts (
+ id INT, ts TIMESTAMP_NTZ
+ ) USING paimon;
+ INSERT INTO paimon.${dbName}.t_frac_ts VALUES
+ (1, TIMESTAMP '2024-01-01 00:00:00.123456'),
+ (2, TIMESTAMP '2024-01-01 00:00:00.123000'),
+ (3, TIMESTAMP '2024-01-02 00:00:00.999999');
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_frac_ts_dim;
+ CREATE TABLE paimon.${dbName}.t_frac_ts_dim (
+ id INT, ts TIMESTAMP_NTZ
+ ) USING paimon;
+ INSERT INTO paimon.${dbName}.t_frac_ts_dim VALUES
+ (1, TIMESTAMP '2024-01-01 00:00:00.123456'),
+ (2, TIMESTAMP '2024-01-02 00:00:00.999999');
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_narrowed;
+ CREATE TABLE paimon.${dbName}.t_narrowed (
+ id INT, v VARCHAR(10), amount DECIMAL(10,3)
+ ) USING paimon TBLPROPERTIES ('file.format' = 'parquet');
+ INSERT INTO paimon.${dbName}.t_narrowed VALUES (1, 'abcdef', 1.234),
(2, 'xyzuvw', 2.345);
+ DROP TABLE IF EXISTS paimon.${dbName}.t_narrowed_dim;
+ CREATE TABLE paimon.${dbName}.t_narrowed_dim (
+ v VARCHAR(3), amount DECIMAL(10,2)
+ ) USING paimon;
+ INSERT INTO paimon.${dbName}.t_narrowed_dim VALUES ('abc', 1.23);
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_nan;
+ CREATE TABLE paimon.${dbName}.t_nan (
+ id INT, d DOUBLE
+ ) USING paimon;
+ INSERT INTO paimon.${dbName}.t_nan VALUES
+ (1, 1.5), (2, CAST('NaN' AS DOUBLE)), (3, NULL);
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_dedup_ignore_del;
+ CREATE TABLE paimon.${dbName}.t_dedup_ignore_del (
+ id INT, v INT
+ ) USING paimon TBLPROPERTIES (
+ 'primary-key' = 'id',
+ 'merge-engine' = 'deduplicate',
+ 'deduplicate.ignore-delete' = 'true',
+ 'file.format' = 'parquet'
+ );
+ INSERT INTO paimon.${dbName}.t_dedup_ignore_del VALUES (1, 11), (2,
22);
+ DELETE FROM paimon.${dbName}.t_dedup_ignore_del WHERE id = 1;
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_pu_remove_record_del;
+ CREATE TABLE paimon.${dbName}.t_pu_remove_record_del (
+ id INT, v INT
+ ) USING paimon TBLPROPERTIES (
+ 'primary-key' = 'id',
+ 'merge-engine' = 'partial-update',
+ 'partial-update.remove-record-on-delete' = 'true',
+ 'file.format' = 'parquet'
+ );
+ INSERT INTO paimon.${dbName}.t_pu_remove_record_del VALUES (1, 11),
(2, 22);
+ DELETE FROM paimon.${dbName}.t_pu_remove_record_del WHERE id = 1;
+
+ DROP TABLE IF EXISTS paimon.${dbName}.t_nested_evo;
+ CREATE TABLE paimon.${dbName}.t_nested_evo (
+ id INT, s STRUCT<a: INT, b: STRING>
+ ) USING paimon TBLPROPERTIES (
+ 'bucket' = '-1',
+ 'file.format' = 'parquet'
+ );
+ INSERT INTO paimon.${dbName}.t_nested_evo VALUES (1, struct(10, 'x')),
(2, struct(20, 'y'));
+ ALTER TABLE paimon.${dbName}.t_nested_evo ADD COLUMN s.c INT;
+ INSERT INTO paimon.${dbName}.t_nested_evo VALUES (3, struct(30, 'z',
33));
+ """
+
+ // 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]
+ def originalRfWait = sql("select @@runtime_filter_wait_infinitely")[0][0]
+ def originalRfType = sql("select @@runtime_filter_type")[0][0]
+ def originalRfPrune = sql("select @@enable_runtime_filter_prune")[0][0]
+ def originalRfMaxIn = sql("select @@runtime_filter_max_in_num")[0][0]
+
+ try {
+ sql """switch ${catalogName}"""
+ sql """use ${dbName}"""
+ // Spark rejects narrowing in its analyzer. Doris forwards these
explicit
+ // casts to Paimon and publishes a new schema without rewriting old
files.
+ sql "ALTER TABLE t_narrowed MODIFY COLUMN v VARCHAR(3) NULL"
+ sql "ALTER TABLE t_narrowed MODIFY COLUMN amount DECIMAL(10,2) NULL"
+ sql """set enable_file_scanner_v2=true"""
+ // These tables are parquet append tables, whose DataSplits convert to
+ // raw native splits; without forcing, getSplits() would hand both legs
+ // to the native reader and bypass the JNI / rust converters entirely.
+ sql """set force_jni_scanner=true"""
+ sql "set enable_paimon_rust_reader=false"
+ // Verify the DDL really narrows historical values before testing
filter pushdown.
+ assertEquals([[1, "abc", "1.23"], [2, "xyz", "2.35"]].toString(),
+ sql("select id,v,cast(amount as string) from t_narrowed order
by id").toString())
+ // Profile capture for the reader-path verification below.
+ sql """set enable_profile=true"""
+ // The join leg must receive its IN runtime filter before the split
+ // opens, so the rust converter sees it in the conjuncts.
+ sql """set runtime_filter_wait_infinitely=true"""
+ // TRuntimeFilterType.IN == 1: force the IN runtime-filter shape.
+ sql """set runtime_filter_type=1"""
+
+ // Reuse the framework's profile readiness polling and configured HTTP
credentials.
+ def profileAction = new ProfileAction(context)
+ def profileTextOf = { String query ->
+ sql(query)
+ def queryId = sql("select last_query_id()")[0][0]
+ profileAction.getProfile(queryId.toString(), ["FileScannerV2"])
Review Comment:
[P2] Wait for the Rust child profile and counters before using this helper's
result. `getProfile` can return as soon as `FileScannerV2` appears, while
`PaimonRustReader` and its predicate/runtime-filter counters are still
unpublished. The positive checks below can fail on a valid Rust scan; absence
and zero-applied checks can pass on a partial profile. Require the child group
and expected counters for positive checks, and a completed profile for negative
checks.
##########
regression-test/suites/external_table_p0/paimon/test_paimon_rust_reader_compatibility.groovy:
##########
@@ -0,0 +1,266 @@
+// 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.
+
+import org.apache.doris.regression.action.ProfileAction
+
+suite("test_paimon_rust_reader_compatibility", "p0,external,paimon") {
+ if
(!"true".equalsIgnoreCase(context.config.otherConfigs.get("enablePaimonTest")))
{
+ return
+ }
+ def endpoint = context.config.otherConfigs.get("externalEnvIp")
+ def port = context.config.otherConfigs.get("iceberg_minio_port")
+ def catalog = "test_paimon_rust_compatibility"
+ def database = "test_paimon_rust_compatibility_db"
+ def settings = ["enable_paimon_rust_reader", "force_jni_scanner",
+ "enable_file_scanner_v2", "enable_profile",
"enable_prune_nested_column",
+ "enable_push_down_no_group_agg", "time_zone"]
+ def saved = settings.collectEntries { [(it): sql("select @@${it}")[0][0]] }
+ sql "DROP CATALOG IF EXISTS ${catalog}"
+ sql """CREATE CATALOG ${catalog} PROPERTIES (
+ 'type'='paimon', 'paimon.catalog.type'='filesystem',
+ 'warehouse'='s3://warehouse/wh',
's3.endpoint'='http://${endpoint}:${port}',
+ 's3.access_key'='admin', 's3.secret_key'='password',
+ 's3.region'='us-east-1', 'use_path_style'='true')"""
+ sql "SWITCH ${catalog}"
+ sql "DROP DATABASE IF EXISTS ${database} FORCE"
+ sql "CREATE DATABASE ${database}"
+ sql "USE ${database}"
+ try {
+ sql "set enable_profile=true"
+ sql "set time_zone='+00:00'"
+ sql "set enable_prune_nested_column=true"
+ sql "set enable_file_scanner_v2=true"
+ sql "set force_jni_scanner=true"
+ def profiles = new ProfileAction(context)
+ def check = { String query, List expected, boolean rustExpected ->
+ sql "set enable_paimon_rust_reader=false"
+ def baseline = sql(query)
+ assertEquals(expected.toString(), baseline.toString())
+ sql "set enable_paimon_rust_reader=true"
+ def actual = sql(query)
+ def queryId = sql("select last_query_id()")[0][0].toString()
+ def profile = profiles.getProfile(queryId, ["FileScannerV2"])
Review Comment:
[P2] Wait for the Rust child profile before asserting reader selection. This
poll asks `ProfileAction` only for `FileScannerV2`; when the profile has no
completion marker yet, it can return before `PaimonRustReader` appears. A valid
Rust scan may fail this assertion, and a fallback check may pass on an
incomplete snapshot. Require both markers for Rust-positive cases and a
completed profile for negative cases.
--
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]