텍스트 파일 포맷
텍스트 파일 포맷 (Text files format)
Flink는 TextLineInputFormat을 사용해 파일에서 텍스트 줄을 읽는 것을 지원해요. 이 포맷은 Java의 내장 InputStreamReader를 사용해 다양한 지원 문자셋 인코딩으로 바이트 스트림을 디코딩해요.
출처: 문서
본문
포맷을 사용하려면 프로젝트에 Flink Connector Files 의존성을 추가해야 해요.
org.apache.flink
flink-connector-files
2.3.0
PyFlink 사용자는 job에서 이 포맷을 직접 사용할 수 있어요.
이 포맷은 배치와 스트리밍 모드 모두에서 사용할 수 있는 새 Source와 호환돼요. 따라서 이 포맷을 두 가지 방식으로 사용할 수 있어요.
- 배치 모드의 유한 읽기 (Bounded read)
- 스트리밍 모드의 연속 읽기 (Continuous read): 디렉터리에 새 파일이 나타나는지 모니터링해요
유한 읽기 예시 (Bounded read example):
이 예시에서는 텍스트 파일의 줄을 String으로 포함하는 DataStream을 만들어요. 레코드가 이벤트 타임스탬프를 포함하지 않으므로 워터마크 전략이 필요 없어요.
final FileSource<String> source =
FileSource.forRecordStreamFormat(new TextLineInputFormat(), /* Flink Path */)
.build();
final DataStream<String> stream =
env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");
source = FileSource.for_record_stream_format(StreamFormat.text_line_format(), *path).build()
stream = env.from_source(source, WatermarkStrategy.no_watermarks(), "file-source")
연속 읽기 예시 (Continuous read example):
이 예시에서는 디렉터리에 새 파일이 추가됨에 따라 무한히 커지는 텍스트 파일의 줄을 String으로 포함하는 DataStream을 만들어요. 매초 새 파일을 모니터링해요. 레코드가 이벤트 타임스탬프를 포함하지 않으므로 워터마크 전략이 필요 없어요.
final FileSource<String> source =
FileSource.forRecordStreamFormat(new TextLineInputFormat(), /* Flink Path */)
.monitorContinuously(Duration.ofSeconds(1L))
.build();
final DataStream<String> stream =
env.fromSource(source, WatermarkStrategy.noWatermarks(), "file-source");
source = FileSource \
.for_record_stream_format(StreamFormat.text_line_format(), *path) \
.monitor_continously(Duration.of_seconds(1)) \
.build()
stream = env.from_source(source, WatermarkStrategy.no_watermarks(), "file-source")