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)를 추가로 호출해야 해요.
Flink RowData
이 예시에서는 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");