Skip to content

Commit b1e2d03

Browse files
GH-3816: Automatically detect list encodings in AvroReadSupport (#3753)
### Rationale for this change parquet-avro supports writing both "old" and "new" list encodings via the [parquet.avro.write-old-list-structure](https://github.com/apache/parquet-java/blob/apache-parquet-1.18.0/parquet-avro/src/main/java/org/apache/parquet/avro/AvroWriteSupport.java#L72-L73) config. "old" encodings (aka "2-level"), which wrap the list in a `repeated group array` schema, are the default; "new" encodings (aka "3-level") are opt-in. On the reader side, if you're using `ParquetAvroReader` to read data that was written using `ParquetAvroWriter`, and don't specify a projection, both type sof list encoding get parsed automatically from a combination of the file schema + the `parquet.avro.schema` metadata key. There's no need to set `parquet.avro.write-old-list-structure` key in your Configuration. However, if you're either: - specifying a projection (`AvroReadSupport.setRequestedProjection(...)`), or - reading data _not_ written using ParquetAvroWriter (and thus not containing the `parquet.avro.schema` metadata key), 3-levle list encodings will not be parsed correctly - the reader will inject an extra nested record, named `element`, into the list item type. As a reader this introduces some pain, since you have to look up the underlying file metadata of the upstream Parquet file, and risk reading incorrect data. This PR attempts to automatically detect new list encodings based on the writer file schema. lmk what you think of this change. Automatic inference is always a bit risky, but I tried to be conservative with the approach (only set the list structure property if _all_ list fields in the schema use 3-level encoding; don't override `parquet.avro.write-old-list-structure` if the user is already setting it). any ideas for a better approach here are welcome - this is becoming more of a pain point as 3-level lists become a more popular option among other writer sdks. ### What changes are included in this PR? A new read configuration property `parquet.avro.read.autoDetectListStructure` (defaulting to true) that will instruct AvroReadSupport to automatically set List configuration properties based on parsing the writer file schema. ### Are these changes tested? Yes, unit tests + locally on real data. ### Are there any user-facing changes? Yes, since the new property defaults to `true` - it would impact anyone who's reading 3-level list data without setting the `parquet.avro.write-old-list-structure` key and who's relying on/working around the incorrectly formatted data (e.g. `{"locations": [{"element": {"latitude": 0.0, "longitude": 180.0}}, ...]}` instead of `{"locations": [{"latitude": 0.0, "longitude": 180.0}, ...]}` . additionally, this change also modifies the underlying Configuration object to add the properties. Closes #3816
1 parent 29347ba commit b1e2d03

3 files changed

Lines changed: 471 additions & 10 deletions

File tree

‎parquet-avro/README.md‎

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -26,14 +26,15 @@ Apache Avro integration
2626

2727
### Configuration for reading
2828

29-
| Name | Type | Description |
30-
|-----------------------------------------|-----------|----------------------------------------------------------------------|
31-
| `parquet.avro.data.supplier` | `Class` | The implementation of the interface org.apache.parquet.avro.AvroDataSupplier. Available implementations in the library: GenericDataSupplier, ReflectDataSupplier, SpecificDataSupplier.<br/>The default value is `org.apache.parquet.avro.SpecificDataSupplier` |
32-
| `parquet.avro.read.schema` | `String` | The Avro schema to be used for reading. It shall be compatible with the file schema. The file schema will be used directly if not set. |
33-
| `parquet.avro.projection` | `String` | The Avro schema to be used for projection. |
34-
| `parquet.avro.compatible` | `boolean` | Flag for compatibility mode. `true` for materializing Avro `IndexedRecord` objects, `false` for materializing the related objects for either generic, specific, or reflect records.<br/>The default value is `true`. |
35-
| `parquet.avro.readInt96AsFixed` | `boolean` | Flag for handling the `INT96` Parquet types. `true` for converting it to the `fixed` Avro type, `false` for not handling `INT96` types (throwing exception).<br/>The default value is `false`.<br/>**NOTE: The `INT96` Parquet type is deprecated. This option is only to support old data.** |
36-
| `parquet.avro.serializable.classes` | `String` | List of the fully qualified class names separated by ',' that may be referenced from the Avro schema by "java-class" or "java-key-class" and are allowed to be loaded. |
29+
| Name | Type | Description |
30+
|---------------------------------------------|-----------|-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
31+
| `parquet.avro.data.supplier` | `Class` | The implementation of the interface org.apache.parquet.avro.AvroDataSupplier. Available implementations in the library: GenericDataSupplier, ReflectDataSupplier, SpecificDataSupplier.<br/>The default value is `org.apache.parquet.avro.SpecificDataSupplier` |
32+
| `parquet.avro.read.schema` | `String` | The Avro schema to be used for reading. It shall be compatible with the file schema. The file schema will be used directly if not set. |
33+
| `parquet.avro.projection` | `String` | The Avro schema to be used for projection. |
34+
| `parquet.avro.compatible` | `boolean` | Flag for compatibility mode. `true` for materializing Avro `IndexedRecord` objects, `false` for materializing the related objects for either generic, specific, or reflect records.<br/>The default value is `true`. |
35+
| `parquet.avro.readInt96AsFixed` | `boolean` | Flag for handling the `INT96` Parquet types. `true` for converting it to the `fixed` Avro type, `false` for not handling `INT96` types (throwing exception).<br/>The default value is `false`.<br/>**NOTE: The `INT96` Parquet type is deprecated. This option is only to support old data.** |
36+
| `parquet.avro.serializable.classes` | `String` | List of the fully qualified class names separated by ',' that may be referenced from the Avro schema by "java-class" or "java-key-class" and are allowed to be loaded. |
37+
| `parquet.avro.read.autoDetectListStructure` | `boolean` | Automatically detect whether the write schema uses 2- or 3-level list encoding and converts the projection/read schema accordingly.<br/>The default value is `true`. |
3738

3839
### Configuration for writing
3940

‎parquet-avro/src/main/java/org/apache/parquet/avro/AvroReadSupport.java‎

Lines changed: 83 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,10 +28,14 @@
2828
import org.apache.hadoop.util.ReflectionUtils;
2929
import org.apache.parquet.conf.HadoopParquetConfiguration;
3030
import org.apache.parquet.conf.ParquetConfiguration;
31+
import org.apache.parquet.conf.PlainParquetConfiguration;
3132
import org.apache.parquet.hadoop.api.ReadSupport;
3233
import org.apache.parquet.hadoop.util.ConfigurationUtil;
3334
import org.apache.parquet.io.api.RecordMaterializer;
35+
import org.apache.parquet.schema.GroupType;
36+
import org.apache.parquet.schema.LogicalTypeAnnotation;
3437
import org.apache.parquet.schema.MessageType;
38+
import org.apache.parquet.schema.Type;
3539
import org.slf4j.Logger;
3640
import org.slf4j.LoggerFactory;
3741

@@ -63,6 +67,11 @@ public class AvroReadSupport<T> extends ReadSupport<T> {
6367
public static final String READ_INT96_AS_FIXED = "parquet.avro.readInt96AsFixed";
6468
public static final boolean READ_INT96_AS_FIXED_DEFAULT = false;
6569

70+
// Automatically detect whether a Parquet file uses 2-level or 3-level encoding;
71+
// Ignored if AvroWriteSupport.WRITE_OLD_LIST_STRUCTURE is also set
72+
public static final String AUTO_DETECT_LIST_STRUCTURE = "parquet.avro.read.autoDetectListStructure";
73+
static final boolean AUTO_DETECT_LIST_STRUCTURE_DEFAULT = true;
74+
6675
/**
6776
* List of the fully qualified class names separated by ',' that may be referenced from the Avro schema by
6877
* "java-class" or "java-key-class" and are allowed to be loaded.
@@ -131,7 +140,8 @@ public ReadContext init(
131140
String requestedProjectionString = configuration.get(AVRO_REQUESTED_PROJECTION);
132141
if (requestedProjectionString != null) {
133142
Schema avroRequestedProjection = new Schema.Parser().parse(requestedProjectionString);
134-
projection = new AvroSchemaConverter(configuration).convert(avroRequestedProjection);
143+
projection = new AvroSchemaConverter(getDerivedListEncodingConf(configuration, fileSchema))
144+
.convert(avroRequestedProjection);
135145
}
136146

137147
String avroReadSchema = configuration.get(AVRO_READ_SCHEMA);
@@ -176,7 +186,8 @@ public RecordMaterializer<T> prepareForRead(
176186
avroSchema = new Schema.Parser().parse(keyValueMetaData.get(OLD_AVRO_SCHEMA_METADATA_KEY));
177187
} else {
178188
// default to converting the Parquet schema into an Avro schema
179-
avroSchema = new AvroSchemaConverter(configuration).convert(parquetSchema);
189+
avroSchema = new AvroSchemaConverter(getDerivedListEncodingConf(configuration, fileSchema))
190+
.convert(parquetSchema);
180191
}
181192

182193
GenericData model = getDataModel(configuration, avroSchema);
@@ -230,4 +241,74 @@ private GenericData getDataModel(ParquetConfiguration conf, Schema schema) {
230241
return ReflectionUtils.newInstance(suppClass, ConfigurationUtil.createHadoopConfiguration(conf))
231242
.get();
232243
}
244+
245+
// Returns a ParquetConfiguration with appropriate list-decoding properties set, inferred from
246+
// the file schema as well as user-supplied Configuration properties.
247+
// If no configuration changes are required, the original ParquetConfiguration object will be;
248+
// returned; otherwise, a copy will be created with the correct properties.
249+
private static ParquetConfiguration getDerivedListEncodingConf(
250+
ParquetConfiguration configuration, MessageType fileSchema) {
251+
final boolean autoDetectListStructure =
252+
configuration.getBoolean(AUTO_DETECT_LIST_STRUCTURE, AUTO_DETECT_LIST_STRUCTURE_DEFAULT);
253+
254+
if (!autoDetectListStructure
255+
|| configuration.get(AvroWriteSupport.WRITE_OLD_LIST_STRUCTURE) != null
256+
|| configuration.get(AvroSchemaConverter.ADD_LIST_ELEMENT_RECORDS) != null
257+
|| !writesNewListStructure(fileSchema)) {
258+
return configuration;
259+
}
260+
261+
// Avoid mutating the original Configuration by creating a copy, with new properties set
262+
final ParquetConfiguration copiedConfiguration = new PlainParquetConfiguration();
263+
for (Map.Entry<String, String> property : configuration) {
264+
copiedConfiguration.set(property.getKey(), property.getValue());
265+
}
266+
copiedConfiguration.setBoolean(AvroWriteSupport.WRITE_OLD_LIST_STRUCTURE, false);
267+
copiedConfiguration.setBoolean(AvroSchemaConverter.ADD_LIST_ELEMENT_RECORDS, false);
268+
269+
return copiedConfiguration;
270+
}
271+
272+
private static boolean writesNewListStructure(MessageType schema) {
273+
return Boolean.TRUE.equals(allListStructuresAreThreeLevel(schema));
274+
}
275+
276+
// Given a Parquet schema, return true only if the schema:
277+
// - contains one or more List fields
278+
// - encodes every List field using 3-level list structure
279+
private static Boolean allListStructuresAreThreeLevel(Type type) {
280+
if (type.isPrimitive()) {
281+
return null;
282+
}
283+
GroupType group = type.asGroupType();
284+
if (group.getLogicalTypeAnnotation() instanceof LogicalTypeAnnotation.ListLogicalTypeAnnotation) {
285+
if (group.isRepetition(Type.Repetition.REPEATED) || group.getFieldCount() != 1) {
286+
return false;
287+
}
288+
Type repeated = group.getType(0);
289+
if (repeated.isPrimitive()
290+
|| !repeated.isRepetition(Type.Repetition.REPEATED)
291+
|| !repeated.getName().equals("list")
292+
|| repeated.asGroupType().getFieldCount() != 1) {
293+
return false;
294+
}
295+
Type element = repeated.asGroupType().getType(0);
296+
if (element.isRepetition(Type.Repetition.REPEATED)
297+
|| !element.getName().equals("element")) {
298+
return false;
299+
}
300+
return !Boolean.FALSE.equals(allListStructuresAreThreeLevel(element));
301+
}
302+
Boolean result = null;
303+
for (Type field : group.getFields()) {
304+
Boolean fieldListStructuresAreThreeLevel = allListStructuresAreThreeLevel(field);
305+
if (Boolean.FALSE.equals(fieldListStructuresAreThreeLevel)) {
306+
return false;
307+
}
308+
if (Boolean.TRUE.equals(fieldListStructuresAreThreeLevel)) {
309+
result = true;
310+
}
311+
}
312+
return result;
313+
}
233314
}

0 commit comments

Comments
 (0)