Parquet format

Flink supports reading Parquet files, producing Flink RowData and producing Avro records. To use the format you need to add the flink-parquet dependency to your project:

  1. <dependency>
  2. <groupId>org.apache.flink</groupId>
  3. <artifactId>flink-parquet</artifactId>
  4. <version>1.16.0</version>
  5. </dependency>

To read Avro records, you will need to add the parquet-avro dependency:

  1. <dependency>
  2. <groupId>org.apache.parquet</groupId>
  3. <artifactId>parquet-avro</artifactId>
  4. <version>1.12.2</version>
  5. <optional>true</optional>
  6. <exclusions>
  7. <exclusion>
  8. <groupId>org.apache.hadoop</groupId>
  9. <artifactId>hadoop-client</artifactId>
  10. </exclusion>
  11. <exclusion>
  12. <groupId>it.unimi.dsi</groupId>
  13. <artifactId>fastutil</artifactId>
  14. </exclusion>
  15. </exclusions>
  16. </dependency>

In order to use the Parquet format in PyFlink jobs, the following dependencies are required:

PyFlink JAR
Download

See Python dependency management for more details on how to use JARs in PyFlink.

This format is compatible with the new Source that can be used in both batch and streaming execution modes. Thus, you can use this format for two kinds of data:

  • Bounded data: lists all files and reads them all.
  • Unbounded data: monitors a directory for new files that appear.

When you start a File Source it is configured for bounded data by default. To configure the File Source for unbounded data, you must additionally call AbstractFileSource.AbstractFileSourceBuilder.monitorContinuously(Duration).

Vectorized reader

Java

  1. // Parquet rows are decoded in batches
  2. FileSource.forBulkFileFormat(BulkFormat,Path...)
  3. // Monitor the Paths to read data as unbounded data
  4. FileSource.forBulkFileFormat(BulkFormat,Path...)
  5. .monitorContinuously(Duration.ofMillis(5L))
  6. .build();

Python

  1. # Parquet rows are decoded in batches
  2. FileSource.for_bulk_file_format(BulkFormat, Path...)
  3. # Monitor the Paths to read data as unbounded data
  4. FileSource.for_bulk_file_format(BulkFormat, Path...) \
  5. .monitor_continuously(Duration.of_millis(5)) \
  6. .build()

Avro Parquet reader

Java

  1. // Parquet rows are decoded in batches
  2. FileSource.forRecordStreamFormat(StreamFormat,Path...)
  3. // Monitor the Paths to read data as unbounded data
  4. FileSource.forRecordStreamFormat(StreamFormat,Path...)
  5. .monitorContinuously(Duration.ofMillis(5L))
  6. .build();

Python

  1. # Parquet rows are decoded in batches
  2. FileSource.for_record_stream_format(StreamFormat, Path...)
  3. # Monitor the Paths to read data as unbounded data
  4. FileSource.for_record_stream_format(StreamFormat, Path...) \
  5. .monitor_continuously(Duration.of_millis(5)) \
  6. .build()

Following examples are all configured for bounded data. To configure the File Source for unbounded data, you must additionally call AbstractFileSource.AbstractFileSourceBuilder.monitorContinuously(Duration).

In this example, you will create a DataStream containing Parquet records as Flink RowDatas. The schema is projected to read only the specified fields (“f7”, “f4” and “f99”).
Flink will read records in batches of 500 records. The first boolean parameter specifies that timestamp columns will be interpreted as UTC. The second boolean instructs the application that the projected Parquet fields names are case-sensitive. There is no watermark strategy defined as records do not contain event timestamps.

Java

  1. final LogicalType[] fieldTypes =
  2. new LogicalType[] {
  3. new DoubleType(), new IntType(), new VarCharType()};
  4. final RowType rowType = RowType.of(fieldTypes, new String[] {"f7", "f4", "f99"});
  5. final ParquetColumnarRowInputFormat<FileSourceSplit> format =
  6. new ParquetColumnarRowInputFormat<>(
  7. new Configuration(),
  8. rowType,
  9. InternalTypeInfo.of(rowType),
  10. 500,
  11. false,
  12. true);
  13. final FileSource<RowData> source =
  14. FileSource.forBulkFileFormat(format, /* Flink Path */)
  15. .build();
  16. final DataStream<RowData> stream =
  17. env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

Python

  1. row_type = DataTypes.ROW([
  2. DataTypes.FIELD('f7', DataTypes.DOUBLE()),
  3. DataTypes.FIELD('f4', DataTypes.INT()),
  4. DataTypes.FIELD('f99', DataTypes.VARCHAR()),
  5. ])
  6. source = FileSource.for_bulk_file_format(ParquetColumnarRowInputFormat(
  7. row_type=row_type,
  8. hadoop_config=Configuration(),
  9. batch_size=500,
  10. is_utc_timestamp=False,
  11. is_case_sensitive=True,
  12. ), PARQUET_FILE_PATH).build()
  13. ds = env.from_source(source, WatermarkStrategy.no_watermarks(), "file-source")

