FileSystem

FileSystem

이 커넥터는 Flink FileSystem 추상화가 지원하는 파일시스템에 (파티션된) 파일을 읽거나 쓰는 BATCHSTREAMING을 위한 통합 Source와 Sink를 제공해요. 이 파일시스템 커넥터는 BATCHSTREAMING 모두에 대해 동일한 보장을 제공하며, STREAMING 실행에 대해 정확히 한 번(exactly-once) 의미론을 제공하도록 설계됐어요.

출처: 문서

본문

이 커넥터는 포맷(format)(예: Avro, CSV, Parquet)과 함께 모든 (분산) 파일시스템(예: POSIX, S3, HDFS)에서 파일 집합을 읽고 쓸 수 있게 해주며, 스트림 또는 레코드를 생성해요.

File Source

File SourceSource API를 기반으로 하는 통합 데이터 소스로, 배치와 스트리밍 모드 모두에서 파일을 읽어요. SplitEnumeratorSourceReader의 두 부분으로 나뉘어요.

  • SplitEnumerator는 읽을 파일을 발견하고 식별하며, 이를 SourceReader에 할당하는 책임을 가져요.
  • SourceReader는 처리해야 할 파일을 요청하고 파일시스템에서 파일을 읽어요.

File Source를 포맷과 결합해야 하며, 이를 통해 CSV를 파싱하고, AVRO를 디코딩하고, Parquet 컬럼형 파일을 읽을 수 있어요.

유한 및 무한 스트림 (Bounded and Unbounded Streams)

유한(bounded) File Source는 모든 파일을 나열하고(SplitEnumerator를 통해 — 숨겨진 파일을 걸러낸 재귀적 디렉터리 나열) 모두 읽어요.

무한(unbounded) File Source는 열거자(enumerator)를 주기적 파일 발견으로 구성할 때 만들어져요. 이 경우 SplitEnumerator는 유한 경우처럼 열거하지만, 일정 간격 후 열거를 반복해요. 반복 열거마다 SplitEnumerator는 이전에 감지된 파일을 걸러내고 새 파일만 SourceReader로 보내요.

사용법 (Usage)

다음 API 호출 중 하나로 File Source 구축을 시작할 수 있어요:

// reads the contents of a file from a file stream. 
FileSource.forRecordStreamFormat(StreamFormat,Path...);
        
// reads batches of records from a file at a time
FileSource.forBulkFileFormat(BulkFormat,Path...);
# reads the contents of a file from a file stream.
FileSource.for_record_stream_format(stream_format, *path)

# reads batches of records from a file at a time
FileSource.for_bulk_file_format(bulk_format, *path)

이것은 File Source의 모든 속성을 구성할 수 있는 FileSource.FileSourceBuilder를 만들어요.

유한/배치의 경우 File Source는 주어진 경로 아래의 모든 파일을 처리해요. 연속/스트리밍의 경우 소스는 새 파일을 위해 경로를 주기적으로 확인하고 그것들을 읽기 시작해요.

File Source 생성을 시작할 때(위 메서드 중 하나로 만들어진 FileSource.FileSourceBuilder를 통해), 소스는 기본적으로 유한/배치 모드예요. AbstractFileSource.AbstractFileSourceBuilder.monitorContinuously(Duration)을 호출해 소스를 연속 스트리밍 모드로 전환할 수 있어요.

final FileSource<String> source =
        FileSource.forRecordStreamFormat(...)
        .monitorContinuously(Duration.ofMillis(5))  
        .build();
source = FileSource.for_record_stream_format(...) \
    .monitor_continously(Duration.of_millis(5)) \
    .build()

포맷 타입 (Format Types)

각 파일의 읽기는 파일 포맷에 의해 정의된 파일 리더를 통해 일어나요. 이들은 파일 내용에 대한 파싱 로직을 정의해요. 소스가 지원하는 클래스가 여러 개 있어요. 인터페이스는 구현의 단순성과 유연성/효율성 사이의 트레이드오프예요.

  • StreamFormat은 파일 스트림에서 파일의 내용을 읽어요. 구현하기 가장 간단한 포맷이며, 즉시 사용할 수 있는 많은 기능(체크포인팅 로직 같은)을 제공하지만, 적용할 수 있는 최적화(객체 재사용, 배칭 등)가 제한돼요.
  • BulkFormat은 한 번에 파일에서 레코드 배치를 읽어요. 구현하기 가장 "저수준" 포맷이지만, 구현을 최적화할 수 있는 가장 큰 유연성을 제공해요.
TextLine 포맷

StreamFormat 리더는 파일에서 텍스트 라인을 포맷해요. 리더는 Java의 내장 InputStreamReader를 사용해 다양한 지원 문자셋 인코딩으로 바이트 스트림을 디코딩해요. 이 포맷은 체크포인트에서의 최적화된 복구를 지원하지 않아요. 복구 시 마지막 체크포인트 이전에 처리된 라인 수를 다시 읽고 버려요. 이것은 라인의 오프셋이 스트림 입력의 내부 버퍼링과 문자셋 디코더 상태를 가진 문자셋 디코더를 통해 추적될 수 없기 때문이에요.

SimpleStreamFormat 추상 클래스

이것은 분할할 수 없는(splittable) 포맷을 위한 StreamFormat의 간단한 버전이에요. SimpleStreamFormat을 구현해 Array 또는 File의 커스텀 읽기를 할 수 있어요:

private static final class ArrayReaderFormat extends SimpleStreamFormat<byte[]> {
    private static final long serialVersionUID = 1L;

