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 사용자의 경우 CsvBulkWritersBulkWriterFactory를 만들어 레코드를 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)

더 알아보기 (Learn more)