sunchao commented on code in PR #6515:
URL: https://github.com/apache/datafusion-comet/pull/6515#discussion_r4161786070
##########
native/core/src/parquet/schema_adapter.rs:
##########
@@ -478,33 +478,79 @@ enum ConversionCheck {
},
}
-/// Apply the rejection matrix of Spark's
`ParquetVectorUpdaterFactory.getUpdater` to a single
-/// physical/logical leaf pair. `column` is the Spark-style column path used
in the error (`a`
-/// for a top-level column, `s, x` for a nested leaf, mirroring
-/// `Arrays.toString(descriptor.getPath())`). The rules and their order are
exactly those the
-/// adapter applies to top-level columns; [`check_conversion`] applies them to
nested leaves.
+/// Whether Parquet stores `data_type` as a group: a struct, a list or a map.
+fn is_complex(data_type: &DataType) -> bool {
+ matches!(
+ data_type,
+ DataType::Struct(_)
+ | DataType::List(_)
+ | DataType::LargeList(_)
+ | DataType::FixedSizeList(_, _)
+ | DataType::ListView(_)
+ | DataType::LargeListView(_)
+ | DataType::Map(_, _)
+ )
+}
+
+/// Check a pair that [`check_conversion`] doesn't walk: two primitives, or
two types of
+/// different shape. `column` is the Spark-style column path used in the error
(`a` for a
+/// top-level column, `s, x` for a nested leaf, mirroring
`Arrays.toString(descriptor.getPath())`).
+///
+/// A shape mismatch (e.g. TIMESTAMP read as ARRAY<TIMESTAMP>, or STRUCT read
as ARRAY) fails
+/// when Spark opens the file if Spark can't clip the file's type to the
requested one: a group
+/// read as another type (`ParquetToSparkSchemaConverter`), or a primitive
read as a struct, or
+/// as an array or map with a complex element
(`ParquetReadSupport.clipParquetType`). Every
+/// other pair Spark rejects, including a primitive read as an array or map of
primitives
+/// (SPARK-45604), is rejected only by `getUpdater`, which Spark calls while
decoding a row
+/// group, so the rejection is deferred to runtime (#6506).
fn check_leaf_conversion(
physical_type: &DataType,
target_type: &DataType,
column: &str,
options: &SparkParquetOptions,
-) -> ConversionCheck {
+) -> DataFusionResult<ConversionCheck> {
if physical_type == target_type {
- return ConversionCheck::Accept;
- }
- let reject = || {
- ConversionCheck::Reject(parquet_schema_convert_err(
- column,
- physical_type,
- target_type,
- ))
- };
- let reject_on_non_empty = || ConversionCheck::RejectOnNonEmpty {
+ return Ok(ConversionCheck::Accept);
+ }
+ if is_complex(physical_type) || is_complex(target_type) {
+ let is_unclipped = !is_complex(physical_type)
+ && match target_type {
+ DataType::List(item)
+ | DataType::LargeList(item)
+ | DataType::FixedSizeList(item, _)
+ | DataType::ListView(item)
+ | DataType::LargeListView(item) =>
!is_complex(item.data_type()),
Review Comment:
[P2] Preserve open-time rejection for legacy LIST encodings. With
`spark.sql.parquet.writeLegacyFormat=true`, a non-null `array<int>` is stored
with a repeated primitive child. Reading it as `array<array<int>>` or
`array<map<int,int>>` makes Spark's `clipParquetListType` fail before decoding,
including for empty files and fully pruned row groups. Arrow normalizes the
file type to a list, so this check treats the nested primitive-to-container
mismatch as deferred. The head now succeeds where both Spark and the base fail,
bypassing schema validation for legacy datasets. Could this classification
retain the original Parquet list encoding and preserve eager rejection for this
shape?
Evidence: Wrote `select 1 as id, array(1) as a` with legacy format enabled,
once with `where false` and once with one row. Read each with `id int, a
array<array<int>>`, applying `id = 100` to the nonempty file. Spark 3.5.9 and
4.1.3 both failed in `clipParquetListType`. The exact-head native scan returned
zero rows without error, while the base adapter rejected both files. Repeated
successfully with `array<map<int,int>>`. Standard LIST encoding remained
readable in Spark and the head. Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3,
and 4.2.0 all contain the relevant repeated-primitive assertion.
##########
native/core/src/parquet/schema_adapter.rs:
##########
@@ -478,33 +478,79 @@ enum ConversionCheck {
},
}
-/// Apply the rejection matrix of Spark's
`ParquetVectorUpdaterFactory.getUpdater` to a single
-/// physical/logical leaf pair. `column` is the Spark-style column path used
in the error (`a`
-/// for a top-level column, `s, x` for a nested leaf, mirroring
-/// `Arrays.toString(descriptor.getPath())`). The rules and their order are
exactly those the
-/// adapter applies to top-level columns; [`check_conversion`] applies them to
nested leaves.
+/// Whether Parquet stores `data_type` as a group: a struct, a list or a map.
+fn is_complex(data_type: &DataType) -> bool {
+ matches!(
+ data_type,
+ DataType::Struct(_)
+ | DataType::List(_)
+ | DataType::LargeList(_)
+ | DataType::FixedSizeList(_, _)
+ | DataType::ListView(_)
+ | DataType::LargeListView(_)
+ | DataType::Map(_, _)
+ )
+}
+
+/// Check a pair that [`check_conversion`] doesn't walk: two primitives, or
two types of
+/// different shape. `column` is the Spark-style column path used in the error
(`a` for a
+/// top-level column, `s, x` for a nested leaf, mirroring
`Arrays.toString(descriptor.getPath())`).
+///
+/// A shape mismatch (e.g. TIMESTAMP read as ARRAY<TIMESTAMP>, or STRUCT read
as ARRAY) fails
+/// when Spark opens the file if Spark can't clip the file's type to the
requested one: a group
+/// read as another type (`ParquetToSparkSchemaConverter`), or a primitive
read as a struct, or
+/// as an array or map with a complex element
(`ParquetReadSupport.clipParquetType`). Every
+/// other pair Spark rejects, including a primitive read as an array or map of
primitives
+/// (SPARK-45604), is rejected only by `getUpdater`, which Spark calls while
decoding a row
+/// group, so the rejection is deferred to runtime (#6506).
fn check_leaf_conversion(
physical_type: &DataType,
target_type: &DataType,
column: &str,
options: &SparkParquetOptions,
-) -> ConversionCheck {
+) -> DataFusionResult<ConversionCheck> {
if physical_type == target_type {
- return ConversionCheck::Accept;
- }
- let reject = || {
- ConversionCheck::Reject(parquet_schema_convert_err(
- column,
- physical_type,
- target_type,
- ))
- };
- let reject_on_non_empty = || ConversionCheck::RejectOnNonEmpty {
+ return Ok(ConversionCheck::Accept);
+ }
+ if is_complex(physical_type) || is_complex(target_type) {
+ let is_unclipped = !is_complex(physical_type)
+ && match target_type {
+ DataType::List(item)
+ | DataType::LargeList(item)
+ | DataType::FixedSizeList(item, _)
+ | DataType::ListView(item)
+ | DataType::LargeListView(item) =>
!is_complex(item.data_type()),
+ DataType::Map(entries, _) => matches!(
+ entries.data_type(),
+ DataType::Struct(kv) if kv.iter().all(|f|
!is_complex(f.data_type()))
+ ),
+ _ => false,
+ };
+ if !is_unclipped {
+ return Err(parquet_schema_convert_err(
+ column,
+ physical_type,
+ target_type,
+ ));
+ }
+ } else if spark_has_updater(physical_type, target_type, options) {
+ return Ok(ConversionCheck::Accept);
+ }
+ Ok(ConversionCheck::RejectOnNonEmpty {
Review Comment:
[P2] Enforce newly deferred conversion errors before pushed row filtering.
With `spark.comet.parquet.rowFilterPushdown.enabled=true`, a file containing
`(id, s) = (1, 'bad'), (3, 'bad')`, read as `id int, s int` with `id = 2`, now
succeeds with zero rows. The row group survives statistics pruning, but row
filtering removes every row before the projection evaluates `RejectOnNonEmpty`.
Spark and the base raise the string-to-int conversion error. Although the new
documentation acknowledges this limitation, the PR introduces it for
conversions that previously failed correctly. Please validate rejected
conversions after format pruning but before row-level selection, or disable
row-filter pushdown for affected files.
Evidence: Created one Spark-written row group with dictionary encoding
disabled and `id` bounds [1, 3], then queried `id = 2` with read schema `id
int, s int`. Spark 3.5.9 and 4.1.3 raised
`SchemaColumnConvertNotSupportedException` for `s`. Disposable native scan
tests showed: base with pushdown enabled = error; head with pushdown disabled =
error; head with pushdown enabled = zero rows and no error. Thus this
BINARY-to-int case is a base-relative regression, independent of the inherited
behavior for other deferred conversions.
--
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]