    @Override
    public Reader<byte[]> createReader(Configuration config, FSDataInputStream stream)
            throws IOException {
        return new ArrayReader(stream);
    }

    @Override
    public TypeInformation<byte[]> getProducedType() {
        return PrimitiveArrayTypeInfo.BYTE_PRIMITIVE_ARRAY_TYPE_INFO;
    }
}

final FileSource<byte[]> source =
                FileSource.forRecordStreamFormat(new ArrayReaderFormat(), path).build();

SimpleStreamFormat의 예로는 CsvReaderFormat이 있어요. 다음과 같이 초기화할 수 있어요:

CsvReaderFormat<SomePojo> csvFormat = CsvReaderFormat.forPojo(SomePojo.class);
FileSource<SomePojo> source = 
        FileSource.forRecordStreamFormat(csvFormat, Path.fromLocalFile(...)).build();

이 경우 CSV 파싱을 위한 스키마는 Jackson 라이브러리를 사용해 SomePojo 클래스의 필드를 기반으로 자동으로 파생돼요. (참고: CSV 파일 컬럼과 정확히 일치하는 필드 순서로 클래스 정의에 @JsonPropertyOrder({field1, field2, ...}) 주석을 추가해야 할 수 있어요.)

CSV 스키마나 파싱 옵션에 대해 더 세밀한 제어가 필요하면, CsvReaderFormat의 더 저수준인 forSchema 정적 팩토리 메서드를 사용해요:

CsvReaderFormat<T> forSchema(Supplier<CsvMapper> mapperFactory, 
                             Function<CsvMapper, CsvSchema> schemaGenerator, 
                             TypeInformation<T> typeInformation)
Bulk 포맷

BulkFormat은 한 번에 레코드 배치를 읽고 디코딩해요. bulk 포맷의 예로는 ORC나 Parquet 같은 포맷이 있어요. 외부 BulkFormat 클래스는 주로 리더의 구성 홀더이자 팩토리 역할을 해요. 실제 읽기는 BulkFormat#createReader(Configuration, FileSourceSplit) 메서드에서 생성되는 BulkFormat.Reader에 의해 수행돼요. 체크포인트된 스트리밍 실행 중 체크포인트를 기반으로 bulk 리더가 생성되면, 리더는 BulkFormat#restoreReader(Configuration, FileSourceSplit) 메서드에서 다시 생성돼요.

SimpleStreamFormatStreamFormatAdapter로 감싸 BulkFormat으로 변환될 수 있어요:

BulkFormat<SomePojo, FileSourceSplit> bulkFormat = 
        new StreamFormatAdapter<>(CsvReaderFormat.forPojo(SomePojo.class));

파일 열거 커스터마이징 (Customizing File Enumeration)

/**
 * A FileEnumerator implementation for hive source, which generates splits based on 
 * HiveTablePartition.
 */
public class HiveSourceFileEnumerator implements FileEnumerator {
    
    // reference constructor
    public HiveSourceFileEnumerator(...) {
        ...
    }

    /***
     * Generates all file splits for the relevant files under the given paths. The {@code
     * minDesiredSplits} is an optional hint indicating how many splits would be necessary to
     * exploit parallelism properly.
     */
    @Override
    public Collection<FileSourceSplit> enumerateSplits(Path[] paths, int minDesiredSplits)
            throws IOException {
        // createInputSplits:splitting files into fragmented collections
        return new ArrayList<>(createInputSplits(...));
    }

    ...

    /***
     * A factory to create HiveSourceFileEnumerator.
     */
    public static class Provider implements FileEnumerator.Provider {

        ...
        @Override
        public FileEnumerator create() {
            return new HiveSourceFileEnumerator(...);
        }
    }
}
// use the customizing file enumeration
new HiveSource<>(
        ...,
        new HiveSourceFileEnumerator.Provider(
        partitions != null ? partitions : Collections.emptyList(),
        new JobConfWrapper(jobConf)),
       ...);

현재 제한 사항 (Current Limitations)

워터마킹은 큰 파일 백로그에 대해 잘 작동하지 않아요. 워터마크가 파일 내에서 열심히(eagerly) 진행되고, 다음 파일이 워터마크보다 늦은 데이터를 포함할 수 있기 때문이에요.

무한 File Source의 경우 열거자는 현재 이미 처리된 모든 파일의 경로를 기억하는데, 이는 어떤 경우에는 꽤 커질 수 있는 상태예요. 미래에는 이미 처리된 파일 추적의 압축된 형태를 추가할 계획이 있어요(예: 수정 타임스탬프를 경계 아래로 유지).

뒤에서 일어나는 일 (Behind the Scenes)

새 데이터 소스 API 설계를 통해 File Source가 어떻게 동작하는지에 관심이 있다면 이 부분을 참조로 읽으면 돼요. 새 데이터 소스 API에 대한 자세한 내용은 데이터 소스 문서FLIP-27를 확인하세요.

File Sink

파일 싱크는 들어오는 데이터를 버킷(bucket)으로 써요. 들어오는 스트림이 무한할 수 있으므로, 각 버킷의 데이터는 유한 크기의 파트 파일(part files)로 구성돼요. 버케팅 동작은 완전히 구성 가능하며, 기본적으로 시간 기반 버케팅으로 매시간 새 버킷 작성을 시작해요. 즉 각 결과 버킷은 스트림에서 1시간 간격으로 수신된 레코드를 담은 파일을 포함하게 돼요.

