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]

Reply via email to