mxtymoshyk opened a new pull request, #40344:
URL: https://github.com/apache/beam/pull/40344
Fixes #27234
`ParquetIO.read(schema)` and `ParquetIO.readFiles(schema)` fail when the
schema has a field the file doesn't have, which is the usual case after you add
a field to a class and read older files. This PR makes ParquetIO use that
schema as the Avro reader schema, so parquet-avro resolves each file record
against it by field name.
## Reproduction
Write a file with two fields, then read it with a schema that adds a
nullable third field:
```java
Schema oldSchema = /* {firstName: string, middleName: string} */;
Schema newSchema = /* same, plus {"name": "lastName", "type": ["null",
"string"], "default": null} */;
// write with oldSchema via FileIO.write().via(ParquetIO.sink(oldSchema)) ...
pipeline.apply(ParquetIO.read(newSchema).from(path));
pipeline.run().waitUntilFinish();
```
On current master this fails with:
```
org.apache.parquet.io.ParquetDecodingException: Can not read value at 1 in
block 1 in file ...
Caused by: java.lang.ArrayIndexOutOfBoundsException: Index 2 out of bounds
for length 2
at org.apache.avro.generic.GenericData$Record.get(GenericData.java:289)
...
at
org.apache.avro.generic.GenericDatumWriter.writeRecord(GenericDatumWriter.java:234)
```
A related case fails without an exception. If the read schema leaves out a
column that isn't the last one (for example, reading `{name, id}` files with
`{id}`), each record comes back with the `name` value in the `id` field.
## Root cause
`ReadFiles` only used the schema to build the output coder (`AvroCoder` or a
Beam `SchemaCoder`). `SplitReadFn` never got it, and nothing called
`AvroReadSupport.setAvroReadSchema`. So parquet-avro built each `GenericRecord`
with the writer schema stored in the file footer, and `AvroCoder` then encoded
it by field position against the user's schema. When the two schemas differ in
length or order, `GenericDatumWriter` reads the wrong positions or runs past
the end of the record.
## Workaround for released versions
Pass the read schema through the Hadoop configuration:
```java
ParquetIO.read(newSchema)
.from(path)
.withConfiguration(
Collections.singletonMap("parquet.avro.read.schema",
newSchema.toString()));
```
This covers added fields. It won't help with a read schema that has fewer
columns than the file, because parquet-avro rejects file columns missing from
the read schema (`Parquet/Avro schema mismatch: Avro field 'x' not found`). Use
`withProjection` for that case.
## Changes
`ReadFiles` passes the output schema to `SplitReadFn`. That's the schema
passed to `readFiles()`/`read()`, or the encoder schema when `withProjection`
is set, the same one the coder uses. `SplitReadFn` sets it as
`parquet.avro.read.schema` unless the user already put that key in the
configuration.
Without a projection, `SplitReadFn` also drops top-level file columns that
match no field name or alias in the read schema. parquet-avro requires every
requested column to exist in the read schema, so without this step a subset
schema that works today by accident (a prefix of the file's columns) would
start failing.
`parseGenericRecords` and `parseFilesGenericRecords` take no schema and are
unchanged.
## Behavior change
Fields now resolve by name (or Avro alias) instead of by position. A
pipeline that reads files whose column names differ from its Avro field names,
for example files written by another tool and read with a hand-written schema,
used to get position-matched data and will now fail with `Avro field 'x' not
found` or with a null in a non-nullable field. I think that failure is correct,
since positional matching gives the wrong data as soon as column order differs.
Reviewers may want this listed under Breaking Changes in `CHANGES.md` instead
of Bugfixes. I'm happy to move it.
## Tests
New tests in `ParquetIOTest`:
- `testReadWithAddedNullableField`: the reproduction above. The new field
reads as null.
- `testReadFilesWithAddedFieldWithDefault`: a non-null default (`"unknown"`)
gets filled in, which shows real schema resolution and not only null padding.
Goes through `readFiles()`.
- `testReadWithSubsetSchemaWithoutProjection`: a schema that drops the first
column returns the right ids.
- `testReadWithProjectionAndAddedField`: `withProjection` with an encoder
schema that has an extra field.
- `testReadSchemaFromConfigurationIsNotOverridden`: a
`parquet.avro.read.schema` set in the configuration wins.
- `testPruneToReadSchema`: unit test for the column pruning, including alias
matching and the no-match case.
With the fix disabled, the first four fail: three with the
`ParquetDecodingException` shown above, and the subset test with wrong values.
The existing 19 tests pass unchanged. `./gradlew :sdks:java:io:parquet:check
:sdks:java:io:parquet:javadoc` passes locally.
## Not covered
Pruning only works on top-level columns. A nested record in the read schema
that has fewer fields than the file's nested group still fails in parquet-avro,
the same as on master.
------------------------
- [x] Mention the appropriate issue in your description (for example:
`addresses #123`), if applicable. This will automatically add a link to the
pull request in the issue. If you would like the issue to automatically close
on merging the pull request, comment `fixes #<ISSUE NUMBER>` instead.
- [x] Update `CHANGES.md` with noteworthy changes.
- [ ] If this contribution is large, please file an Apache [Individual
Contributor License Agreement](https://www.apache.org/licenses/icla.pdf).
See the [Contributor Guide](https://beam.apache.org/contribute) for more
tips on [how to make review process
smoother](https://github.com/apache/beam/blob/master/CONTRIBUTING.md#make-the-reviewers-job-easier).
--
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]