버킷 디렉터리 안의 데이터는 파트 파일로 분할돼요. 각 버킷은 해당 버킷의 데이터를 받은 싱크의 각 서브태스크에 대해 최소 하나의 파트 파일을 포함해요. 추가 파트 파일은 구성 가능한 롤링 정책(rolling policy)에 따라 생성돼요. 행 인코딩 포맷(Row-encoded Formats)(see File Formats)의 경우 기본 정책은 크기, 파일이 열릴 수 있는 최대 기간을 지정하는 타임아웃, 파일이 닫히는 최대 비활성(inactivity) 타임아웃에 따라 파트 파일을 롤해요. Bulk 인코딩 포맷(Bulk-encoded Formats)의 경우 매 체크포인트마다 롤하며, 사용자는 크기나 시간에 기반한 추가 조건을 지정할 수 있어요.

중요: STREAMING 모드에서 FileSink를 사용할 때는 체크포인팅을 활성화해야 해요. 파트 파일은 성공적인 체크포인트에서만 최종화(finalized)될 수 있어요. 체크포인팅이 비활성화되면 파트 파일은 영원히 in-progress 또는 pending 상태로 남게 되며, 다운스트림 시스템이 안전하게 읽을 수 없어요.

streamfilesink_bucketing

포맷 타입 (Format Types)

FileSinkApache Parquet 같은 행 및 bulk 인코딩 포맷을 모두 지원해요. 이 두 변형은 다음 정적 메서드로 만들 수 있는 각각의 빌더를 제공해요:

  • 행 인코딩 싱크: FileSink.forRowFormat(basePath, rowEncoder)
  • Bulk 인코딩 싱크: FileSink.forBulkFormat(basePath, bulkWriterFactory)

행 또는 bulk 인코딩 싱크를 만들 때 버킷이 저장될 기본 경로와 데이터에 대한 인코딩 로직을 지정해야 해요.

모든 구성 옵션과 다양한 데이터 포맷의 구현에 대한 더 많은 문서는 FileSink의 JavaDoc을 확인하세요.

행 인코딩 포맷 (Row-encoded Formats)

행 인코딩 포맷은 in-progress 파트 파일의 OutputStream으로 개별 행을 직렬화하는 데 사용되는 Encoder를 지정해야 해요.

버킷 할당기(bucket assigner)에 더해, RowFormatBuilder는 사용자가 다음을 지정할 수 있게 해줘요:

  • Custom RollingPolicy : DefaultRollingPolicy를 재정의하는 롤링 정책
  • bucketCheckInterval (기본값 = 1 min) : 시간 기반 롤링 정책을 확인하는 간격

String 요소를 쓰는 기본 사용법은 다음과 같아요:

import org.apache.flink.api.common.serialization.SimpleStringEncoder;
import org.apache.flink.core.fs.Path;
import org.apache.flink.configuration.MemorySize;
import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.streaming.api.functions.sink.filesystem.rollingpolicies.DefaultRollingPolicy;

import java.time.Duration;

DataStream<String> input = ...;

final FileSink<String> sink = FileSink
    .forRowFormat(new Path(outputPath), new SimpleStringEncoder<String>("UTF-8"))
    .withRollingPolicy(
        DefaultRollingPolicy.builder()
            .withRolloverInterval(Duration.ofMinutes(15))
            .withInactivityInterval(Duration.ofMinutes(5))
            .withMaxPartSize(MemorySize.ofMebiBytes(1024))
            .build())
    .build();

input.sinkTo(sink);
data_stream = ...

sink = FileSink \
    .for_row_format(OUTPUT_PATH, Encoder.simple_string_encoder("UTF-8")) \
    .with_rolling_policy(RollingPolicy.default_rolling_policy(
        part_size=1024 ** 3, rollover_interval=15 * 60 * 1000, inactivity_interval=5 * 60 * 1000)) \
    .build()

data_stream.sink_to(sink)

이 예시는 레코드를 기본 한 시간 시간 버킷에 할당하는 간단한 싱크를 만들어요. 또한 다음 세 조건 중 하나에서 in-progress 파트 파일을 롤하는 롤링 정책을 지정해요:

  • 최소 15분 분량의 데이터를 포함
  • 지난 5분 동안 새 레코드를 받지 않음
  • 파일 크기가 1GB에 도달(마지막 레코드 작성 후)
Bulk 인코딩 포맷 (Bulk-encoded Formats)

Bulk 인코딩 싱크는 행 인코딩 싱크와 유사하게 생성되지만, Encoder를 지정하는 대신 BulkWriter.Factory를 지정해야 해요. BulkWriter 로직은 새 요소가 어떻게 추가되고 플러시되는지, 그리고 추가 인코딩을 위해 레코드 배치가 어떻게 최종화되는지 정의해요.

Flink는 다섯 가지 내장 BulkWriter 팩토리와 함께 제공돼요:

  • ParquetWriterFactory
  • AvroWriterFactory
  • SequenceFileWriterFactory
  • CompressWriterFactory
  • OrcBulkWriterFactory

중요 Bulk 포맷은 CheckpointRollingPolicy를 확장하는 롤링 정책만 가질 수 있어요. 후자는 매 체크포인트마다 롤해요. 정책은 크기나 프로세싱 타임에 기반해 추가로 롤할 수 있어요.

Parquet 포맷

Flink에는 Avro 데이터에 대한 Parquet 라이터 팩토리를 만드는 내장 편의 메서드가 포함돼요. 이 메서드와 관련 문서는 AvroParquetWriters 클래스에서 찾을 수 있어요.

