Gabriel39 commented on code in PR #66227: URL: https://github.com/apache/doris/pull/66227#discussion_r4119333299
########## regression-test/suites/external_table_p0/paimon/test_paimon_rust_reader_eq_for_null.groovy: ########## @@ -0,0 +1,459 @@ +// 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. +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") + + // Table is created via Spark because Doris does not support Paimon DDL. + // 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. The rust predicate converter's + // sub-millisecond preservation (paimon_datum int_val2 / nanos) is covered + // by the PaimonRustPredicateConverterTest unit tests; this suite exercises + // the full equality and runtime-filter-join pushdown chains with the + // precision the plan can actually deliver. + // ---- 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_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 ( + 'primary-key' = 'id', + '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] + + try { + sql """switch ${catalogName}""" + sql """use ${dbName}""" + 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""" + // 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""" + + // 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", "") + if (code == 0 && out.contains("FileScannerV2")) { + return out + } + Thread.sleep(1000) + } + throw new Exception("profile not available for query ${queryId}") + } + + // The join must generate an IN runtime filter onto the probe ts. The + // LIMIT on the build side is required for the planner to assign a + // runtime filter to the paimon scan, and the RF lines only appear in + // the verbose explain. + def joinExplain = sql( + """explain verbose select p.id from t_frac_ts p join + (select ts from t_frac_ts_dim limit 10) d on p.ts = d.ts order by p.id""") + .flatten().join("\n") + assertTrue(joinExplain.contains("runtime filters") && joinExplain.contains("[in]"), + "the join must plan an IN runtime filter on the probe scan") + // The null-safe twin: the build side contains a NULL, so the probe's + // runtime filter is null-aware (EQ_FOR_NULL) — its residual execution + // restores NULL probes to true. + def nullAwareJoinExplain = sql( + """explain verbose select p.a, p.b from t_eq_null p join + (select a from t_eq_null limit 3) d on p.a <=> d.a + order by p.a nulls last, p.b nulls last""") + .flatten().join("\n") + assertTrue(nullAwareJoinExplain.contains("runtime filters") + && nullAwareJoinExplain.contains("[in]"), + "the null-safe join must plan a null-aware IN runtime filter on the probe scan") + + def testQueries = [ + // Column-to-column: the only form that reaches the BE as + // EQ_FOR_NULL. Must keep (NULL, NULL) and (1, 1), drop (1, 2). + """select * from t_eq_null where a <=> b order by a nulls last, b nulls last""", + // FE-rewritten forms; also exercise the equality / IS NULL + // pushdown paths of the rust predicate converter. + """select * from t_eq_null where a <=> 1 order by a, b""", + """select * from t_eq_null where a <=> NULL order by a, b""", + // Fractional TIMESTAMP(6) equality. The literal is deliberately + // 3-digit: FE truncates plan-time pushed-down timestamp + // predicates to milliseconds, so a 6-digit literal would reach + // the readers truncated and the exact residual would then drop + // every row — for the JNI, rust and native readers alike. The + // table data keeps 6-digit values (see the join below), and + // the rust converter's sub-millisecond handling is covered by + // the PaimonRustPredicateConverterTest unit tests. + """select id from t_frac_ts where ts = '2024-01-01 00:00:00.123' order by id""", + // The join form exercises the timestamp conversion through a + // runtime-filter IN predicate on the probe scan (t_frac_ts): + // runtime filters are built at runtime from the build side's + // actual values, so they bypass the plan-time millisecond + // truncation and carry the full 6-digit precision through the + // rust converter — both dim values must match their probe + // rows. The LIMIT on the build side is what makes the planner + // assign the runtime filter to the paimon scan (see the + // explain check above); runtime_filter_wait_infinitely + // guarantees the filter has arrived before the split opens. + """select p.id from t_frac_ts p join (select ts from t_frac_ts_dim limit 10) d + on p.ts = d.ts order by p.id""", + // Null-safe join with NULLs on both sides: the arrived runtime + // filter is null-aware and its residual execution restores + // the NULL probes to true, so (null, null) must survive. The + // rust pushdown must not unwrap the wrapper into the ordinary + // IN set — the rebuilt set carries only the concrete member + // (1) and would prune the NULL probe before the join sees it + // (BE unit test: NullAwareRuntimeFilterStaysResidual). + """select p.a, p.b from t_eq_null p join (select a from t_eq_null limit 3) d + on p.a <=> d.a order by p.a nulls last, p.b nulls last""", + // NaN total-ordering differential (see the t_nan setup): NaN + // is greater than every finite value, so `d > 1.0` keeps the + // NaN row and `d < 2.0` does not; `d = 'NaN'` matches it. + """select id from t_nan where d > 1.0 order by id""", + """select id from t_nan where d < 2.0 order by id""", + """select id from t_nan where d = cast('NaN' as double) order by id""", + // deduplicate.ignore-delete=true: the DELETE of (1, 11) writes a + // retract record into a new, uncompacted file, and Java's + // DeduplicateMergeFunction skips it — the row must survive. The + // pinned rust deduplicate merge has no option channel and would + // pick the retract as the latest row, silently dropping the key, + // so the FE gate keeps this table on JNI (verified through the + // profile below). + """select id, v from t_dedup_ignore_del order by id""", + // partial-update.remove-record-on-delete (non-DV): Java honors + // it — the delete removes the whole (1, 11) record — but the + // pinned rust PartialUpdateConfig read validation returns + // Unsupported for the key, so the FE gate keeps the table on + // JNI (verified through the profile below). + """select id, v from t_pu_remove_record_del order by id""", + // Nested schema evolution: ALTER ADD COLUMN s.c landed after + // rows 1-2 were written, so their files predate the child. The + // paimon-rust reader reconciles nested children by field id and + // NULL-fills the added child — the same semantics as Java's + // SchemaEvolutionUtil — since paimon-rust 381a1ad + // "fix(read): null-fill nested fields a data file predates + // (#805)", first included in the baac87c pin (the previous + // cabdeb9 pin cast the whole StructArray through arrow-cast and + // failed; reproduced live before the upgrade). + // The rust leg must actually run the rust reader (profile + // below): this differential is the capability guard for the + // crate upgrade. + """select id, s from t_nested_evo order by id""" + ] + def expectedResults = [ + [[1, 1], [null, null]], + [[1, 1], [1, 2]], + [[null, null]], + [[2]], + [[1], [3]], + [[1, 1], [1, 1], [1, 2], [1, 2], [null, null]], + [[1], [2]], + [[1]], + [[2]], + [[1, 11], [2, 22]], + [[2, 22]], + [[1, '{"a":10, "b":"x", "c":null}'], + [2, '{"a":20, "b":"y", "c":null}'], + [3, '{"a":30, "b":"z", "c":33}']] + ] + // Representative converter query reused for the reader-path checks. + String pushdownQuery = testQueries[3] + + sql """set enable_paimon_rust_reader=false""" + def jniResults = testQueries.collect { query -> sql(query) } + // The JNI leg must ride the logical-split JNI reader: the profile of a + // representative query must not contain the rust reader's timer. + def jniProfile = profileTextOf(pushdownQuery) + assertFalse(jniProfile.contains("PaimonRustReader"), "JNI leg must not use the rust reader") + + sql """set enable_paimon_rust_reader=true""" + def rustResults = testQueries.collect { query -> sql(query) } + // The rust leg must actually run the rust reader: its profile carries + // the PaimonRustReader timer group, which only the rust reader creates. + def rustProfile = profileTextOf(pushdownQuery) + assertTrue(rustProfile.contains("PaimonRustReader"), Review Comment: Addressed in 6009fdb795. Added input, converted, and successfully applied Rust predicate profile counters, plus separate arrived/applied runtime-filter counts. Applied counters advance only after paimon_read_builder_with_filter succeeds. The regression uses integer equality/IN positive controls and unsupported DOUBLE/timestamp negative controls; converter tests also cover counter reset and residual exclusion. ########## be/src/format_v2/table/paimon_rust_predicate_converter.cpp: ########## @@ -0,0 +1,832 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "format_v2/table/paimon_rust_predicate_converter.h" + +#include <algorithm> +#include <cctype> +#include <memory> +#include <utility> + +#include "common/logging.h" +#include "core/column/column_const.h" +#include "core/column/column_nullable.h" +#include "core/data_type/data_type.h" +#include "core/data_type/data_type_nullable.h" +#include "core/field.h" +#include "core/types.h" +#include "core/value/decimalv2_value.h" +#include "core/value/timestamptz_value.h" +#include "core/value/vdatetime_value.h" +#include "exprs/runtime_filter_expr.h" +#include "exprs/vcompound_pred.h" +#include "exprs/vdirect_in_predicate.h" +#include "exprs/vectorized_fn_call.h" +#include "exprs/vexpr.h" +#include "exprs/vin_predicate.h" +#include "exprs/vliteral.h" +#include "exprs/vslot_ref.h" + +namespace doris { + +namespace { +// paimon_datum tags (see paimon.h / bindings/c/src/table.rs::datum_from_c). +constexpr int32_t kTagBool = 0; +constexpr int32_t kTagTinyInt = 1; +constexpr int32_t kTagSmallInt = 2; +constexpr int32_t kTagInt = 3; +constexpr int32_t kTagLong = 4; +constexpr int32_t kTagDouble = 6; +constexpr int32_t kTagString = 7; +constexpr int32_t kTagDate = 8; +constexpr int32_t kTagTimestamp = 10; +constexpr int32_t kTagDecimal = 12; +constexpr int32_t kTagBytes = 13; + +// paimon decimal precision ceiling (paimon::Decimal::MAX_PRECISION). +constexpr int32_t kPaimonDecimalMaxPrecision = 38; + +// RAII for an owned paimon_predicate*. and/or/not consume their inputs, so we +// release() before handing pointers to them. +struct predicate_deleter { + void operator()(paimon_predicate* p) const { + if (p) { + paimon_predicate_free(p); + } + } +}; +using predicate_ptr = std::unique_ptr<paimon_predicate, predicate_deleter>; + +// RAII for an owned paimon_error*. +struct error_deleter { + void operator()(paimon_error* p) const { + if (p) { + paimon_error_free(p); + } + } +}; +using error_ptr = std::unique_ptr<paimon_error, error_deleter>; + +// Render a paimon_error into a string. Takes ownership of `err` via RAII so it +// is freed on every return path. Safe to call with nullptr. +std::string consume_predicate_error(paimon_error* err) { + error_ptr owned(err); + if (!owned) { + return "unknown error"; + } + std::string msg; + if (owned->message.data != nullptr && owned->message.len > 0) { + msg.assign(reinterpret_cast<const char*>(owned->message.data), owned->message.len); + } + return "code=" + std::to_string(owned->code) + ", msg=" + msg; +} +} // namespace + +PaimonRustPredicateConverter::PaimonRustPredicateConverter( + const std::vector<std::string>& column_names, const std::vector<DataTypePtr>& column_types, + const paimon_table* table) + : _table(table) { + DORIS_CHECK(column_names.size() == column_types.size()); + _columns_by_name.reserve(column_names.size()); + for (size_t i = 0; i < column_names.size(); ++i) { + _columns_by_name.emplace(_normalize_name(column_names[i]), + std::make_pair(column_names[i], column_types[i])); + } + // Paimon TIMESTAMP (wall clock) is stored as epoch-millis-of-the-wall-time + // and the DateTimeV2 serde decodes timezone-naive arrow values in UTC, so + // timestamp literals convert wall->epoch in UTC. utc_time_zone() needs no + // tzdata lookup, so the conversion cannot silently fall back to a + // machine-local zone. + _utc_tz = cctz::utc_time_zone(); +} + +paimon_predicate* PaimonRustPredicateConverter::build(const VExprContextSPtrs& conjuncts) { + if (_table == nullptr) { + return nullptr; + } + predicate_ptr result; + for (const auto& conjunct : conjuncts) { + if (!conjunct || !conjunct->root()) { + continue; + } + auto root = conjunct->root(); + if (root->is_rf_wrapper()) { + if (auto impl = root->get_impl()) { + // A null-aware runtime filter (an EQ_FOR_NULL join) must stay + // residual: its wrapper execution restores NULL probe rows to + // true (RuntimeFilterExpr::change_null_to_true), while the + // unwrapped impl — rebuilt as an ordinary IN predicate through + // VDirectInPredicate::get_slot_in_expr — treats NULL as + // not-in-set and would prune the NULL probes permanently + // before the join sees them. Keep the wrapper itself: it fails + // every dispatch below, so the conjunct stays in the residual, + // and it is safe to execute on selected rows, so later + // conjuncts keep pushing. The lance pushdown declines + // is_null_aware() filters for the same reason. + // is_null_aware() is concrete on RuntimeFilterExpr (the only + // class whose is_rf_wrapper() is true), so the dynamic_cast + // never fails in practice; the null guard keeps the unwrap for + // any future wrapper shape. + auto* rf_wrapper = dynamic_cast<RuntimeFilterExpr*>(root.get()); + if (rf_wrapper == nullptr || !rf_wrapper->is_null_aware()) { + root = impl; + } + } + } + // Preserve a safe prefix of the conjunct order: a later pushed + // predicate (e.g. an arrived IN runtime filter) could otherwise prune + // rows on which an earlier error-preserving conjunct — + // assert_true(...), a failing cast, ... — must still raise. The v1 + // partition-pruning path (FileScanner::_init_runtime_filter_partition_ + // prune_ctxs) stops at is_safe_to_execute_on_selected_rows() for the + // same reason, so a convertible predicate after an unsafe conjunct + // must not be pushed. Safe conjuncts that cannot be converted keep + // the old skip: they cannot raise, so pruning rows before they are + // evaluated as the residual never loses an error. + if (!_is_safe_to_push(root)) { + break; + } + predicate_ptr pred(_convert_expr(root)); + if (!pred) { + continue; + } + if (!result) { + result = std::move(pred); + } else { + // and consumes both inputs regardless of success. + result.reset(paimon_predicate_and(result.release(), pred.release())); + if (!result) { + return nullptr; + } + } + } + return result.release(); +} + +bool PaimonRustPredicateConverter::_is_safe_to_push(const VExprSPtr& expr) { + if (expr->is_safe_to_execute_on_selected_rows()) { + return true; + } + // VectorizedFnCall::is_safe_to_execute_on_selected_rows() admits a fixed + // whitelist that does not include `like`, so a like conjunct would always + // stop the safe prefix here and never reach _convert_like. A like call + // whose children are themselves safe is just as total as the whitelisted + // comparisons: it only compares strings, and its regex is built by escaping + // the pattern operand (FunctionLike::convert_like_pattern), so no scanned + // value can make it raise — a malformed pattern or escape operand fails + // at open() on every execution, regardless of which rows survive. Admit + // it here rather than widening the generic whitelist, which gates shared + // BE pushdown paths outside this reader's scope. + auto* fn = dynamic_cast<VectorizedFnCall*>(expr.get()); + if (fn == nullptr || _normalize_name(fn->function_name()) != "like") { + return false; + } + for (uint16_t i = 0; i < fn->get_num_children(); ++i) { + if (!fn->get_child(i)->is_safe_to_execute_on_selected_rows()) { + return false; + } + } + return true; +} + +paimon_predicate* PaimonRustPredicateConverter::_convert_expr(const VExprSPtr& expr) { + if (!expr) { + return nullptr; + } + + // Casts are not unwrapped anywhere (predicate root included): a cast node + // fails every dispatch below and the conjunct stays in the Doris residual, + // mirroring the FE converter, which keeps casted expressions unconverted. + if (auto* direct_in = dynamic_cast<VDirectInPredicate*>(expr.get())) { + VExprSPtr in_expr; + if (direct_in->get_slot_in_expr(in_expr)) { + return _convert_in(in_expr); + } + return nullptr; + } + + if (dynamic_cast<VInPredicate*>(expr.get()) != nullptr) { + return _convert_in(expr); + } + + switch (expr->op()) { + case TExprOpcode::COMPOUND_AND: + case TExprOpcode::COMPOUND_OR: + return _convert_compound(expr); + case TExprOpcode::COMPOUND_NOT: + return nullptr; + case TExprOpcode::EQ: + case TExprOpcode::EQ_FOR_NULL: + case TExprOpcode::NE: + case TExprOpcode::GE: + case TExprOpcode::GT: + case TExprOpcode::LE: + case TExprOpcode::LT: + return _convert_binary(expr); + default: + break; + } + + if (auto* fn = dynamic_cast<VectorizedFnCall*>(expr.get())) { + auto fn_name = _normalize_name(fn->function_name()); + if (fn_name == "is_null_pred" || fn_name == "is_not_null_pred") { + return _convert_is_null(expr, fn_name); + } + if (fn_name == "like") { + return _convert_like(expr); + } + } + + return nullptr; +} + +paimon_predicate* PaimonRustPredicateConverter::_convert_compound(const VExprSPtr& expr) { + if (!expr || expr->get_num_children() != 2) { + return nullptr; + } + predicate_ptr left(_convert_expr(expr->get_child(0))); + if (!left) { + return nullptr; + } + predicate_ptr right(_convert_expr(expr->get_child(1))); + if (!right) { + return nullptr; + } + + if (expr->op() == TExprOpcode::COMPOUND_AND) { + return paimon_predicate_and(left.release(), right.release()); + } + if (expr->op() == TExprOpcode::COMPOUND_OR) { + return paimon_predicate_or(left.release(), right.release()); + } + return nullptr; +} + +paimon_predicate* PaimonRustPredicateConverter::_convert_in(const VExprSPtr& expr) { + auto* in_pred = dynamic_cast<VInPredicate*>(expr.get()); + if (!in_pred || expr->get_num_children() < 2) { + return nullptr; + } + auto field_meta = _resolve_field(expr->get_child(0)); + if (!field_meta) { + return nullptr; + } + + const auto num_values = expr->get_num_children() - 1; + // Reserve up front so the backing strings never reallocate: each datum's + // str_data points into storages[i], which must stay stable. + std::vector<std::string> storages; + std::vector<paimon_datum> datums; + storages.reserve(num_values); + datums.reserve(num_values); + for (uint16_t i = 1; i < expr->get_num_children(); ++i) { + // Casted list values are rejected by _convert_literal (the same rule + // as the binary RHS): with debug_skip_fold_constant the cast reaches + // the BE un-folded, and unwrapping it would filter on the pre-cast + // value — in `amount IN (CAST(1.24 AS DECIMAL(10,1)))` Doris keeps + // the 1.2 rows while the unwrapped 1.24 push removes them. Rejecting + // the value rejects the whole predicate; the residual applies the + // cast correctly. + auto holder = _convert_literal(expr->get_child(i), field_meta->type); + if (!holder) { + return nullptr; + } + storages.emplace_back(std::move(holder->storage)); + paimon_datum datum = holder->datum; + _bind_datum_storage(&datum, storages.back()); + datums.emplace_back(datum); + } + + if (datums.empty()) { + return nullptr; + } + if (in_pred->is_not_in()) { + return _take(paimon_predicate_is_not_in(_table, field_meta->column.c_str(), datums.data(), + datums.size())); + } + return _take(paimon_predicate_is_in(_table, field_meta->column.c_str(), datums.data(), + datums.size())); +} + +paimon_predicate* PaimonRustPredicateConverter::_convert_binary(const VExprSPtr& expr) { + if (!expr || expr->get_num_children() != 2) { + return nullptr; + } + auto field_meta = _resolve_field(expr->get_child(0)); + if (!field_meta) { + return nullptr; + } + const char* column = field_meta->column.c_str(); + + // Convert the RHS first so EQ_FOR_NULL (<=>) only converts when the RHS is + // a convertible literal, mirroring the FE converter, which rejects a + // non-literal RHS. A column-to-column `a <=> b` must therefore stay in the + // Doris residual: it has no single-column rust predicate, and pushing + // `a IS NULL` would wrongly discard rows like (1, 1) — rows dropped by the + // pushed filter cannot be recovered by the residual conjunct. + auto holder = _convert_literal(expr->get_child(1), field_meta->type); + if (!holder) { + return nullptr; + } + + if (expr->op() == TExprOpcode::EQ_FOR_NULL) { + return _take(paimon_predicate_is_null(_table, column)); + } + + // `holder` is a local, so its storage stays put for the duration of the call. + _bind_datum_storage(&holder->datum, holder->storage); + const paimon_datum& datum = holder->datum; + + switch (expr->op()) { + case TExprOpcode::EQ: + return _take(paimon_predicate_equal(_table, column, datum)); + case TExprOpcode::NE: + return _take(paimon_predicate_not_equal(_table, column, datum)); + case TExprOpcode::GE: + return _take(paimon_predicate_greater_or_equal(_table, column, datum)); + case TExprOpcode::GT: + return _take(paimon_predicate_greater_than(_table, column, datum)); + case TExprOpcode::LE: + return _take(paimon_predicate_less_or_equal(_table, column, datum)); + case TExprOpcode::LT: + return _take(paimon_predicate_less_than(_table, column, datum)); + default: + break; + } + return nullptr; +} + +paimon_predicate* PaimonRustPredicateConverter::_convert_is_null(const VExprSPtr& expr, + const std::string& fn_name) { + if (!expr || expr->get_num_children() != 1) { + return nullptr; + } + auto field_meta = _resolve_field(expr->get_child(0)); + if (!field_meta) { + return nullptr; + } + if (fn_name == "is_not_null_pred") { + return _take(paimon_predicate_is_not_null(_table, field_meta->column.c_str())); + } + return _take(paimon_predicate_is_null(_table, field_meta->column.c_str())); +} + +paimon_predicate* PaimonRustPredicateConverter::_convert_like(const VExprSPtr& expr) { + if (!expr || expr->get_num_children() < 2) { + return nullptr; + } + auto field_meta = _resolve_field(expr->get_child(0)); + if (!field_meta || !_is_string_type(field_meta->type->get_primitive_type())) { + return nullptr; + } + + auto pattern_opt = _extract_string_literal(expr->get_child(1)); + if (!pattern_opt) { + return nullptr; + } + const std::string& pattern = *pattern_opt; + + // Doris's 3-arg like(col, pattern, escape) form: the pinned rust + // predicate only supports the backslash escape (its builder rejects any + // other escape character), so only the default-escape shape converts. + // The 2-arg form always carries the Doris/SQL default '\'. + if (expr->get_num_children() >= 3) { + auto escape_opt = _extract_string_literal(expr->get_child(2)); + if (!escape_opt || escape_opt->size() != 1 || (*escape_opt)[0] != '\\') { + return nullptr; + } + } + + // The pinned rust like implements SQL LIKE (% = any run, _ = one char, + // \X = literal X for any X — arrow's like kernel), and its builder + // optimizes `prefix%` / `%suffix` / `%mid%` shapes internally, so every + // pattern without a backslash matches Doris semantics exactly. The one + // divergence is the escape handling for a backslash before an ordinary + // character: Doris keeps both characters (only \%, \_ and \\ are + // escapes; `a\qb` matches a backslash followed by q b) while rust + // consumes the backslash (`a\qb` matches aqb) — pushing such a pattern + // would wrongly prune the Doris-matching rows before the residual can + // see them. Reject exactly the divergent shapes: every backslash must + // precede %, _ or \ (identical literal semantics on both sides) or end + // the pattern (literal backslash on both sides); anything else stays in + // the residual. + for (size_t i = 0; i < pattern.size(); ++i) { + if (pattern[i] == '\\') { + if (i + 1 >= pattern.size()) { + continue; // trailing backslash: literal on both sides + } + char next = pattern[i + 1]; + if (next != '%' && next != '_' && next != '\\') { + return nullptr; + } + ++i; + } + } + + paimon_datum datum {}; + datum.tag = kTagString; + // `pattern` outlives the build call (it is a const ref into the caller's + // optional), so the datum may point straight at it. + _bind_datum_storage(&datum, pattern); + return _take(paimon_predicate_like(_table, field_meta->column.c_str(), datum, '\\')); +} + +std::optional<PaimonRustPredicateConverter::FieldMeta> PaimonRustPredicateConverter::_resolve_field( + const VExprSPtr& expr) const { + if (!expr) { + return std::nullopt; + } + // Mirror the FE converter's convertDorisExprToSlotRef: a casted column is + // rejected, never unwrapped. Stripping a lossy cast changes which rows match + // — for a DECIMAL(10,2) column, CAST(amount AS DECIMAL(10,1)) = 1.2 keeps + // the row 1.24 while the unwrapped `amount = 1.2` prunes it — and rows + // pruned by the pushed filter cannot be recovered by the Doris residual. + // The conjunct stays in the residual instead. + auto* slot_ref = dynamic_cast<VSlotRef*>(expr.get()); + if (!slot_ref) { + return std::nullopt; + } + // FileScannerV2 rewrites conjunct VSlotRefs to table global indices, so slot_id + // is a position, not a slot id; resolve by the carried column name against the + // projected-column registry instead of the desc table. + auto it = _columns_by_name.find(_normalize_name(slot_ref->column_name())); + if (it == _columns_by_name.end()) { + return std::nullopt; + } + const auto& [column, type] = it->second; + if (!_is_supported_slot_type(type->get_primitive_type(), type->get_precision())) { + return std::nullopt; + } + return FieldMeta {column, type}; +} + +std::optional<PaimonRustPredicateConverter::DatumHolder> +PaimonRustPredicateConverter::_convert_literal(const VExprSPtr& expr, + const DataTypePtr& column_type) const { + // A casted literal is rejected, never unwrapped: the cast is not executed + // here (with constant folding disabled — debug_skip_fold_constant — it + // reaches the BE un-folded), so unwrapping would push the pre-cast value + // while the Doris residual compares against the cast result. A + // scale-reducing cast makes the two disagree — `amount = + // CAST(1.24 AS DECIMAL(10,1))` keeps a stored 1.20 row in Doris (1.24 + // rounds down to 1.2) but the unwrapped `amount = 1.24` push prunes it, and + // rows pruned by the pushed filter cannot be recovered by the residual. + // The conjunct stays in the residual, where the cast is evaluated exactly. + if (expr->node_type() == TExprNodeType::CAST_EXPR) { + return std::nullopt; + } + auto* literal = dynamic_cast<VLiteral*>(expr.get()); + if (!literal) { + return std::nullopt; + } + + auto literal_type = remove_nullable(literal->get_data_type()); + PrimitiveType literal_primitive = literal_type->get_primitive_type(); + PrimitiveType slot_primitive = column_type->get_primitive_type(); + + ColumnPtr col = literal->get_column_ptr()->convert_to_full_column_if_const(); + if (const auto* nullable = check_and_get_column<ColumnNullable>(*col)) { + if (nullable->is_null_at(0)) { + return std::nullopt; + } + col = nullable->get_nested_column_ptr(); + } + + Field field; + col->get(0, field); + + DatumHolder holder; + paimon_datum& datum = holder.datum; + + switch (slot_primitive) { + case TYPE_BOOLEAN: { + if (literal_primitive != TYPE_BOOLEAN) { + return std::nullopt; + } + datum.tag = kTagBool; + datum.int_val = static_cast<bool>(field.get<TYPE_BOOLEAN>()) ? 1 : 0; + return holder; + } + case TYPE_TINYINT: + case TYPE_SMALLINT: + case TYPE_INT: + case TYPE_BIGINT: { + if (!_is_integer_type(literal_primitive)) { + return std::nullopt; + } + int64_t value = 0; + switch (literal_primitive) { + case TYPE_TINYINT: + value = field.get<TYPE_TINYINT>(); + break; + case TYPE_SMALLINT: + value = field.get<TYPE_SMALLINT>(); + break; + case TYPE_INT: + value = field.get<TYPE_INT>(); + break; + case TYPE_BIGINT: + value = field.get<TYPE_BIGINT>(); + break; + default: + return std::nullopt; + } + datum.int_val = value; + switch (slot_primitive) { + case TYPE_TINYINT: + datum.tag = kTagTinyInt; + break; + case TYPE_SMALLINT: + datum.tag = kTagSmallInt; + break; + case TYPE_INT: + datum.tag = kTagInt; + break; + default: + datum.tag = kTagLong; + break; + } + return holder; + } + case TYPE_DOUBLE: { + if (literal_primitive != TYPE_DOUBLE && literal_primitive != TYPE_FLOAT) { + return std::nullopt; + } + datum.tag = kTagDouble; + datum.double_val = literal_primitive == TYPE_FLOAT + ? static_cast<double>(field.get<TYPE_FLOAT>()) + : field.get<TYPE_DOUBLE>(); + return holder; + } + case TYPE_DATE: + case TYPE_DATEV2: { + if (!_is_date_type(literal_primitive)) { + return std::nullopt; + } + int64_t seconds = 0; + if (literal_primitive == TYPE_DATE) { + const auto& dt = field.get<TYPE_DATE>(); + if (!dt.is_valid_date()) { + return std::nullopt; + } + dt.unix_timestamp(&seconds, _utc_tz); + } else { + const auto& dt = field.get<TYPE_DATEV2>(); + if (!dt.is_valid_date()) { + return std::nullopt; + } + dt.unix_timestamp(&seconds, _utc_tz); + } + datum.tag = kTagDate; + datum.int_val = _seconds_to_days(seconds); + return holder; + } + case TYPE_DATETIME: + case TYPE_DATETIMEV2: { + if (!_is_datetime_type(literal_primitive)) { + return std::nullopt; + } + datum.tag = kTagTimestamp; + if (literal_primitive == TYPE_DATETIME) { + const auto& dt = field.get<TYPE_DATETIME>(); + if (!dt.is_valid_date()) { + return std::nullopt; + } + int64_t seconds = 0; + dt.unix_timestamp(&seconds, _utc_tz); + // No sub-second part in datetime (v1): millis only, nanos stays 0. + datum.int_val = seconds * 1000; + } else { + const auto& dt = field.get<TYPE_DATETIMEV2>(); + if (!dt.is_valid_date()) { + return std::nullopt; + } + // ts is (seconds since epoch, the microsecond-of-second part); + // split it into the rust timestamp's (millis, nanos): truncating + // the sub-millisecond remainder would make an equality/IN predicate + // more selective than the original conjunct (e.g. .123456 pushed as + // .123000), wrongly discarding matching rows. DATETIMEV2 carries at + // most microseconds, and micros % 1000 * 1000 <= 999000 fits the + // nanos field, so the value is always representable exactly. + std::pair<int64_t, int64_t> ts; + dt.unix_timestamp(&ts, _utc_tz); + datum.int_val = ts.first * 1000 + ts.second / 1000; + datum.int_val2 = (ts.second % 1000) * 1000; + } + return holder; + } + case TYPE_VARCHAR: + case TYPE_STRING: { + if (!_is_string_type(literal_primitive)) { + return std::nullopt; + } + const auto& value = field.get<TYPE_STRING>(); + datum.tag = kTagString; + holder.storage.assign(value.data(), value.size()); + return holder; + } + case TYPE_DECIMALV2: + case TYPE_DECIMAL32: + case TYPE_DECIMAL64: + case TYPE_DECIMAL128I: + case TYPE_DECIMAL256: { + if (!_is_decimal_type(literal_primitive)) { + return std::nullopt; + } + int32_t precision = static_cast<int32_t>(literal_type->get_precision()); + int32_t scale = static_cast<int32_t>(literal_type->get_scale()); + if (precision <= 0 || precision > kPaimonDecimalMaxPrecision) { + return std::nullopt; + } + + __int128 value = 0; + switch (literal_primitive) { + case TYPE_DECIMALV2: + value = field.get<TYPE_DECIMALV2>().value(); + break; + case TYPE_DECIMAL32: + value = field.get<TYPE_DECIMAL32>().value; + break; + case TYPE_DECIMAL64: + value = field.get<TYPE_DECIMAL64>().value; + break; + case TYPE_DECIMAL128I: + value = field.get<TYPE_DECIMAL128I>().value; + break; + default: + return std::nullopt; + } + datum.tag = kTagDecimal; + // rust reassembles as ((int_val2 as i128) << 64) | (int_val as u64 as i128): + // int_val holds the low 64 bits, int_val2 the (sign-extended) high 64 bits. + datum.int_val = static_cast<int64_t>(static_cast<uint64_t>(value)); + datum.int_val2 = static_cast<int64_t>(value >> 64); + datum.uint_val = static_cast<uint32_t>(precision); + datum.uint_val2 = static_cast<uint32_t>(scale); + return holder; + } + default: + break; + } + return std::nullopt; +} + +std::optional<std::string> PaimonRustPredicateConverter::_extract_string_literal( + const VExprSPtr& expr) const { + // Same rule as _convert_literal: a bare string literal converts, a casted + // one is rejected — the cast is not executed here, so the pushed pattern + // must come from the literal's own value. A truncating cast (e.g. to + // CHAR(k)) would otherwise push a different prefix than the residual + // matches. Cast trees stay in the residual. + if (expr->node_type() == TExprNodeType::CAST_EXPR) { + return std::nullopt; + } + auto* literal = dynamic_cast<VLiteral*>(expr.get()); + if (!literal) { + return std::nullopt; + } + auto literal_type = remove_nullable(literal->get_data_type()); + if (!_is_string_type(literal_type->get_primitive_type())) { + return std::nullopt; + } + + ColumnPtr col = literal->get_column_ptr()->convert_to_full_column_if_const(); + if (const auto* nullable = check_and_get_column<ColumnNullable>(*col)) { + if (nullable->is_null_at(0)) { + return std::nullopt; + } + col = nullable->get_nested_column_ptr(); + } + Field field; + col->get(0, field); + const auto& value = field.get<TYPE_STRING>(); + return std::string(value.data(), value.size()); +} + +paimon_predicate* PaimonRustPredicateConverter::_take(paimon_result_predicate result) { + if (result.error != nullptr) { + // A single leaf failing to build is not fatal: that conjunct is simply + // dropped from the pushed-down filter and the engine still re-applies it. + // Log at WARNING so it is visible without verbose logging enabled. + LOG(WARNING) << "paimon-rust build predicate failed: " + << consume_predicate_error(result.error); + return nullptr; + } + return result.predicate; +} + +void PaimonRustPredicateConverter::_bind_datum_storage(paimon_datum* datum, + const std::string& storage) { + if (datum->tag == kTagString || datum->tag == kTagBytes) { + datum->str_data = reinterpret_cast<const uint8_t*>(storage.data()); + datum->str_len = storage.size(); + } +} + +std::string PaimonRustPredicateConverter::_normalize_name(std::string_view name) { + std::string out(name); + std::transform(out.begin(), out.end(), out.begin(), + [](unsigned char c) { return static_cast<char>(std::tolower(c)); }); + return out; +} + +int32_t PaimonRustPredicateConverter::_seconds_to_days(int64_t seconds) { + static constexpr int64_t kSecondsPerDay = 24 * 60 * 60; + int64_t days = seconds / kSecondsPerDay; + if (seconds < 0 && seconds % kSecondsPerDay != 0) { + --days; + } + return static_cast<int32_t>(days); +} + +bool PaimonRustPredicateConverter::_is_integer_type(PrimitiveType type) { + switch (type) { + case TYPE_TINYINT: + case TYPE_SMALLINT: + case TYPE_INT: + case TYPE_BIGINT: + return true; + default: + return false; + } +} + +bool PaimonRustPredicateConverter::_is_string_type(PrimitiveType type) { + return type == TYPE_CHAR || type == TYPE_VARCHAR || type == TYPE_STRING; +} + +bool PaimonRustPredicateConverter::_is_decimal_type(PrimitiveType type) { + switch (type) { + case TYPE_DECIMALV2: + case TYPE_DECIMAL32: + case TYPE_DECIMAL64: + case TYPE_DECIMAL128I: + case TYPE_DECIMAL256: + return true; + default: + return false; + } +} + +bool PaimonRustPredicateConverter::_is_date_type(PrimitiveType type) { + return type == TYPE_DATE || type == TYPE_DATEV2; +} + +bool PaimonRustPredicateConverter::_is_datetime_type(PrimitiveType type) { + return type == TYPE_DATETIME || type == TYPE_DATETIMEV2; +} + +bool PaimonRustPredicateConverter::_is_supported_slot_type(PrimitiveType type, uint32_t precision) { + switch (type) { + case TYPE_BOOLEAN: + case TYPE_TINYINT: + case TYPE_SMALLINT: + case TYPE_INT: + case TYPE_BIGINT: + case TYPE_VARCHAR: + case TYPE_STRING: + case TYPE_DATE: + case TYPE_DATEV2: + case TYPE_DATETIME: + case TYPE_DATETIMEV2: + return true; + case TYPE_DECIMALV2: Review Comment: Addressed in 6009fdb795. Decimal predicates and arrived IN runtime filters now remain residual because current slot metadata cannot establish historical file scale equivalence. Added converter coverage and a DECIMAL(10,3) to DECIMAL(10,2) historical-file join differential. The fixture performs the narrowing through Doris Paimon ALTER since Spark rejects it in its analyzer. -- 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]
