Parquet 포맷

Parquet 포맷 (Parquet format)

Flink는 Parquet 파일을 읽어 Flink RowData를 생성하고 Avro 레코드를 생성하는 것을 지원해요. 배치와 스트리밍 실행 모드 모두에서 사용할 수 있는 새 Source와 호환돼요.

출처: 문서

본문

포맷을 사용하려면 프로젝트에 flink-parquet 의존성을 추가해야 해요.

	org.apache.flink
	flink-parquet
	2.3.0

Avro 레코드를 읽으려면 parquet-avro 의존성을 추가해야 해요.

    org.apache.parquet
    parquet-avro
    1.12.2
    true
    
        
            org.apache.hadoop
            hadoop-client
        
        
            it.unimi.dsi
            fastutil
        
    

PyFlink job에서 Parquet 포맷을 사용하려면 다음 의존성이 필요해요. PyFlink에서 JAR을 사용하는 방법에 대한 자세한 내용은 Python dependency management를 참고하세요.

이 포맷은 배치 및 스트리밍 실행 모드 모두에서 사용할 수 있는 새 Source와 호환돼요. 따라서 이 포맷을 두 종류의 데이터에 사용할 수 있어요.

  • Bounded data (유한 데이터): 모든 파일을 나열하고 모두 읽어요.
  • Unbounded data (무한 데이터): 디렉터리에 새 파일이 나타나는지 모니터링해요.

벡터화 리더 (Vectorized reader)

// Parquet rows are decoded in batches
FileSource.forBulkFileFormat(BulkFormat,Path...)

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

# Monitor the Paths to read data as unbounded data
FileSource.for_bulk_file_format(BulkFormat, Path...) \
          .monitor_continuously(Duration.of_millis(5)) \
          .build()

Avro Parquet 리더 (Avro Parquet reader)

// Parquet rows are decoded in batches
FileSource.forRecordStreamFormat(StreamFormat,Path...)

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

# Monitor the Paths to read data as unbounded data
FileSource.for_record_stream_format(StreamFormat, Path...) \
          .monitor_continuously(Duration.of_millis(5)) \
          .build()

다음 예시들은 모두 유한(bounded) 데이터용으로 구성되었어요. File Source를 무한 데이터용으로 구성하려면 AbstractFileSource.AbstractFileSourceBuilder.monitorContinuously(Duration)를 추가로 호출해야 해요.

이 예시에서는 Parquet 레코드를 Flink RowData로 포함하는 DataStream을 만들 거예요. 스키마는 지정된 필드("f7", "f4", "f99")만 읽도록 프로젝션돼요.

Flink는 500개 레코드 단위로 읽어요. 첫 번째 boolean 매개변수는 타임스탬프 컬럼이 UTC로 해석될 것임을 지정해요. 두 번째 boolean은 프로젝션된 Parquet 필드 이름이 대소문자를 구분함을 애플리케이션에 지시해요. 레코드가 이벤트 타임스탬프를 포함하지 않으므로 워터마크 전략이 정의되어 있지 않아요.

final LogicalType[] fieldTypes =
        new LogicalType[] {
                new DoubleType(), new IntType(), new VarCharType()};
final RowType rowType = RowType.of(fieldTypes, new String[] {"f7", "f4", "f99"});

final ParquetColumnarRowInputFormat<FileSourceSplit> format =
        new ParquetColumnarRowInputFormat<>(
                new Configuration(),
                rowType,
                InternalTypeInfo.of(rowType),
                500,
                false,
                true);
final FileSource<RowData> source =
        FileSource.forBulkFileFormat(format,  /* Flink Path */)
                .build();
final DataStream<RowData> stream =
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");
row_type = DataTypes.ROW([
    DataTypes.FIELD('f7', DataTypes.DOUBLE()),
    DataTypes.FIELD('f4', DataTypes.INT()),
    DataTypes.FIELD('f99', DataTypes.VARCHAR()),
])
source = FileSource.for_bulk_file_format(ParquetColumnarRowInputFormat(
    row_type=row_type,
    hadoop_config=Configuration(),
    batch_size=500,
    is_utc_timestamp=False,
    is_case_sensitive=True,
), PARQUET_FILE_PATH).build()
ds = env.from_source(source, WatermarkStrategy.no_watermarks(), "file-source")

Avro 레코드 (Avro Records)

Flink는 Parquet 파일을 읽어 세 가지 유형의 Avro 레코드를 생성하는 것을 지원해요 (PyFlink에서는 Generic record만 지원돼요).

  • Generic record
  • Specific record
  • Reflect record

