사용자 정의 함수
사용자 정의 함수 (User-Defined Functions)
대부분의 연산은 사용자 정의 함수를 필요로 해요. 이 섹션은 함수를 지정할 수 있는 다양한 방법을 나열해요. 또한 Flink 애플리케이션에 대한 통찰을 얻는 데 사용할 수 있는 Accumulator(Acumulator)도 다뤄요.
출처: 문서
본문
인터페이스 구현하기 (Implementing an interface)
가장 기본적인 방법은 제공된 인터페이스 중 하나를 구현하는 것이에요:
class MyMapFunction implements MapFunction<String, Integer> {
public Integer map(String value) { return Integer.parseInt(value); }
}
data.map(new MyMapFunction());
익명 클래스 (Anonymous classes)
함수를 익명 클래스로 전달할 수 있어요:
data.map(new MapFunction<String, Integer> () {
public Integer map(String value) { return Integer.parseInt(value); }
});
Java 8 람다 (Java 8 Lambdas)
Flink는 Java API에서 Java 8 람다도 지원해요.
data.filter(s -> s.startsWith("http://"));
data.reduce((i1,i2) -> i1 + i2);
리치 함수 (Rich functions)
사용자 정의 함수가 필요한 모든 변환은 대신 rich 함수를 인자로 받을 수 있어요. 예를 들어 다음 대신
class MyMapFunction implements MapFunction<String, Integer> {
public Integer map(String value) { return Integer.parseInt(value); }
}
이렇게 작성할 수 있어요:
class MyMapFunction extends RichMapFunction<String, Integer> {
public Integer map(String value) { return Integer.parseInt(value); }
}
그리고 평소처럼 map 변환에 함수를 전달해요:
data.map(new MyMapFunction());
Rich 함수는 익명 클래스로도 정의할 수 있어요:
data.map (new RichMapFunction<String, Integer>() {
public Integer map(String value) { return Integer.parseInt(value); }
});
Accumulator와 카운터 (Accumulators & Counters)
Accumulator는 **추가 연산(add operation)**과 작업 종료 후 사용할 수 있는 **최종 누적 결과(final accumulated result)**를 가진 간단한 구조체예요. 가장 간단한 accumulator는 counter로, Accumulator.add(V value) 메서드로 증가시킬 수 있어요. 작업이 끝나면 Flink는 모든 부분 결과를 합산(병합)해 결과를 클라이언트로 보내요. Accumulator는 디버깅 중이거나 데이터에 대해 빨리 더 알고 싶을 때 유용해요.
Flink는 현재 다음 내장 accumulator를 가지고 있어요. 각각은 Accumulator 인터페이스를 구현해요.
- IntCounter, LongCounter, 그리고 DoubleCounter: 카운터를 사용하는 예제는 아래를 참고해요.
- Histogram: 이산적인 수의 빈(bin)에 대한 히스토그램 구현이에요. 내부적으로는 Integer에서 Integer로의 맵일 뿐이에요. 이를 사용해 값의 분포를 계산할 수 있어요(예: 워드 카운트 프로그램의 행당 단어 수 분포).
accumulator 사용 방법:
먼저 사용하려는 사용자 정의 변환 함수에 accumulator 객체(여기서는 카운터)를 만들어야 해요.
private IntCounter numLines = new IntCounter();
둘째로, 일반적으로 rich 함수의 open() 메서드에서 accumulator 객체를 등록해야 해요. 여기서 이름도 정의해요.
getRuntimeContext().addAccumulator("num-lines", this.numLines);
이제 open()과 close() 메서드를 포함해 연산자 함수 어디에서나 accumulator를 사용할 수 있어요.
this.numLines.add(1);
전체 결과는 실행 환경의 execute() 메서드가 반환하는 JobExecutionResult 객체에 저장돼요(현재는 실행이 작업 완료를 기다릴 때만 작동해요).
myJobExecutionResult.getAccumulatorResult("num-lines");
모든 accumulator는 작업당 하나의 단일 네임스페이스를 공유해요. 따라서 작업의 다른 연산자 함수에서 같은 accumulator를 사용할 수 있어요. Flink는 같은 이름의 모든 accumulator를 내부적으로 병합해요.
사용자 정의 accumulator (Custom accumulators):
자신만의 accumulator를 구현하려면 Accumulator 인터페이스의 구현을 작성하면 돼요. 사용자 정의 accumulator가 Flink와 함께 배포되어야 한다고 생각하면 pull request를 만들어도 좋아요. Accumulator 또는 SimpleAccumulator 중 하나를 구현할 수 있어요.
Accumulator<V,R>이 가장 유연해요: 추가할 값의 타입 V와 최종 결과의 결과 타입 R을 정의해요. 예를 들어 히스토그램의 경우 V는 숫자이고 R은 히스토그램이에요. SimpleAccumulator는 두 타입이 같은 경우(카운터 등)를 위한 것이에요.