Avro Records

Flink supports producing three types of Avro records by reading Parquet files (Only Generic record is supported in PyFlink):

Generic record

Avro schemas are defined using JSON. You can get more information about Avro schemas and types from the Avro specification. This example uses an Avro schema example similar to the one described in the official Avro tutorial:

  1. {"namespace": "example.avro",
  2. "type": "record",
  3. "name": "User",
  4. "fields": [
  5. {"name": "name", "type": "string"},
  6. {"name": "favoriteNumber", "type": ["int", "null"]},
  7. {"name": "favoriteColor", "type": ["string", "null"]}
  8. ]
  9. }

This schema defines a record representing a user with three fields: name, favoriteNumber, and favoriteColor. You can find more details at record specification for how to define an Avro schema.

In the following example, you will create a DataStream containing Parquet records as Avro Generic records. It will parse the Avro schema based on the JSON string. There are many other ways to parse a schema, e.g. from java.io.File or java.io.InputStream. Please refer to Avro Schema for details. After that, you will create an AvroParquetRecordFormat via AvroParquetReaders for Avro Generic records.

Java

  1. // parsing avro schema
  2. final Schema schema =
  3. new Schema.Parser()
  4. .parse(
  5. "{\"type\": \"record\", "
  6. + "\"name\": \"User\", "
  7. + "\"fields\": [\n"
  8. + " {\"name\": \"name\", \"type\": \"string\" },\n"
  9. + " {\"name\": \"favoriteNumber\", \"type\": [\"int\", \"null\"] },\n"
  10. + " {\"name\": \"favoriteColor\", \"type\": [\"string\", \"null\"] }\n"
  11. + " ]\n"
  12. + " }");
  13. final FileSource<GenericRecord> source =
  14. FileSource.forRecordStreamFormat(
  15. AvroParquetReaders.forGenericRecord(schema), /* Flink Path */)
  16. .build();
  17. final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.enableCheckpointing(10L);
  19. final DataStream<GenericRecord> stream =
  20. env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

Python

  1. # parsing avro schema
  2. schema = AvroSchema.parse_string("""
  3. {
  4. "type": "record",
  5. "name": "User",
  6. "fields": [
  7. {"name": "name", "type": "string"},
  8. {"name": "favoriteNumber", "type": ["int", "null"]},
  9. {"name": "favoriteColor", "type": ["string", "null"]}
  10. ]
  11. }
  12. """)
  13. source = FileSource.for_record_stream_format(
  14. AvroParquetReaders.for_generic_record(schema), # file paths
  15. ).build()
  16. env = StreamExecutionEnvironment.get_execution_environment()
  17. env.enable_checkpointing(10)
  18. stream = env.from_source(source, WatermarkStrategy.no_watermarks(), "file-source")

Specific record

Based on the previously defined schema, you can generate classes by leveraging Avro code generation. Once the classes have been generated, there is no need to use the schema directly in your programs. You can either use avro-tools.jar to generate code manually or you could use the Avro Maven plugin to perform code generation on any .avsc files present in the configured source directory. Please refer to Avro Getting Started for more information.

The following example uses the example schema testdata.avsc :

  1. [
  2. {"namespace": "org.apache.flink.formats.parquet.generated",
  3. "type": "record",
  4. "name": "Address",
  5. "fields": [
  6. {"name": "num", "type": "int"},
  7. {"name": "street", "type": "string"},
  8. {"name": "city", "type": "string"},
  9. {"name": "state", "type": "string"},
  10. {"name": "zip", "type": "string"}
  11. ]
  12. }
  13. ]

You will use the Avro Maven plugin to generate the Address Java class:

  1. @org.apache.avro.specific.AvroGenerated
  2. public class Address extends org.apache.avro.specific.SpecificRecordBase implements org.apache.avro.specific.SpecificRecord {
  3. // generated code...
  4. }

You will create an AvroParquetRecordFormat via AvroParquetReaders for Avro Specific record and then create a DataStream containing Parquet records as Avro Specific records.

  1. final FileSource<GenericRecord> source =
  2. FileSource.forRecordStreamFormat(
  3. AvroParquetReaders.forSpecificRecord(Address.class), /* Flink Path */)
  4. .build();
  5. final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  6. env.enableCheckpointing(10L);
  7. final DataStream<GenericRecord> stream =
  8. env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

Reflect record

Beyond Avro Generic and Specific record that requires a predefined Avro schema, Flink also supports creating a DataStream from Parquet files based on existing Java POJO classes. In this case, Avro will use Java reflection to generate schemas and protocols for these POJO classes. Java types are mapped to Avro schemas, please refer to the Avro reflect documentation for more details.

