manuzhang opened a new issue, #25801: URL: https://github.com/apache/datafusion/issues/25801
### Is your feature request related to a problem or challenge? When a table schema declares a string type for a column that the Parquet file stores as `BINARY` without the `UTF8` annotation, the Parquet opener always rewrites the file schema so the reader decodes the column straight into a string array (`apply_file_schema_type_coercions` in `opener/mod.rs`). The Parquet reader only validates UTF-8 for annotated columns (`ByteArrayColumnValueDecoder` sets `validate_utf8` from `converted_type() == UTF8`). So when the bytes are not valid UTF-8, the scan returns an invalid `StringArray` in release builds and panics in debug builds, in `OffsetBuffer::into_array`. The module docs in `schema_coercion.rs` describe this, and #25342 left map children out of the coercion for the same reason. There is no way to opt out. The coercion runs before the `PhysicalExprAdapter`, so a custom adapter never sees a binary column that it could convert by its own rules. Apache DataFusion Comet runs into this with Spark's `spark.sql.parquet.binaryAsString`. Spark types unannotated `BINARY` columns as strings and replaces invalid bytes with U+FFFD. To get those semantics, Comet has to rewrite the table schema to `Binary` and add its own conversion projection. Filters are typed against the string schema, so that workaround also stops Comet from pushing any Parquet filters while the setting is on (apache/datafusion-comet#6246). ### Describe the solution you'd like Add a Parquet read option, `datafusion.execution.parquet.coerce_binary_to_string`, defaulting to `true`. When it is `false`, the opener leaves binary file columns as binary and the scan's `PhysicalExprAdapter` casts them to the table's string type. With the default adapter that cast rejects invalid UTF-8, and a custom adapter can convert the bytes however it needs to. ### Describe alternatives you've considered - Validate UTF-8 in the Parquet reader whenever the requested Arrow type is a string. That belongs in arrow-rs and would turn invalid bytes into errors, but it would still give callers no way to apply their own conversion. - A builder method on `ParquetSource` instead of a config option. That is smaller, but unlike the neighbouring `binary_as_string`, `schema_force_view_types` and `coerce_int96` options it would not be settable through SQL or serialized with the plan. -- 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]