다른 Parquet 호환 데이터 포맷으로 쓰려면, 사용자는 ParquetBuilder 인터페이스의 커스텀 구현으로 ParquetWriterFactory를 만들어야 해요.

애플리케이션에서 Parquet bulk 인코더를 사용하려면 다음 의존성을 추가해야 해요:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-parquet</artifactId>
    <version>2.3.0</version>
</dependency>

PyFlink 작업에서 Parquet 포맷을 사용하려면 다음 의존성이 필요해요: flink-sql-parquet 2.3.0 (Download). PyFlink에서 JAR을 사용하는 방법에 대한 자세한 내용은 Python dependency management를 참조하세요.

Avro 데이터를 Parquet 포맷으로 쓰는 FileSink는 다음과 같이 만들 수 있어요:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.formats.parquet.avro.AvroParquetWriters;
import org.apache.avro.Schema;

Schema schema = ...;
DataStream<GenericRecord> input = ...;

final FileSink<GenericRecord> sink = FileSink
    .forBulkFormat(outputBasePath, AvroParquetWriters.forGenericRecord(schema))
    .build();

input.sinkTo(sink);
schema = AvroSchema.parse_string(JSON_SCHEMA)
# The element could be vanilla Python data structure matching the schema,
# which is annotated with default Types.PICKLED_BYTE_ARRAY()
data_stream = ...

avro_type_info = GenericRecordAvroTypeInfo(schema)
sink = FileSink \
    .for_bulk_format(OUTPUT_BASE_PATH, AvroParquetWriters.for_generic_record(schema)) \
    .build()

# A map to indicate its Avro type info is necessary for serialization
data_stream.map(lambda e: e, output_type=avro_type_info).sink_to(sink)

유사하게, Protobuf 데이터를 Parquet 포맷으로 쓰는 FileSink는 다음과 같이 만들 수 있어요:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.formats.parquet.protobuf.ParquetProtoWriters;

// ProtoRecord is a generated protobuf Message class.
DataStream<ProtoRecord> input = ...;

final FileSink<ProtoRecord> sink = FileSink
    .forBulkFormat(outputBasePath, ParquetProtoWriters.forType(ProtoRecord.class))
    .build();

input.sinkTo(sink);

PyFlink 사용자는 ParquetBulkWriters를 사용해 Row들을 Parquet 파일로 쓰는 BulkWriterFactory를 만들 수 있어요.

row_type = DataTypes.ROW([
    DataTypes.FIELD('string', DataTypes.STRING()),
    DataTypes.FIELD('int_array', DataTypes.ARRAY(DataTypes.INT()))
])

sink = FileSink.for_bulk_format(
    OUTPUT_DIR, ParquetBulkWriters.for_row_type(
        row_type,
        hadoop_config=Configuration(),
        utc_timestamp=True,
    )
).build()

ds.sink_to(sink)
Avro 포맷

Flink는 또한 Avro 파일로 데이터를 쓰는 내장 지원을 제공해요. Avro 라이터 팩토리를 만드는 편의 메서드 목록과 관련 문서는 AvroWriters 클래스에서 찾을 수 있어요.

애플리케이션에서 Avro 라이터를 사용하려면 다음 의존성을 추가해야 해요:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-avro</artifactId>
    <version>2.3.0</version>
</dependency>

PyFlink 작업에서 Avro 포맷을 사용하려면 다음 의존성이 필요해요: flink-sql-avro 2.3.0 (Download).

데이터를 Avro 파일로 쓰는 FileSink는 다음과 같이 만들 수 있어요:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.formats.avro.AvroWriters;
import org.apache.avro.Schema;

Schema schema = ...;
DataStream<GenericRecord> input = ...;

final FileSink<GenericRecord> sink = FileSink
    .forBulkFormat(outputBasePath, AvroWriters.forGenericRecord(schema))
    .build();

input.sinkTo(sink);
schema = AvroSchema.parse_string(JSON_SCHEMA)
# The element could be vanilla Python data structure matching the schema,
# which is annotated with default Types.PICKLED_BYTE_ARRAY()
data_stream = ...

avro_type_info = GenericRecordAvroTypeInfo(schema)
sink = FileSink \
    .for_bulk_format(OUTPUT_BASE_PATH, AvroBulkWriters.for_generic_record(schema)) \
    .build()

# A map to indicate its Avro type info is necessary for serialization
data_stream.map(lambda e: e, output_type=avro_type_info).sink_to(sink)

커스텀 Avro 라이터(예: 압축 활성화)를 만들려면 사용자는 AvroBuilder 인터페이스의 커스텀 구현으로 AvroWriterFactory를 만들어야 해요:

AvroWriterFactory factory = new AvroWriterFactory<>((AvroBuilder<Address>) out -> {
    Schema schema = ReflectData.get().getSchema(Address.class);
    DatumWriter<Address> datumWriter = new ReflectDatumWriter<>(schema);

    DataFileWriter<Address> dataFileWriter = new DataFileWriter<>(datumWriter);
    dataFileWriter.setCodec(CodecFactory.snappyCodec());
    dataFileWriter.create(schema, out);
    return dataFileWriter;
});

DataStream<Address> stream = ...
stream.sinkTo(FileSink.forBulkFormat(
    outputBasePath,
    factory).build());
ORC 포맷

데이터를 ORC 포맷으로 bulk 인코딩할 수 있게 하기 위해, Flink는 구체적인 Vectorizer 구현을 받는 OrcBulkWriterFactory를 제공해요.

