사이드 아웃풋
사이드 아웃풋 (Side Outputs)
DataStream 연산 결과로 나오는 메인 스트림 외에도 여러 개의 추가 사이드 아웃풋 결과 스트림을 생성할 수 있습니다. 결과 스트림의 데이터 타입은 메인 스트림의 데이터 타입과 일치할 필요가 없으며, 서로 다른 사이드 아웃풋의 타입도 서로 달라도 됩니다. 이 연산은 보통 스트림을 복제한 뒤 각 스트림에서 원하지 않는 데이터를 걸러내야 하는 상황에서 데이터 스트림을 분할할 때 유용합니다.
출처: 문서
본문
사이드 아웃풋을 사용할 때 먼저 사이드 아웃풋 스트림을 식별하는 OutputTag를 정의해야 합니다:
// 타입을 분석할 수 있도록 익명 내부 클래스여야 합니다
OutputTag<String> outputTag = new OutputTag<String>("side-output") {};
output_tag = OutputTag("side-output", Types.STRING())
OutputTag가 사이드 아웃풋 스트림에 포함된 요소의 타입에 따라 타입이 지정되는 방식에 주목하세요.
다음 함수들에서 사이드 아웃풋으로 데이터를 방출할 수 있습니다:
- ProcessFunction
- KeyedProcessFunction
- CoProcessFunction
- KeyedCoProcessFunction
- ProcessWindowFunction
- ProcessAllWindowFunction
위 함수들에서 사용자에게 노출되는 Context 매개변수를 사용해 OutputTag로 식별되는 사이드 아웃풋에 데이터를 방출할 수 있습니다. ProcessFunction에서 사이드 아웃풋 데이터를 방출하는 예제는 다음과 같습니다:
DataStream<Integer> input = ...;
final OutputTag<String> outputTag = new OutputTag<String>("side-output"){};
SingleOutputStreamOperator<Integer> mainDataStream = input
.process(new ProcessFunction<Integer, Integer>() {
@Override
public void processElement(
Integer value,
Context ctx,
Collector<Integer> out) throws Exception {
// 일반 출력으로 데이터 방출
out.collect(value);
// 사이드 아웃풋으로 데이터 방출
ctx.output(outputTag, "sideout-" + String.valueOf(value));
}
});
input = ... # type: DataStream
output_tag = OutputTag("side-output", Types.STRING())
class MyProcessFunction(ProcessFunction):
def process_element(self, value: int, ctx: ProcessFunction.Context):
# 일반 출력으로 데이터 방출
yield value
# 사이드 아웃풋으로 데이터 방출
yield output_tag, "sideout-" + str(value)
main_data_stream = input \
.process(MyProcessFunction(), Types.INT())
사이드 아웃풋 스트림을 검색하려면 DataStream 연산의 결과에 getSideOutput(OutputTag)를 사용합니다. 이는 사이드 아웃풋 스트림 결과로 타입이 지정된 DataStream을 반환합니다:
final OutputTag<String> outputTag = new OutputTag<String>("side-output"){};
SingleOutputStreamOperator<Integer> mainDataStream = ...;
DataStream<String> sideOutputStream = mainDataStream.getSideOutput(outputTag);
output_tag = OutputTag("side-output", Types.STRING())
main_data_stream = ... # type: DataStream
side_output_stream = main_data_stream.get_side_output(output_tag) # type: DataStream
참고: Python API에서 사이드 아웃풋을 생성한다면
get_side_output(OutputTag)를 반드시 호출해야 합니다. 그렇지 않으면 사이드 아웃풋 스트림의 결과가 메인 스트림으로 출력되어 예상치 못한 결과가 되고, 데이터 타입이 다르면 작업이 실패할 수 있습니다.