This example uses a simple Java POJO class Datum :

  1. public class Datum implements Serializable {
  2. public String a;
  3. public int b;
  4. public Datum() {}
  5. public Datum(String a, int b) {
  6. this.a = a;
  7. this.b = b;
  8. }
  9. @Override
  10. public boolean equals(Object o) {
  11. if (this == o) {
  12. return true;
  13. }
  14. if (o == null || getClass() != o.getClass()) {
  15. return false;
  16. }
  17. Datum datum = (Datum) o;
  18. return b == datum.b && (a != null ? a.equals(datum.a) : datum.a == null);
  19. }
  20. @Override
  21. public int hashCode() {
  22. int result = a != null ? a.hashCode() : 0;
  23. result = 31 * result + b;
  24. return result;
  25. }
  26. }

You will create an AvroParquetRecordFormat via AvroParquetReaders for Avro Reflect record and then create a DataStream containing Parquet records as Avro Reflect records.

  1. final FileSource<GenericRecord> source =
  2. FileSource.forRecordStreamFormat(
  3. AvroParquetReaders.forReflectRecord(Datum.class), /* Flink Path */)
  4. .build();
  5. final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  6. env.enableCheckpointing(10L);
  7. final DataStream<GenericRecord> stream =
  8. env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

Prerequisite for Parquet files

In order to support reading Avro reflect records, the Parquet file must contain specific meta information. The Avro schema used for creating the Parquet data must contain a namespace, which will be used by the program to identify the concrete Java class for the reflection process.

The following example shows the User schema used previously. But this time it contains a namespace pointing to the location(in this case the package), where the User class for the reflection could be found.

  1. // avro schema with namespace
  2. final String schema =
  3. "{\"type\": \"record\", "
  4. + "\"name\": \"User\", "
  5. + "\"namespace\": \"org.apache.flink.formats.parquet.avro\", "
  6. + "\"fields\": [\n"
  7. + " {\"name\": \"name\", \"type\": \"string\" },\n"
  8. + " {\"name\": \"favoriteNumber\", \"type\": [\"int\", \"null\"] },\n"
  9. + " {\"name\": \"favoriteColor\", \"type\": [\"string\", \"null\"] }\n"
  10. + " ]\n"
  11. + " }";

Parquet files created with this schema will contain meta information like:

  1. creator: parquet-mr version 1.12.2 (build 77e30c8093386ec52c3cfa6c34b7ef3321322c94)
  2. extra: parquet.avro.schema =
  3. {"type":"record","name":"User","namespace":"org.apache.flink.formats.parquet.avro","fields":[{"name":"name","type":"string"},{"name":"favoriteNumber","type":["int","null"]},{"name":"favoriteColor","type":["string","null"]}]}
  4. extra: writer.model.name = avro
  5. file schema: org.apache.flink.formats.parquet.avro.User
  6. --------------------------------------------------------------------------------
  7. name: REQUIRED BINARY L:STRING R:0 D:0
  8. favoriteNumber: OPTIONAL INT32 R:0 D:1
  9. favoriteColor: OPTIONAL BINARY L:STRING R:0 D:1
  10. row group 1: RC:3 TS:143 OFFSET:4
  11. --------------------------------------------------------------------------------
  12. name: BINARY UNCOMPRESSED DO:0 FPO:4 SZ:47/47/1.00 VC:3 ENC:PLAIN,BIT_PACKED ST:[min: Jack, max: Tom, num_nulls: 0]
  13. favoriteNumber: INT32 UNCOMPRESSED DO:0 FPO:51 SZ:41/41/1.00 VC:3 ENC:RLE,PLAIN,BIT_PACKED ST:[min: 1, max: 3, num_nulls: 0]
  14. favoriteColor: BINARY UNCOMPRESSED DO:0 FPO:92 SZ:55/55/1.00 VC:3 ENC:RLE,PLAIN,BIT_PACKED ST:[min: green, max: yellow, num_nulls: 0]

With the User class defined in the package org.apache.flink.formats.parquet.avro:

  1. public class User {
  2. private String name;
  3. private Integer favoriteNumber;
  4. private String favoriteColor;
  5. public User() {}
  6. public User(String name, Integer favoriteNumber, String favoriteColor) {
  7. this.name = name;
  8. this.favoriteNumber = favoriteNumber;
  9. this.favoriteColor = favoriteColor;
  10. }
  11. public String getName() {
  12. return name;
  13. }
  14. public Integer getFavoriteNumber() {
  15. return favoriteNumber;
  16. }
  17. public String getFavoriteColor() {
  18. return favoriteColor;
  19. }
  20. }

you can write the following program to read Avro Reflect records of User type from parquet files:

  1. final FileSource<GenericRecord> source =
  2. FileSource.forRecordStreamFormat(
  3. AvroParquetReaders.forReflectRecord(User.class), /* Flink Path */)
  4. .build();
  5. final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  6. env.enableCheckpointing(10L);
  7. final DataStream<GenericRecord> stream =
  8. env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");