Dustin Smith created SPARK-60050:
------------------------------------
Summary: Convert only timestamp-bearing ORC file columns when
checking timestamp compatibility per split
Key: SPARK-60050
URL: https://issues.apache.org/jira/browse/SPARK-60050
Project: Spark
Issue Type: Improvement
Components: SQL
Affects Versions: 5.0.0
Reporter: Dustin Smith
{{OrcUtils.requestedColumnIds}} runs once per file split in the ORC readers (V1
{{OrcFileFormat}} and V2 {{OrcPartitionReaderFactory}}). It calls
{{checkTimestampCompatibility(toCatalystSchema(orcSchema), dataSchema)}} to
reject timestamp family or precision mismatches with a clear error.
{{toCatalystSchema}} converts the whole file schema and runs
{{CatalystSqlParser.parseDataType}} for every primitive column, so the cost
grows with the width of the file, not with the columns the query reads, and is
paid on every split.
The check can only fail for a column whose read type is, or contains, a
timestamp type, so converting the other columns is wasted work.
Measured on 400 ORC files with 1000 INT columns (local[1], JDK 17):
* {{toCatalystSchema}} on the 1000-column file schema: 0.75 ms per call
* {{SELECT c0}}: 715 ms before, 384 ms after converting only timestamp-bearing
columns
Proposed change: in {{requestedColumnIds}}, convert and check only the file
columns whose read type contains a timestamp. The errors for timestamp
mismatches stay the same. As a side effect, a file column whose ORC type the
SQL parser cannot handle (such as a Hive {{uniontype}}) no longer fails the
read when its read type has no timestamp.
The filter pushdown path still converts the full file schema per split; that is
left for a separate change.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]