CSV 포맷
CSV 포맷
Flink에서 CSV 파일을 읽고 쓰는 방법을 설명해요. CsvReaderFormat을 사용한 CSV 읽기와 PyFlink에서의 사용법, CSV 스키마를 위한 Jackson 설정을 다룹니다.
출처: CSV format
본문
CSV 포맷을 사용하려면 프로젝트에 Flink CSV 의존성을 추가해야 해요:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-csv</artifactId>
<version>${flink.version}</version>
</dependency>
PyFlink 사용자는 잡에서 직접 사용할 수 있어요.
Flink는 CsvReaderFormat을 사용해 CSV 파일 읽기를 지원해요. 이 리더는 Jackson 라이브러리를 사용하며 CSV 스키마와 파싱 옵션에 해당하는 구성을 전달할 수 있어요.
CsvReaderFormat은 다음과 같이 초기화하고 사용할 수 있어요:
CsvReaderFormat<SomePojo> csvFormat = CsvReaderFormat.forPojo(SomePojo.class);
FileSource<SomePojo> source =
FileSource.forRecordStreamFormat(csvFormat, Path.fromLocalFile(...)).build();
이 경우 CSV 파싱을 위한 스키마는 Jackson 라이브러리를 사용해 SomePojo 클래스의 필드에 기반해 자동으로 유도돼요.
정보: 클래스 정의에
@JsonPropertyOrder({field1, field2, ...})애노테이션을 추가해야 할 수 있어요. 필드 순서가 CSV 파일 컬럼과 정확히 일치해야 해요.
Advanced configuration
CSV 스키마나 파싱 옵션을 더 세밀하게 제어해야 한다면 CsvReaderFormat의 더 저수준인 forSchema 정적 팩토리 메서드를 사용하세요:
CsvReaderFormat<T> forSchema(Supplier<CsvMapper> mapperFactory,
Function<CsvMapper, CsvSchema> schemaGenerator,
TypeInformation<T> typeInformation)
다음은 커스텀 컬럼 구분자를 사용해 POJO를 읽는 예시예요:
//Has to match the exact order of columns in the CSV file
@JsonPropertyOrder({"city","lat","lng","country","iso2",
"adminName","capital","population"})
public static class CityPojo {
public String city;
public BigDecimal lat;
public BigDecimal lng;
public String country;
public String iso2;
public String adminName;
public String capital;
public long population;
}
Function<CsvMapper, CsvSchema> schemaGenerator = mapper ->
mapper.schemaFor(CityPojo.class).withoutQuoteChar().withColumnSeparator('|');
CsvReaderFormat<CityPojo> csvFormat =
CsvReaderFormat.forSchema(() -> new CsvMapper(), schemaGenerator, TypeInformation.of(CityPojo.class));
FileSource<CityPojo> source =
FileSource.forRecordStreamFormat(csvFormat, Path.fromLocalFile(...)).build();
해당 CSV 파일:
Berlin|52.5167|13.3833|Germany|DE|Berlin|primary|3644826
San Francisco|37.7562|-122.443|United States|US|California||3592294
Beijing|39.905|116.3914|China|CN|Beijing|primary|19433000
세밀한 Jackson 설정으로 더 복잡한 데이터 타입을 읽는 것도 가능해요:
public static class ComplexPojo {
private long id;
private int[] array;
}
CsvReaderFormat<ComplexPojo> csvFormat =
CsvReaderFormat.forSchema(
CsvSchema.builder()
.addColumn(
new CsvSchema.Column(0, "id", CsvSchema.ColumnType.NUMBER))
.addColumn(
new CsvSchema.Column(4, "array", CsvSchema.ColumnType.ARRAY)
.withArrayElementSeparator("#"))
.build(),
TypeInformation.of(ComplexPojo.class));
PyFlink 사용자의 경우 컬럼을 수동으로 추가해 csv 스키마를 정의할 수 있고, csv 소스의 출력 타입은 각 컬럼이 필드에 매핑된 Row가 돼요.
schema = CsvSchema.builder() \
.add_number_column('id', number_type=DataTypes.BIGINT()) \
.add_array_column('array', separator='#', element_type=DataTypes.INT()) \
.set_column_separator(',') \
.build()
source = FileSource.for_record_stream_format(
CsvReaderFormat.for_schema(schema), CSV_FILE_PATH).build()
# the type of record will be Types.ROW_NAMED(['id', 'array'], [Types.LONG(), Types.LIST(Types.INT())])
ds = env.from_source(source, WatermarkStrategy.no_watermarks(), 'csv-source')
해당 CSV 파일:
0,1#2#3
1,
2,1
TextLineInputFormat와 유사하게 CsvReaderFormat은 연속(continues)·배치(batch) 두 모드 모두에서 사용할 수 있어요(TextLineInputFormat 예제 참고).
PyFlink 사용자의 경우 CsvBulkWriters로 BulkWriterFactory를 만들어 레코드를 CSV 포맷으로 파일에 기록할 수 있어요.
schema = CsvSchema.builder() \
.add_number_column('id', number_type=DataTypes.BIGINT()) \
.add_array_column('array', separator='#', element_type=DataTypes.INT()) \
.set_column_separator(',') \
.build()
sink = FileSink.for_bulk_format(
OUTPUT_DIR, CsvBulkWriters.for_schema(schema)).build()
ds.sink_to(sink)