데이터를 bulk 방식으로 인코딩하는 다른 컬럼형 포맷처럼, Flink의 OrcBulkWriter는 입력 요소를 배치로 써요. 이것은 ORC의 VectorizedRowBatch를 사용해 달성돼요.

입력 요소가 VectorizedRowBatch로 변환되어야 하므로, 사용자는 추상 Vectorizer 클래스를 확장하고 vectorize(T element, VectorizedRowBatch batch) 메서드를 오버라이드해야 해요. 이 메서드는 사용자가 직접 사용할 VectorizedRowBatch 인스턴스를 제공하므로, 사용자는 입력 elementColumnVectors로 변환하고 제공된 VectorizedRowBatch 인스턴스에 설정하는 로직만 작성하면 돼요.

예를 들어, 입력 요소가 다음과 같은 Person 타입이라면:

class Person {
    private final String name;
    private final int age;
    ...
}

Person 타입 요소를 변환하고 VectorizedRowBatch에 설정하는 자식 구현은 다음과 같을 수 있어요:

import org.apache.hadoop.hive.ql.exec.vector.BytesColumnVector;
import org.apache.hadoop.hive.ql.exec.vector.LongColumnVector;

import java.io.IOException;
import java.io.Serializable;
import java.nio.charset.StandardCharsets;

public class PersonVectorizer extends Vectorizer<Person> implements Serializable {	
    public PersonVectorizer(String schema) {
        super(schema);
    }
    @Override
    public void vectorize(Person element, VectorizedRowBatch batch) throws IOException {
        BytesColumnVector nameColVector = (BytesColumnVector) batch.cols[0];
        LongColumnVector ageColVector = (LongColumnVector) batch.cols[1];
        int row = batch.size++;
        nameColVector.setVal(row, element.getName().getBytes(StandardCharsets.UTF_8));
        ageColVector.vector[row] = element.getAge();
    }
}

애플리케이션에서 ORC bulk 인코더를 사용하려면 다음 의존성을 추가해야 해요:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-orc</artifactId>
    <version>2.3.0</version>
</dependency>

그리고 데이터를 ORC 포맷으로 쓰는 FileSink는 다음과 같이 만들 수 있어요:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.orc.writer.OrcBulkWriterFactory;

String schema = "struct<_col0:string,_col1:int>";
DataStream<Person> input = ...;

final OrcBulkWriterFactory<Person> writerFactory = new OrcBulkWriterFactory<>(new PersonVectorizer(schema));

final FileSink<Person> sink = FileSink
    .forBulkFormat(outputBasePath, writerFactory)
    .build();

input.sinkTo(sink);

OrcBulkWriterFactory는 커스텀 Hadoop 구성과 ORC 라이터 속성을 제공할 수 있도록 Hadoop ConfigurationProperties를 받을 수도 있어요.

String schema = ...;
Configuration conf = ...;
Properties writerProperties = new Properties();

writerProperties.setProperty("orc.compress", "LZ4");
// Other ORC supported properties can also be set similarly.

final OrcBulkWriterFactory<Person> writerFactory = new OrcBulkWriterFactory<>(
    new PersonVectorizer(schema), writerProperties, conf);

ORC 라이터 속성의 전체 목록은 여기에서 찾을 수 있어요.

ORC 파일에 사용자 메타데이터를 추가하려는 사용자는 오버라이드된 vectorize(...) 메서드 안에서 addUserMetadata(...)를 호출해 그렇게 할 수 있어요.

public class PersonVectorizer extends Vectorizer<Person> implements Serializable {	
    @Override
    public void vectorize(Person element, VectorizedRowBatch batch) throws IOException {
        ...
        String metadataKey = ...;
        ByteBuffer metadataValue = ...;
        this.addUserMetadata(metadataKey, metadataValue);
    }
}

PyFlink 사용자는 OrcBulkWriters를 사용해 Orc 포맷으로 파일에 레코드를 쓰는 BulkWriterFactory를 만들 수 있어요. 필요한 의존성은 flink-sql-orc 2.3.0 (Download)예요.

row_type = DataTypes.ROW([
    DataTypes.FIELD('name', DataTypes.STRING()),
    DataTypes.FIELD('age', DataTypes.INT()),
])

sink = FileSink.for_bulk_format(
    OUTPUT_DIR,
    OrcBulkWriters.for_row_type(
        row_type=row_type,
        writer_properties=Configuration(),
        hadoop_config=Configuration(),
    )
).build()

ds.sink_to(sink)
Hadoop SequenceFile 포맷

애플리케이션에서 SequenceFile bulk 인코더를 사용하려면 다음 의존성을 추가해야 해요:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-sequence-file</artifactId>
    <version>2.3.0</version>
</dependency>

간단한 SequenceFile 라이터는 다음과 같이 만들 수 있어요:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.configuration.GlobalConfiguration;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.SequenceFile;
import org.apache.hadoop.io.Text;

DataStream<Tuple2<LongWritable, Text>> input = ...;
Configuration hadoopConf = HadoopUtils.getHadoopConfiguration(GlobalConfiguration.loadConfiguration());
final FileSink<Tuple2<LongWritable, Text>> sink = FileSink
  .forBulkFormat(
    outputBasePath,
    new SequenceFileWriterFactory<>(hadoopConf, LongWritable.class, Text.class))
    .build();

input.sinkTo(sink);

SequenceFileWriterFactory는 압축 설정을 지정하기 위한 추가 생성자 파라미터를 지원해요.