Generic record

Avro 스키마는 JSON을 사용해 정의돼요. Avro 스키마와 유형에 대한 더 많은 정보는 Avro specification에서 얻을 수 있어요. 이 예시는 공식 Avro 튜토리얼에 설명된 것과 유사한 Avro 스키마 예시를 사용해요.

{"namespace": "example.avro",
 "type": "record",
 "name": "User",
 "fields": [
     {"name": "name", "type": "string"},
     {"name": "favoriteNumber",  "type": ["int", "null"]},
     {"name": "favoriteColor", "type": ["string", "null"]}
 ]
}

이 스키마는 name, favoriteNumber, favoriteColor 세 필드를 가진 사용자를 나타내는 레코드를 정의해요. Avro 스키마를 정의하는 방법에 대한 더 자세한 내용은 record specification에서 찾을 수 있어요.

다음 예시에서는 Parquet 레코드를 Avro Generic 레코드로 포함하는 DataStream을 만들 거예요. JSON 문자열을 기반으로 Avro 스키마를 파싱할 거예요. 스키마를 파싱하는 방법에는 java.io.File이나 java.io.InputStream 등 여러 가지가 있어요. 자세한 내용은 Avro Schema를 참고하세요. 그다음 Avro Generic 레코드를 위해 AvroParquetReaders를 통해 AvroParquetRecordFormat을 만들 거예요.

// parsing avro schema
final Schema schema =
        new Schema.Parser()
            .parse(
                    "{\"type\": \"record\", "
                        + "\"name\": \"User\", "
                        + "\"fields\": [\n"
                        + "        {\"name\": \"name\", \"type\": \"string\" },\n"
                        + "        {\"name\": \"favoriteNumber\",  \"type\": [\"int\", \"null\"] },\n"
                        + "        {\"name\": \"favoriteColor\", \"type\": [\"string\", \"null\"] }\n"
                        + "    ]\n"
                        + "    }");

final FileSource<GenericRecord> source =
        FileSource.forRecordStreamFormat(
                AvroParquetReaders.forGenericRecord(schema), /* Flink Path */)
        .build();

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(10L);
        
final DataStream<GenericRecord> stream =
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");
# parsing avro schema
schema = AvroSchema.parse_string("""
{
    "type": "record",
    "name": "User",
    "fields": [
        {"name": "name", "type": "string"},
        {"name": "favoriteNumber",  "type": ["int", "null"]},
        {"name": "favoriteColor", "type": ["string", "null"]}
    ]
}
""")

source = FileSource.for_record_stream_format(
    AvroParquetReaders.for_generic_record(schema), # file paths
).build()

env = StreamExecutionEnvironment.get_execution_environment()
env.enable_checkpointing(10)

stream = env.from_source(source, WatermarkStrategy.no_watermarks(), "file-source")

Specific record

이전에 정의한 스키마를 기반으로 Avro 코드 생성을 활용해 클래스를 생성할 수 있어요. 클래스가 생성되면 프로그램에서 스키마를 직접 사용할 필요가 없어요.

avro-tools.jar를 사용해 수동으로 코드를 생성하거나, Avro Maven 플러그인을 사용해 구성된 소스 디렉터리에 있는 .avsc 파일에 대해 코드 생성을 수행할 수 있어요. 더 많은 정보는 Avro Getting Started를 참고하세요.

다음 예시는 예시 스키마 testdata.avsc를 사용해요.

[
  {"namespace": "org.apache.flink.formats.parquet.generated",
    "type": "record",
    "name": "Address",
    "fields": [
      {"name": "num", "type": "int"},
      {"name": "street", "type": "string"},
      {"name": "city", "type": "string"},
      {"name": "state", "type": "string"},
      {"name": "zip", "type": "string"}
    ]
  }
]

Avro Maven 플러그인을 사용해 Address Java 클래스를 생성할 거예요.

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

Avro Specific 레코드를 위해 AvroParquetReaders를 통해 AvroParquetRecordFormat을 만들고, Avro Specific 레코드로 Parquet 레코드를 포함하는 DataStream을 만들 거예요.

final FileSource<GenericRecord> source =
        FileSource.forRecordStreamFormat(
                AvroParquetReaders.forSpecificRecord(Address.class), /* Flink Path */)
        .build();

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(10L);
        
final DataStream<GenericRecord> stream =
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

Reflect record

