andygrove opened a new issue, #6504:
URL: https://github.com/apache/datafusion-comet/issues/6504
### Describe the bug
The native Iceberg scan can't read a data file that was written before a
nested field was added to a struct column, or to the struct inside an array or
map. iceberg-rust's `RecordBatchTransformer` casts the file's struct to the
table's wider struct with arrow's struct cast, which can't add a field, so the
scan fails with `Iceberg scan error: ... Incorrect number of arrays for
StructArray fields, expected 2 got 1`. This is apache/iceberg-rust#2617. A
plain `SELECT s FROM t` on such a table already fails this way in 1.0.0.
1.1.0 makes it reach more queries. 1.0.0 sent an Iceberg scan to Spark
whenever its pushed filters held `IS NULL` or `IS NOT NULL` on a struct, array
or map column. #5732 removed that fallback, so these queries now run natively
and fail where 1.0.0 returned the right rows:
- `WHERE s IS NOT NULL` or `WHERE s IS NULL` on the evolved column
- every non-outer `explode`, `posexplode` or `inline` of an evolved array or
map, because Spark infers `isnotnull(items)` for it and pushes it into the scan
- joins on an evolved struct column, for the same reason
### Steps to reproduce
```sql
CREATE TABLE cat.db.evo_s (id INT, s STRUCT<a: INT>) USING iceberg;
INSERT INTO cat.db.evo_s VALUES (1, named_struct('a', 1)), (2, null);
ALTER TABLE cat.db.evo_s ADD COLUMN s.b INT;
INSERT INTO cat.db.evo_s VALUES (3, named_struct('a', 3, 'b', 30));
SELECT id, s FROM cat.db.evo_s WHERE s IS NOT NULL;
CREATE TABLE cat.db.evo_l (id INT, items ARRAY<STRUCT<a: INT>>) USING
iceberg;
INSERT INTO cat.db.evo_l VALUES (1, array(named_struct('a', 10))), (2, null);
ALTER TABLE cat.db.evo_l ADD COLUMN items.element.b INT;
INSERT INTO cat.db.evo_l VALUES (3, array(named_struct('a', 30, 'b', 300)));
SELECT id, e.a, e.b FROM cat.db.evo_l LATERAL VIEW explode(items) x AS e;
```
With the native Iceberg scan on, which is the default, both queries fail on
1.1.0-rc1. On 1.0.0 both fall back to Spark, with the reason "IS NULL / IS NOT
NULL predicates on complex type columns (struct/array/map) are not yet
supported by iceberg-rust", and return the right rows. `SELECT id, s FROM
cat.db.evo_s` fails on both.
### Expected behavior
Spark's results, which fill the missing nested field with NULL: `(1, {1,
null})` and `(3, {3, 30})` for the first query, and `(1, 10, null)` and `(3,
30, 300)` for the second.
### Workaround
`spark.comet.scan.icebergNative.enabled=false` reads every Iceberg table
with Spark's reader.
### Additional context
Verified on Spark 4.1 with Iceberg 1.11.0. Until iceberg-rust can add fields
in this cast, Comet could decline the native scan when a projected complex
column has gained nested fields over the table's schema history. That would
also fix the plain read that already fails in 1.0.0.
Found by the 1.1.0 regression audit (#6399) and tracked in #6402.
--
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]