이벤트 기반 애플리케이션
이벤트 기반 애플리케이션 (Event-driven Applications)
ProcessFunction은 이벤트 처리와 타이머, 상태를 결합해 스트림 처리 애플리케이션을 위한 강력한 구성 요소가 됩니다. 이것은 Flink로 이벤트 기반 애플리케이션을 만드는 기반입니다. DataStream API의 KeyedProcessFunction을 사용하든 Table API의 ProcessTableFunction을 사용하든 개념은 비슷합니다: 이벤트를 하나씩 처리하고, 상태를 유지하며, 미래 콜백을 위한 타이머를 등록할 수 있습니다.
출처: 문서
본문
타이머를 이용한 사기 탐지 (Fraud Detection with Timers)
DataStream API 튜토리얼을 완료했다면 작은 금액 후 큰 금액이 오는 거래 패턴을 식별하는 사기 탐지기를 만들었을 것입니다. 하지만 그 구현에는 한계가 있습니다: 시간을 고려하지 않는다는 점이죠. 실제 사기범은 테스트 거래와 큰 구매 사이에 오래 기다리지 않습니다 — 테스트 거래가 발각될 가능성을 최소화하려 하기 때문입니다.
사기 탐지기를 개선해 1분 이내에 발생하는 거래만 플래그하도록 만들어 봅시다.
타이머 추가 (Adding a Timer)
Flink의 KeyedProcessFunction을 사용하면 미래의 어떤 시점에 콜백 메서드를 호출하는 타이머를 설정할 수 있습니다. 요구 사항은 다음과 같습니다:
- 플래그가
true로 설정될 때마다 1분 후의 타이머도 설정합니다. - 타이머가 실행되면 상태를 지워 플래그를 재설정합니다.
- 플래그가 지워지면 타이머는 취소되어야 합니다.
타이머를 취소하려면 타이머가 설정된 시간을 기억해야 하며, 기억한다는 것은 상태를 의미하므로 먼저 플래그 상태와 함께 타이머 상태를 만들기 시작합니다.
private transient ValueState<Boolean> flagState;
private transient ValueState<Long> timerState;
@Override
public void open(OpenContext openContext) {
ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
"flag",
Types.BOOLEAN);
flagState = getRuntimeContext().getState(flagDescriptor);
ValueStateDescriptor<Long> timerDescriptor = new ValueStateDescriptor<>(
"timer-state",
Types.LONG);
timerState = getRuntimeContext().getState(timerDescriptor);
}
타이머 등록 (Registering the Timer)
KeyedProcessFunction#processElement는 타이머 서비스를 포함하는 Context와 함께 호출됩니다. 타이머 서비스는 현재 시간을 조회하고, 타이머를 등록하고, 타이머를 삭제하는 데 사용할 수 있습니다. 이를 통해 플래그가 설정될 때마다 1분 후의 타이머를 설정하고 타임스탬프를 timerState에 저장할 수 있습니다.
if (transaction.getAmount() < SMALL_AMOUNT) {
// 플래그를 true로 설정
flagState.update(true);
// 타이머와 타이머 상태 설정
long timer = context.timerService().currentProcessingTime() + ONE_MINUTE;
context.timerService().registerProcessingTimeTimer(timer);
timerState.update(timer);
}
처리 시간(processing time)은 벽시계 시간이며, 연산자를 실행하는 머신의 시스템 클록에 의해 결정됩니다.
타이머 콜백 처리 (Handling Timer Callbacks)
타이머가 실행되면 KeyedProcessFunction#onTimer를 호출합니다. 이 메서드를 재정의해 플래그를 재설정하는 콜백을 구현할 수 있습니다.
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<Alert> out) {
// 1분 후 플래그 제거
timerState.clear();
flagState.clear();
}
상태 정리와 타이머 취소 (Cleaning Up State and Canceling Timers)
마지막으로, 타이머를 취소하려면 등록된 타이머를 삭제하고 타이머 상태를 삭제해야 합니다. 이것을 헬퍼 메서드로 감싸고 flagState.clear() 대신 이 메서드를 호출할 수 있습니다.
private void cleanUp(Context ctx) throws Exception {
// 타이머 삭제
Long timer = timerState.value();
ctx.timerService().deleteProcessingTimeTimer(timer);
// 모든 상태 정리
timerState.clear();
flagState.clear();
}
타이머가 있는 완전한 사기 탐지기 (Complete Fraud Detector with Timers)
다음은 타이머 기반 사기 탐지를 포함한 완전한 구현입니다:
import org.apache.flink.api.common.functions.OpenContext;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.util.Collector;
import org.apache.flink.walkthrough.common.entity.Alert;
import org.apache.flink.walkthrough.common.entity.Transaction;
public class FraudDetector extends KeyedProcessFunction<Long, Transaction, Alert> {
private static final long serialVersionUID = 1L;
private static final double SMALL_AMOUNT = 1.00;
private static final double LARGE_AMOUNT = 500.00;
private static final long ONE_MINUTE = 60 * 1000;
private transient ValueState<Boolean> flagState;
private transient ValueState<Long> timerState;
@Override
public void open(OpenContext openContext) {
ValueStateDescriptor<Boolean> flagDescriptor = new ValueStateDescriptor<>(
"flag",
Types.BOOLEAN);
flagState = getRuntimeContext().getState(flagDescriptor);
ValueStateDescriptor<Long> timerDescriptor = new ValueStateDescriptor<>(
"timer-state",
Types.LONG);
timerState = getRuntimeContext().getState(timerDescriptor);
}
@Override
public void processElement(
Transaction transaction,
Context context,
Collector<Alert> collector) throws Exception {
// 현재 키에 대한 현재 상태 가져오기
Boolean lastTransactionWasSmall = flagState.value();
// 플래그가 설정되었는지 확인
if (lastTransactionWasSmall != null) {
if (transaction.getAmount() > LARGE_AMOUNT) {
// 다운스트림으로 경고 출력
Alert alert = new Alert();
alert.setId(transaction.getAccountId());
collector.collect(alert);
}
// 상태 정리
cleanUp(context);
}
if (transaction.getAmount() < SMALL_AMOUNT) {
// 플래그를 true로 설정
flagState.update(true);
long timer = context.timerService().currentProcessingTime() + ONE_MINUTE;
context.timerService().registerProcessingTimeTimer(timer);
timerState.update(timer);
}
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<Alert> out) {
// 1분 후 플래그 제거
timerState.clear();
flagState.clear();
}
private void cleanUp(Context ctx) throws Exception {
// 타이머 삭제
Long timer = timerState.value();
ctx.timerService().deleteProcessingTimeTimer(timer);
// 모든 상태 정리
timerState.clear();
flagState.clear();
}
}
이 구현을 사용하면 1분 이내에 작은 금액 뒤에 큰 금액이 오는 거래만 경고를 트리거합니다.
프로세스 함수 (Process Functions)
소개 (Introduction)
ProcessFunction은 이벤트 처리와 타이머, 상태를 결합해 스트림 처리 애플리케이션을 위한 강력한 구성 요소가 됩니다. 이것은 Flink로 이벤트 기반 애플리케이션을 만드는 기반입니다. RichFlatMapFunction과 매우 유사하지만 타이머가 추가됩니다.
예제 (Example)
Streaming Analytics 훈련의 hands-on 연습을 마쳤다면, TumblingEventTimeWindow를 사용해 시간당 각 드라이버의 팁 합계를 계산했던 것을 기억할 것입니다:
// 각 드라이버의 시간당 팁 합계 계산
DataStream<Tuple3<Long, Long, Float>> hourlyTips = fares
.keyBy((TaxiFare fare) -> fare.driverId)
.window(TumblingEventTimeWindows.of(Duration.ofHours(1)))
.process(new AddTips());
동일한 작업을 KeyedProcessFunction으로 하는 것도 상당히 간단하고 교육적입니다. 위 코드를 다음과 같이 바꾸는 것부터 시작해 봅시다:
// 각 드라이버의 시간당 팁 합계 계산
DataStream<Tuple3<Long, Long, Float>> hourlyTips = fares
.keyBy((TaxiFare fare) -> fare.driverId)
.process(new PseudoWindow(Duration.ofHours(1)));
이 코드 조각에서는 PseudoWindow라는 KeyedProcessFunction이 키드 스트림에 적용되고, 그 결과는 DataStream<Tuple3<Long, Long, Float>>입니다(Flink 내장 시간 윈도우를 사용하는 구현이 만드는 것과 같은 종류의 스트림).
PseudoWindow의 전체 구조는 다음과 같은 모양입니다:
// 시간 길이의 윈도우에서 각 드라이버의 팁 합계를 계산합니다.
// 키는 driverId입니다.
public static class PseudoWindow extends
KeyedProcessFunction<Long, TaxiFare, Tuple3<Long, Long, Float>> {
private final long durationMsec;
public PseudoWindow(Time duration) {
this.durationMsec = duration.toMilliseconds();
}
@Override
// 초기화 중 한 번 호출됩니다.
public void open(OpenContext ctx) {
. . .
}
@Override
// 각 요금이 처리되도록 도착할 때 호출됩니다.
public void processElement(
TaxiFare fare,
Context ctx,
Collector<Tuple3<Long, Long, Float>> out) throws Exception {
. . .
}
@Override
// 현재 워터마크가 윈도우가 이제 완료됨을 나타낼 때 호출됩니다.
public void onTimer(long timestamp,
OnTimerContext context,
Collector<Tuple3<Long, Long, Float>> out) throws Exception {
. . .
}
}
알아야 할 사항:
-
여러 유형의 ProcessFunction이 있습니다 — 이것은
KeyedProcessFunction이지만CoProcessFunctions,BroadcastProcessFunctions등도 있습니다. -
KeyedProcessFunction은RichFunction의 일종입니다.RichFunction이므로 관리되는 키드 상태를 다루는 데 필요한open과getRuntimeContext메서드에 접근할 수 있습니다. -
구현할 콜백은 두 가지입니다:
processElement와onTimer.processElement는 각 수신 이벤트와 함께 호출되고,onTimer는 타이머가 실행될 때 호출됩니다. 이들은 이벤트 시간 또는 처리 시간 타이머일 수 있습니다.processElement와onTimer모두TimerService와 상호작용하는 데 사용할 수 있는 컨텍스트 객체가 제공됩니다. 두 콜백 모두 결과를 방출하는 데 사용할 수 있는Collector도 전달됩니다.
open() 메서드
// 키드, 관리 상태, 각 윈도우에 대한 항목, 윈도우 종료 시간으로 키 지정
// 각 드라이버에 대해 별도의 MapState 객체가 있습니다.
private transient MapState<Long, Float> sumOfTips;
@Override
public void open(OpenContext ctx) {
MapStateDescriptor<Long, Float> sumDesc =
new MapStateDescriptor<>("sumOfTips", Long.class, Float.class);
sumOfTips = getRuntimeContext().getMapState(sumDesc);
}
요금 이벤트는 순서가 어긋나 도착할 수 있으므로, 이전 시간의 결과 계산을 끝내기 전에 한 시간의 이벤트를 처리해야 하는 경우가 있습니다. 실제로 워터마킹 지연이 윈도우 길이보다 훨씬 길다면, 단지 두 개가 아니라 여러 윈도우가 동시에 열려 있을 수 있습니다. 이 구현은 각 윈도우의 끝 타임스탬프를 그 윈도우의 팁 합계에 매핑하는 MapState를 사용해 이를 지원합니다.
processElement() 메서드
public void processElement(
TaxiFare fare,
Context ctx,
Collector<Tuple3<Long, Long, Float>> out) throws Exception {
long eventTime = fare.getEventTime();
TimerService timerService = ctx.timerService();
if (eventTime <= timerService.currentWatermark()) {
// 이 이벤트는 늦었습니다. 해당 윈도우는 이미 트리거되었습니다.
} else {
// eventTime을 이 이벤트를 포함하는 윈도우의 끝으로 올림
long endOfWindow = (eventTime - (eventTime % durationMsec) + durationMsec - 1);
// 윈도우가 완료되었을 때 콜백을 예약
timerService.registerEventTimeTimer(endOfWindow);
// 이 요금의 팁을 해당 윈도우의 누계에 추가
Float sum = sumOfTips.get(endOfWindow);
if (sum == null) {
sum = 0.0F;
}
sum += fare.tip;
sumOfTips.put(endOfWindow, sum);
}
}
고려할 사항:
-
늦은 이벤트는 어떻게 되나요? 워터마크보다 뒤처진(즉, 늦은) 이벤트는 버려집니다. 이보다 더 나은 방법을 원한다면 다음 섹션에서 설명하는 사이드 아웃풋을 고려하세요.
-
이 예제는 키가 타임스탬프인
MapState를 사용하고 같은 타임스탬프에Timer를 설정합니다. 이것은 일반적인 패턴입니다. 타이머가 실행될 때 관련 정보를 쉽고 효율적으로 조회하게 해줍니다.
onTimer() 메서드
public void onTimer(
long timestamp,
OnTimerContext context,
Collector<Tuple3<Long, Long, Float>> out) throws Exception {
long driverId = context.getCurrentKey();
// 방금 끝난 시간의 결과 조회
Float sumOfTips = this.sumOfTips.get(timestamp);
Tuple3<Long, Long, Float> result = Tuple3.of(driverId, timestamp, sumOfTips);
out.collect(result);
this.sumOfTips.remove(timestamp);
}
관찰 사항:
-
onTimer에 전달되는OnTimerContext context는 현재 키를 결정하는 데 사용할 수 있습니다. -
우리의 의사 윈도우는 현재 워터마크가 매시간 끝에 도달할 때 트리거되며, 그 시점에
onTimer가 호출됩니다. 이 onTimer 메서드는sumOfTips에서 관련 항목을 제거하는데, 이는 늦은 이벤트를 수용할 수 없게 만드는 효과가 있습니다. 이것은 Flink 시간 윈도우에서 allowedLateness를 0으로 설정하는 것과 동등합니다.
성능 고려 사항 (Performance Considerations)
Flink는 RocksDB에 최적화된 MapState와 ListState 타입을 제공합니다. 가능하면 어떤 종류의 컬렉션을 보유하는 ValueState 객체 대신 이것들을 사용해야 합니다. RocksDB 상태 백엔드는 (역)직렬화를 거치지 않고 ListState에 append할 수 있으며, MapState의 경우 각 키/값 쌍이 별도의 RocksDB 객체이므로 MapState를 효율적으로 접근하고 업데이트할 수 있습니다.
사이드 아웃풋 (Side Outputs)
소개 (Introduction)
Flink 연산자에서 출력 스트림이 둘 이상 있기를 원할 이유는 여러 가지가 있습니다. 예를 들어:
- 예외(exceptions)
- 잘못된 이벤트(malformed events)
- 늦은 이벤트(late events)
- 외부 서비스로의 타임아웃 연결 같은 운영 경고(operational alerts)
사이드 아웃풋은 이를 위한 편리한 방법입니다. 오류 보고 외에도 사이드 아웃풋은 스트림의 n-way 분할을 구현하는 좋은 방법입니다.
예제 (Example)
이제 이전 섹션에서 무시했던 늦은 이벤트로 무언가를 할 수 있는 위치에 있습니다.
사이드 아웃풋 채널은 OutputTag와 연관됩니다. 이 태그들은 사이드 아웃풋 DataStream의 타입에 해당하는 제네릭 타입을 가지며, 이름을 갖습니다.
private static final OutputTag<TaxiFare> lateFares = new OutputTag<TaxiFare>("lateFares") {};
위는 정적 OutputTag로, PseudoWindow의 processElement 메서드에서 늦은 이벤트를 방출할 때:
if (eventTime <= timerService.currentWatermark()) {
// 이 이벤트는 늦었습니다. 해당 윈도우는 이미 트리거되었습니다.
ctx.output(lateFares, fare);
} else {
. . .
}
그리고 작업의 main 메서드에서 이 사이드 아웃풋의 스트림에 접근할 때 참조할 수 있습니다:
// 각 드라이버의 시간당 팁 합계 계산
SingleOutputStreamOperator hourlyTips = fares
.keyBy((TaxiFare fare) -> fare.driverId)
.process(new PseudoWindow(Duration.ofHours(1)));
hourlyTips.getSideOutput(lateFares).print();
또는 같은 이름의 두 OutputTag를 사용해 같은 사이드 아웃풋을 가리킬 수 있습니다. 하지만 그렇게 하려면 같은 타입이어야 합니다.
마무리 발언 (Closing Remarks)
이 예제에서 ProcessFunction을 사용해 간단한 시간 윈도우를 재구현하는 방법을 보았습니다. 물론 Flink 내장 windowing API가 요구 사항을 충족한다면 망설이지 말고 사용하세요. 하지만 Flink 윈도우로 이상한 일을 해야겠다고 생각한다면 직접 만드는 것을 두려워하지 마세요.
또한 ProcessFunction은 분석 계산 외에도 많은 다른 사용 사례에 유용합니다. 아래 hands-on 연습은 완전히 다른 무언가의 예를 제공합니다.
ProcessFunction의 또 다른 일반적인 사용 사례는 오래된(stale) 상태를 만료시키는 것입니다. Rides and Fares Exercise를 떠올려 보면, RichCoFlatMapFunction을 사용해 간단한 join을 계산했고, 샘플 솔루션은 각 rideId에 대해 TaxiRides와 TaxiFares가 일대일로 완벽하게 매칭된다고 가정합니다. 이벤트가 손실되면 같은 rideId의 다른 이벤트가 상태에 영원히 보관됩니다. 이것은 KeyedCoProcessFunction으로 구현하고, 타이머를 사용해 오래된 상태를 감지하고 지울 수 있습니다.
hands-on
이 섹션과 함께하는 hands-on 연습은 Long Ride Alerts Exercise입니다.