미리 정의된 Avro 스키마가 필요한 Avro Generic 및 Specific 레코드 외에도, Flink는 기존 Java POJO 클래스를 기반으로 Parquet 파일에서 DataStream을 생성하는 것도 지원해요. 이 경우 Avro는 Java 리플렉션을 사용해 이러한 POJO 클래스에 대한 스키마와 프로토콜을 생성해요. Java 유형은 Avro 스키마에 매핑되는데, 자세한 내용은 Avro reflect documentation을 참고하세요.

이 예시는 간단한 Java POJO 클래스 Datum을 사용해요.

public class Datum implements Serializable {

    public String a;
    public int b;

    public Datum() {}

    public Datum(String a, int b) {
        this.a = a;
        this.b = b;
    }

    @Override
    public boolean equals(Object o) {
        if (this == o) {
            return true;
        }
        if (o == null || getClass() != o.getClass()) {
            return false;
        }

        Datum datum = (Datum) o;
        return b == datum.b && (a != null ? a.equals(datum.a) : datum.a == null);
    }

    @Override
    public int hashCode() {
        int result = a != null ? a.hashCode() : 0;
        result = 31 * result + b;
        return result;
    }
}

Avro Reflect 레코드를 위해 AvroParquetReaders를 통해 AvroParquetRecordFormat을 만들고, Avro Reflect 레코드로 Parquet 레코드를 포함하는 DataStream을 만들 거예요.

final FileSource<GenericRecord> source =
        FileSource.forRecordStreamFormat(
                AvroParquetReaders.forReflectRecord(Datum.class), /* Flink Path */)
        .build();

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(10L);
        
final DataStream<GenericRecord> stream =
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

Parquet 파일의 사전 요구사항 (Prerequisite for Parquet files)

Avro reflect 레코드 읽기를 지원하려면 Parquet 파일이 특정 메타 정보를 포함해야 해요. Parquet 데이터를 만드는 데 사용된 Avro 스키마는 namespace를 포함해야 하는데, 프로그램이 리플렉션 프로세스를 위해 구체적인 Java 클래스를 식별하는 데 이 namespace를 사용해요.

다음 예시는 이전에 사용한 User 스키마를 보여줘요. 하지만 이번에는 리플렉션에 사용될 User 클래스를 찾을 수 있는 위치(이 경우 패키지)를 가리키는 namespace를 포함해요.

// avro schema with namespace
final String schema = 
                    "{\"type\": \"record\", "
                        + "\"name\": \"User\", "
                        + "\"namespace\": \"org.apache.flink.formats.parquet.avro\", "
                        + "\"fields\": [\n"
                        + "        {\"name\": \"name\", \"type\": \"string\" },\n"
                        + "        {\"name\": \"favoriteNumber\",  \"type\": [\"int\", \"null\"] },\n"
                        + "        {\"name\": \"favoriteColor\", \"type\": [\"string\", \"null\"] }\n"
                        + "    ]\n"
                        + "    }";

이 스키마로 생성된 Parquet 파일은 다음과 같은 메타 정보를 포함할 거예요.

creator:        parquet-mr version 1.12.2 (build 77e30c8093386ec52c3cfa6c34b7ef3321322c94)
extra:          parquet.avro.schema =
{"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"]}]}
extra:          writer.model.name = avro

file schema:    org.apache.flink.formats.parquet.avro.User
--------------------------------------------------------------------------------
name:           REQUIRED BINARY L:STRING R:0 D:0
favoriteNumber: OPTIONAL INT32 R:0 D:1
favoriteColor:  OPTIONAL BINARY L:STRING R:0 D:1

row group 1:    RC:3 TS:143 OFFSET:4
--------------------------------------------------------------------------------
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]
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]
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]

org.apache.flink.formats.parquet.avro 패키지에 User 클래스가 정의되어 있다면,

public class User {
        private String name;
        private Integer favoriteNumber;
        private String favoriteColor;

        public User() {}

        public User(String name, Integer favoriteNumber, String favoriteColor) {
            this.name = name;
            this.favoriteNumber = favoriteNumber;
            this.favoriteColor = favoriteColor;
        }

        public String getName() {
            return name;
        }

        public Integer getFavoriteNumber() {
            return favoriteNumber;
        }

        public String getFavoriteColor() {
            return favoriteColor;
        }
    }

다음 프로그램을 작성해 parquet 파일에서 User 유형의 Avro Reflect 레코드를 읽을 수 있어요.

final FileSource<GenericRecord> source =
        FileSource.forRecordStreamFormat(
        AvroParquetReaders.forReflectRecord(User.class), /* Flink Path */)
        .build();

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(10L);

final DataStream<GenericRecord> stream =
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");

더 알아보기 (Learn more)