zhuxiangyi opened a new pull request, #9980: URL: https://github.com/apache/paimon/pull/9980
# [spark][core] Preserve SQL semantics for CAST predicate pushdown ## Purpose The Spark converter currently accepts any CAST that the storage CastExecutors can resolve. Those casts do not carry Spark's ANSI mode or session time zone and are not SQL-semantically interchangeable. Partition-only predicates are consumed by the connector, so differences can cause wrong results, a suppressed ANSI overflow, or a NULL-related planning exception. Examples: `CAST(p AS BIGINT) IS NULL` throws on a NULL partition; `CAST(128 AS TINYINT)=-128` returns a row instead of throwing in ANSI mode; casting the partition string `'1970'` to DATE drops a matching row; TIMESTAMP-to-STRING differs in formatting/time zone. ## Changes - Propagate NULL in CastTransform without changing non-NULL storage conversion rules. - Gate Spark CAST conversion to an explicit safe subset: integral widening, integral/boolean to unbounded STRING, and eligible identities. - Reject narrowing, parsed strings, decimal/temporal cross-casts and floating-result casts. Keep their predicates in Spark. - Check exposed Spark type identities carefully: decline sub-microsecond timestamp identities and non-identity storage/exposed representation mismatches such as TIME exposed as INT. - Do not modify global CastExecutors, floating comparators, or literal/nested-field handling. Floating-result casts are deliberately conservative: matching scalar conversion values is not sufficient when storage comparison distinguishes signed zeros and Spark comparison does not. A separate comparator change is outside scope. ## Validation - Original four bugs reproduced on Spark 3.5.8 before changes. - Red/green Common NULL cases and SQL/converter matrix; high-precision timestamp identity and floating-result guard regressions added after review. - **310 passing executions on Spark 3.5.8 / Scala 2.12.18** and **310 on Spark 4.1.2 / Scala 2.13.17**, including Common208, Spark Java7, new CAST suite32 and related PushDown/Converter suites63. Shared suites are repeated across profiles, not 620 unique tests. - Normal non-fast-build checks used for final runs. Actual JDK8 javac compilation of changed Common Java and tests succeeded; runtime validation used JDK17. - Parquet/ORC × partition/data column × ANSI on/off; overflow/invalid input error behavior, NULLs, decimal narrowing, time zones, legacy timestamp exposure, original safe-pruning behavior, serialization and copying. ## Accuracy and performance evidence One immutable local Parquet fixture reused by both implementations and Spark versions. Native Spark oracle rebuilt from fully read input, with no Paimon relation in its lineage. Exact full-row comparison or structured ANSI error comparison on every sample. Six workloads × 3 warmups + 11 measurements: all **168 fixed samples** matched the oracle (including 28 expected overflow errors); baseline had **84 mismatches** across string-date loss, NULL failures and suppressed overflow. Safe integer-to-string filtering still reads 1 of 64 partition files. Unsafe string-to-date now reads 4 files and returns 1,024 rows rather than reading 1 file and returning only 512. This fallback may increase scan costs; it is not a no-regression claim. Wrong answers/exceptions are excluded from timing ratios. Even unchanged control-query timings varied strongly across JVMs, so local median/P95 numbers are descriptive, not a speedup claim or production forecast. Benchmark scripts/evidence remain local review artifacts, not production dependencies. No full repository, remote-storage or Spark 3.2/3.3 runtime matrix was run. Existing plain-field signed-zero, typed-null literal and nested-name issues are outside this patch. -- 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]