버킷 할당 (Bucket Assignment)

버케팅 로직은 기본 출력 디렉터리 내의 하위 디렉터리로 데이터가 어떻게 구조화될지 정의해요.

행과 bulk 포맷 모두(see File Formats) 기본 할당기로 DateTimeBucketAssigner를 사용해요. 기본적으로 DateTimeBucketAssigner는 시스템 기본 시간대를 기준으로 yyyy-MM-dd--HH 형식의 시간별 버킷을 만들어요. 날짜 포맷(즉 버킷 크기)과 시간대 모두 수동으로 구성할 수 있어요.

포맷 빌더에서 .withBucketAssigner(assigner)을 호출해 커스텀 BucketAssigner를 지정할 수 있어요.

Flink는 두 가지 내장 BucketAssigner와 함께 제공돼요:

  • DateTimeBucketAssigner : 기본 시간 기반 할당기
  • BasePathBucketAssigner : 모든 파트 파일을 기본 경로에 저장하는 할당기(단일 전역 버킷)

참고: PyFlink는 DateTimeBucketAssignerBasePathBucketAssigner만 지원해요.

롤링 정책 (Rolling Policy)

RollingPolicy는 주어진 in-progress 파트 파일이 언제 닫히고 pending으로, 나중에 finished 상태로 이동할지 정의해요. "finished" 상태의 파트 파일은 보기에 준비된 것이며 실패 시 되돌려지지 않는 유효한 데이터를 포함한다는 것이 보장돼요. STREAMING 모드에서 롤링 정책은 체크포인팅 간격과 결합되어(pending 파일은 다음 체크포인트에서 finished가 돼) 파트 파일이 다운스트림 리더에게 얼마나 빨리 사용 가능해지는지와 이러한 파트들의 크기 및 개수를 제어해요. BATCH 모드에서 파트 파일은 작업이 끝날 때 표시되지만 롤링 정책이 최대 크기를 제어할 수 있어요.

Flink는 두 가지 내장 RollingPolicy와 함께 제공돼요:

  • DefaultRollingPolicy
  • OnCheckpointRollingPolicy

참고: PyFlink는 DefaultRollingPolicyOnCheckpointRollingPolicy만 지원해요.

파트 파일 수명 주기 (Part file lifecycle)

다운스트림 시스템에서 FileSink의 출력을 사용하려면, 생성되는 출력 파일의 명명과 수명 주기를 이해해야 해요.

파트 파일은 세 가지 상태 중 하나일 수 있어요:

  1. In-progress : 현재 쓰여지고 있는 파트 파일은 in-progress예요.
  2. Pending : 닫힌(지정된 롤링 정책으로 인해) in-progress 파일로, 커밋되기를 기다리는 것.
  3. Finished : 성공적인 체크포인트(STREAMING) 또는 입력의 끝(BATCH)에서 pending 파일이 "Finished"로 전환.

다운스트림 시스템이 안전하게 읽을 수 있는 것은 finished 파일뿐이에요. 그것들은 나중에 수정되지 않는다고 보장되기 때문이에요.

각 라이터 서브태스크는 활성 버킷마다 주어진 시점에 단일 in-progress 파트 파일을 가지지만, 여러 pending과 finished 파일이 있을 수 있어요.

파트 파일 예시

이 파일들의 수명 주기를 더 잘 이해하기 위해 2개의 싱크 서브태스크가 있는 간단한 예시를 살펴보죠:

└── 2019-08-25--12
    ├── part-4005733d-a830-4323-8291-8866de98b582-0.inprogress.bd053eb0-5ecf-4c85-8433-9eff486ac334
    └── part-81fc4980-a6af-41c8-9937-9939408a734b-0.inprogress.ea65a428-a1d0-4a0b-bbc5-7a436a75e575

파트 파일 part-81fc4980-a6af-41c8-9937-9939408a734b-0이 롤되면(너무 커졌다고 가정), pending이 되지만 이름이 바뀌지는 않아요. 그런 다음 싱크는 새 파트 파일 part-81fc4980-a6af-41c8-9937-9939408a734b-1을 열어요:

└── 2019-08-25--12
    ├── part-4005733d-a830-4323-8291-8866de98b582-0.inprogress.bd053eb0-5ecf-4c85-8433-9eff486ac334
    ├── part-81fc4980-a6af-41c8-9937-9939408a734b-0.inprogress.ea65a428-a1d0-4a0b-bbc5-7a436a75e575
    └── part-81fc4980-a6af-41c8-9937-9939408a734b-1.inprogress.bc279efe-b16f-47d8-b828-00ef6e2fbd11

part-81fc4980-a6af-41c8-9937-9939408a734b-0이 이제 pending 완료이므로, 다음 성공적인 체크포인트 후 최종화됩니다:

└── 2019-08-25--12
    ├── part-4005733d-a830-4323-8291-8866de98b582-0.inprogress.bd053eb0-5ecf-4c85-8433-9eff486ac334
    ├── part-81fc4980-a6af-41c8-9937-9939408a734b-0
    └── part-81fc4980-a6af-41c8-9937-9939408a734b-1.inprogress.bc279efe-b16f-47d8-b828-00ef6e2fbd11

버케팅 정책이 지시하는 대로 새 버킷이 생성되며, 이것은 현재 in-progress 파일에 영향을 주지 않아요:

└── 2019-08-25--12
    ├── part-4005733d-a830-4323-8291-8866de98b582-0.inprogress.bd053eb0-5ecf-4c85-8433-9eff486ac334
    ├── part-81fc4980-a6af-41c8-9937-9939408a734b-0
    └── part-81fc4980-a6af-41c8-9937-9939408a734b-1.inprogress.bc279efe-b16f-47d8-b828-00ef6e2fbd11
└── 2019-08-25--13
    └── part-4005733d-a830-4323-8291-8866de98b582-0.inprogress.2b475fec-1482-4dea-9946-eb4353b475f1

버케팅 정책은 레코드 단위로 평가되므로, 오래된 버킷도 여전히 새 레코드를 받을 수 있어요.

파트 파일 구성 (Part file configuration)

Finished 파일은 명명 스킴에 의해서만 in-progress 파일과 구분될 수 있어요.

기본적으로 파일 명명 전략은 다음과 같아요:

  • In-progress / Pending: part-<random-uuid>-<task-id>.inprogress.<uid>
  • Finished: part-<random-uuid>-<task-id> 여기서 uid는 서브태스크가 인스턴스화될 때 싱크의 서브태스크에 할당된 무작위 id예요. 이 uid는 결함 허용(fault-tolerant)되지 않으므로 서브태스크가 실패에서 복구될 때 다시 생성돼요.

Flink는 사용자가 자신의 파트 파일에 접두사 및/또는 접미사를 지정할 수 있게 해줘요. 이것은 OutputFileConfig를 사용해 수행될 수 있어요. 예를 들어 접두사 "prefix"와 접미사 ".ext"의 경우 싱크는 다음 파일을 만들어요:

└── 2019-08-25--12
    ├── prefix-4005733d-a830-4323-8291-8866de98b582-0.ext
    ├── prefix-4005733d-a830-4323-8291-8866de98b582-1.ext.inprogress.bd053eb0-5ecf-4c85-8433-9eff486ac334
    ├── prefix-81fc4980-a6af-41c8-9937-9939408a734b-0.ext
    └── prefix-81fc4980-a6af-41c8-9937-9939408a734b-1.ext.inprogress.bc279efe-b16f-47d8-b828-00ef6e2fbd11

사용자는 다음과 같은 방식으로 OutputFileConfig를 지정할 수 있어요:

OutputFileConfig config = OutputFileConfig
 .builder()
 .withPartPrefix("prefix")
 .withPartSuffix(".ext")
 .build();
            
