[Bug]: ParquetIO failing to read with specified schema
- Dominant language
- Java
- Stars
- 8.7k
- Forks
- 4.7k
- Avg merge
- 1d 20h
- Merged PRs (30d)
- 196
Description
### What happened?
Reading from parquet file is failing with following exception
```
java.lang.ArrayIndexOutOfBoundsException: 3
org.apache.beam.sdk.Pipeline$PipelineExecutionException: java.lang.ArrayIndexOutOfBoundsException: 3
at org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:371)
at org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:339)
at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:219)
at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:67)
```
I'm trying to read a parquet file in Apache beam using following code snippet
```
String path = "pathToFile";
Schema schema = ReflectData.get().getSchema(EntityClass.class)
PCollection records = pipeline.apply("read a file", ParquetIO.read(schema).from(path));
pipeline.run().waitUntilFinish();
```
The above code is failing to read the after adding new fields to the java class.
Previous java class:
```
@With
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
@ParametersAreNonnullByDefault
@DefaultSchema(JavaBeanSchema.class)
public class EntityClass {
private String firstName;
private String middleName;
}
```
Unable to read after making this change to the java class
```
@With
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
@ParametersAreNonnullByDefault
@DefaultSchema(JavaBeanSchema.class)
public class EntityClass {
private String firstName;
private String middleName;
@Nullable
@org.apache.avro.reflect.Nullable
private String lastName;
}
```
The schema generated from reflection seems to be fine.
```
{
"type": "record",
"name": "EntityClass",
"namespace": "com.test.entityclass",
"fields": [{
"name": "firstName",
"type": "string"
}, {
"name": "middleName",
"type": "string"
}, {
"name": "lastName",
"type": ["null", "string"],
"default": null
}]
}
```
When trying to debug, in [ParquetIO.java#L797](https://github.com/apache/beam/blob/ea9147ad2946f72f7d52924cba2820e9aae7cd91/sdks/java/io/parquet/src/main/java/org/apache/beam/sdk/io/parquet/ParquetIO.java#L797) the here is still picked up from parquet file instead of the specified schema.
However, when trying to use AvroReader outside Beam, specifying the schema in AvroParquetSupport is required to support schema change as shown below:
```
AvroReadSupport.setAvroReadSchema(conf, NEW_SCHEMA);
try (ParquetReader reader = AvroParquetReader.builder(dataFile)
.withConf(conf)
.build()) {
GenericRecord genericRecord;
while ((genericRecord = reader.read()) != null) {
//System.out.println(genericRecord.getSchema());
T result = apply((GenericRecord) genericRecord.get("payload"), clazz);
resultFiles.add(result);
}
}
```
Specifying the schema in AvroReadSupport could be a potential solution if this is a genuine bug.
### Issue Priority
Priority: 2 (default / most bugs should be filed as P2)
### Issue Components
- [ ] Component: Python SDK
- [X] Component: Java SDK
- [ ] Component: Go SDK
- [ ] Component: Typescript SDK
- [x] Component: IO connector
- [ ] Component: Beam examples
- [ ] Component: Beam playground
- [ ] Component: Beam katas
- [ ] Component: Website
- [ ] Component: Spark Runner
- [ ] Component: Flink Runner
- [ ] Component: Samza Runner
- [ ] Component: Twister2 Runner
- [ ] Component: Hazelcast Jet Runner
- [ ] Component: Google Cloud Dataflow Runner
Contributor guide
Assessment
This issue has not been assessed yet.