FileSink<Tuple2<Integer, Integer>> sink = FileSink
 .forRowFormat((new Path(outputPath), new SimpleStringEncoder<>("UTF-8"))
 .withBucketAssigner(new KeyBucketAssigner())
 .withRollingPolicy(OnCheckpointRollingPolicy.build())
 .withOutputFileConfig(config)
 .build();
config = OutputFileConfig \
    .builder() \
    .with_part_prefix("prefix") \
    .with_part_suffix(".ext") \
    .build()

sink = FileSink \
    .for_row_format(OUTPUT_PATH, Encoder.simple_string_encoder("UTF-8")) \
    .with_bucket_assigner(BucketAssigner.base_path_bucket_assigner()) \
    .with_rolling_policy(RollingPolicy.on_checkpoint_rolling_policy()) \
    .with_output_file_config(config) \
    .build()

컴팩션 (Compaction)

버전 1.15부터 FileSinkpending 파일의 컴팩션(compaction)을 지원해요. 이는 애플리케이션이 많은 작은 파일을 생성하지 않고 더 작은 체크포인트 간격을 가질 수 있게 해주며, 특히 체크포인트 시점에 롤해야 하는 bulk 인코딩 포맷을 사용할 때 유용해요.

컴팩션은 다음과 같이 활성화할 수 있어요:

FileSink<Integer> fileSink=
    FileSink.forRowFormat(new Path(path),new SimpleStringEncoder<Integer>())
        .enableCompact(
            FileCompactStrategy.Builder.newBuilder()
                .setSizeThreshold(1024)
                .enableCompactionOnCheckpoint(5)
                .build(),
            new RecordWiseFileCompactor<>(
                new DecoderBasedReader.Factory<>(SimpleStringDecoder::new)))
        .build();
file_sink = FileSink \
    .for_row_format(PATH, Encoder.simple_string_encoder()) \
    .enable_compact(
        FileCompactStrategy.builder()
            .set_size_threshold(1024)
            .enable_compaction_on_checkpoint(5)
            .build(),
        FileCompactor.concat_file_compactor()) \
    .build()

활성화되면 컴팩션은 파일이 pending이 되고 커밋되는 사이에 일어나요. pending 파일은 먼저 경로가 .로 시작하는 임시 파일로 커밋돼요. 그런 다음 이 파일들은 사용자가 지정한 컴팩터에 의해 전략에 따라 컴팩션되고, 새 컴팩션된 pending 파일이 생성돼요. 그런 다음 이 pending 파일들이 커미터로 방출되어 공식 파일로 커밋돼요. 그 후 소스 파일은 제거돼요.

컴팩션을 활성화할 때 FileCompactStrategyFileCompactor를 지정해야 해요.

FileCompactStrategy는 언제 어떤 파일이 컴팩션되는지 지정해요. 현재 두 가지 병렬 조건이 있어요: 대상 파일 크기와 전달된 체크포인트 수. 캐시된 파일의 총 크기가 크기 임계값에 도달하거나 마지막 컴팩션 이후 체크포인트 수가 지정된 수에 도달하면, 캐시된 파일들이 컴팩션되도록 스케줄링돼요.

FileCompactor는 주어진 Path 목록을 어떻게 컴팩션하고 결과 파일을 쓸지 지정해요. 파일을 쓰는 방식에 따라 두 가지 유형으로 분류될 수 있어요:

  • OutputStreamBasedFileCompactor: 사용자는 컴팩션된 결과를 출력 스트림으로 쓸 수 있어요. 이것은 사용자가 입력 파일에서 레코드를 읽고 싶지 않거나 읽을 수 없을 때 유용해요. 예로는 파일 목록을 직접 연결하는 ConcatFileCompactor가 있어요.
  • RecordWiseFileCompactor: 컴팩터는 입력 파일에서 레코드를 하나씩 읽고 FileWriter처럼 결과 파일로 쓸 수 있어요. 예로는 소스 파일에서 레코드를 읽고 CompactingFileWriter로 쓰는 RecordWiseFileCompactor가 있어요. 사용자는 소스 파일에서 레코드를 읽는 방법을 지정해야 해요.

중요 참고 1 컴팩션이 활성화되면, 컴팩션을 비활성화하려면 FileSink를 구축할 때 명시적으로 disableCompact를 호출해야 해요.

중요 참고 2 컴팩션이 활성화되면, 작성된 파일이 표시되기 전에 더 오래 기다려야 해요.

참고: PyFlink는 ConcatFileCompactorIdenticalFileCompactor만 지원해요.

중요한 고려 사항 (Important Considerations)

일반 (General)

중요 참고 1: Hadoop < 2.7을 사용할 때는 매 체크포인트마다 파트 파일을 롤하는 OnCheckpointRollingPolicy를 사용하세요. 그 이유는 파트 파일이 체크포인트 간격을 "가로지르면", 실패에서 복구 시 FileSink가 in-progress 파일에서 커밋되지 않은 데이터를 버리기 위해 파일시스템의 truncate() 메서드를 사용할 수 있기 때문이에요. 이 메서드는 2.7 이전 Hadoop 버전에서 지원되지 않으며 Flink가 예외를 던질 거예요.

중요 참고 2: Flink 싱크와 UDF는 일반적으로 정상 작업 종료(예: 유한 입력 스트림)와 실패로 인한 종료를 구분하지 않으므로, 작업의 정상 종료 시 마지막 in-progress 파일은 "finished" 상태로 전환되지 않아요.

중요 참고 3: Flink와 FileSink는 커밋된 데이터를 결코 덮어쓰지 않아요. 따라서 후속 성공적인 체크포인트에 의해 커밋된 in-progress 파일을 가정하는 오래된 체크포인트/savepoint에서 복원하려고 할 때, FileSink는 in-progress 파일을 찾을 수 없으므로 복원을 거부하고 예외를 던질 거예요.

중요 참고 4: 현재 FileSink는 다섯 개의 파일시스템만 지원해요: HDFS, S3, OSS, ABFS, Local. Flink는 지원되지 않는 파일시스템을 런타임에서 사용하면 예외를 던질 거예요.

BATCH 전용

중요 참고 1: Writer는 사용자가 지정한 병렬도로 실행되지만, Committer는 1과 같은 병렬도로 실행돼요.

중요 참고 2: Pending 파일은 전체 입력이 처리된 후 커밋, 즉 Finished 상태로 전환돼요.

중요 참고 3: 고가용성(High-Availability)이 활성화되면, Committer가 커밋하는 동안 JobManager 실패가 발생하면 중복이 발생할 수 있어요. 이것은 미래 Flink 버전에서 수정될 예정이에요 (FLIP-147에서 진행 상황 참고).

S3 전용

중요 참고 1: S3의 경우 FileSinkHadoop 기반 FileSystem 구현만 지원하며, Presto 기반 구현은 지원하지 않아요. 작업이 FileSink로 S3에 쓰지만 체크포인팅에는 Presto 기반을 사용하고 싶다면, 싱크의 대상 경로로는 명시적으로 "s3a://"(Hadoop용)를, 체크포인팅에는 "s3p://"(Presto용)를 사용하는 것이 좋아요. 싱크와 체크포인팅 모두에 *"s3://"*를 사용하면 두 구현 모두 그 스킴을 "듣기" 때문에 예측할 수 없는 동작이 발생할 수 있어요.

중요 참고 2: 효율적이면서 exactly-once 의미론을 보장하기 위해 FileSink는 S3의 Multi-part Upload 기능(이하 MPU)을 사용해요. 이 기능은 파일을 독립적인 청크로 업로드할 수 있게 하며("multi-part"), MPU의 모든 부분이 성공적으로 업로드되면 원래 파일로 결합될 수 있어요. 비활성 MPU의 경우, S3는 시작 후 지정된 일 수 내에 완료되지 않는 multipart 업로드를 중단하는 데 사용할 수 있는 버킷 수명 주기 규칙을 지원해요. 이것은 이 규칙을 공격적으로 설정하고 일부 파트 파일이 완전히 업로드되지 않은 상태로 savepoint를 찍으면, 작업이 재시작되기 전에 관련 MPU가 타임아웃될 수 있다는 것을 의미해요. 이로 인해 pending 파트 파일이 더 이상 없으므로 작업이 그 savepoint에서 복원할 수 없게 되고, Flink는 그것을 가져오려다 실패하며 예외를 던질 거예요.

OSS 전용

중요 참고: 효율적이면서 exactly-once 의미론을 보장하기 위해 FileSink는 OSS의 Multi-part Upload 기능도 사용해요(Similar with S3).

더 알아보기 (Learn more)