이벤트 타이머 서비스
이벤트 타이머 서비스
참고: DataStream API V2는 기존 DataStream API를 점진적으로 대체할 새로운 API 집합입니다. 현재 실험 단계이며 프로덕션에는 완전히 사용할 수 없습니다.
이벤트 타이머 서비스는 Flink가 제공하는 DataStream API의 하이레벨 확장입니다. 사용자가 특정 이벤트 시간 시점에 계산을 실행하기 위한 타이머를 등록할 수 있게 하고, Flink 프레임워크 내에서 윈도우를 언제 트리거할지 결정하는 데 도움을 줍니다.
이벤트 시간에 대한 포괄적인 설명은 Notions of Time: Event Time and Processing Time 섹션을 참조하세요.
이 섹션에서는 Flink DataStream API 내에서 이벤트 타이머 서비스를 활용하는 방법을 소개합니다.
우리는 스트림에서 이벤트 시간의 진행을 나타내는 특별한 유형의 Watermark를 사용하며, 이를 이벤트 시간 워터마크(event time watermark)라고 부릅니다. 또한 idle 입력 또는 소스를 처리하기 위해 입력이나 소스가 idle임을 나타내는 또 다른 유형의 Watermark를 구현합니다. 이를 idle 상태 워터마크(idle status watermark)라고 합니다. 자세한 내용은 Dealing With Idle Inputs / Sources를 참조하세요. 이 문서에서는 이벤트 시간 워터마크와 idle 상태 워터마크를 통칭하여 이벤트 시간 관련 워터마크(event-time related watermarks)라고 합니다.
이벤트 타이머 서비스의 핵심은 스트림을 통해 이벤트 시간 관련 워터마크를 생성하고 전파하는 데 있습니다. 이를 위해 두 가지 측면을 다뤄야 합니다.
- 이벤트 시간 관련 워터마크를 생성하는 방법.
- 이벤트 시간 관련 워터마크를 처리하는 방법.
아래에서 이 두 가지 측면을 소개합니다.
출처: 문서
본문
이벤트 시간 관련 워터마크 생성
이벤트 시간으로 작업하려면 Flink가 이벤트의 타임스탬프를 알아야 합니다. 즉 스트림의 각 요소에 이벤트 타임스탬프가 할당되어야 합니다. 이는 보통 요소의 어떤 필드에서 타임스탬프를 접근/추출하여 수행됩니다.
타임스탬프가 추출되면 Flink는 이벤트 시간 워터마크를 생성합니다. 이를 수행하는 방법은 두 가지입니다. 하나는 Flink가 제공하는 EventTimeWatermarkGeneratorBuilder를 사용하는 것이고, 다른 하나는 사용자 지정 ProcessFunction을 구현하는 것입니다.
사용자는 이 두 접근 방식 중 하나를 요구에 맞게 선택해야 합니다.
WatermarkGeneratorBuilder로 워터마크 생성
사용자는 EventTimeWatermarkGeneratorBuilder를 활용해 이벤트 시간 관련 워터마크를 생성할 수 있습니다. EventTimeWatermarkGeneratorBuilder에서 사용자가 구성할 수 있는 네 가지 측면이 있습니다.
-
[필수]
EventTimeExtractor이 함수는 각 레코드에서 이벤트 시간을 추출하는 방법을 Flink에 지시합니다.
-
[선택] 입력 Idle 타임아웃(Input Idle Timeout)
입력 스트림이 지정된 시간 동안 idle로 유지되면, Flink는 이벤트 시간 관련 워터마크를 결합할 때 이 입력을 무시하여 이벤트 시간 진행이 지연되는 것을 방지합니다. 이 타임아웃의 기본값은 0입니다.
자세한 내용은 Dealing With Idle Inputs / Sources를 참조하세요.
-
[선택] 무질서 시간(Out-of-Order Time)
입력 레코드의 무질서를 수용하기 위해 사용자는 이벤트 시간 워터마크에 대한 최대 무질서 시간을 설정할 수 있습니다. 이 설정의 기본값도 0입니다.
자세한 내용은 Fixed Amount of Lateness를 참조하세요.
-
[선택] 생성 빈도(Generation Frequency)
Flink는 이벤트 시간 관련 워터마크를 생성하는 세 가지 시나리오를 제공합니다.
- EventTimeWatermark를 생성하고 방출하지 않음.
- EventTimeWatermark를 주기적으로 생성하고 방출함.
- EventTimeWatermark를 각 이벤트에 대해 생성하고 방출함.
기본적으로 Flink는 두 번째 접근 방식을 채택하며, 이는 이벤트 시간 관련 워터마크를 주기적으로 생성하고 방출합니다. 이 주기적 생성의 간격은 구성 "pipeline.auto-watermark-interval"에 의해 결정됩니다.
EventTimeWatermarkGeneratorBuilder 인스턴스를 구성하고 얻은 후 사용자는 이를 활용해 ProcessFunction을 구축할 수 있습니다. 이 ProcessFunction은 각 레코드에서 이벤트 시간을 추출하고 해당 이벤트 시간 관련 워터마크를 생성합니다.
EventTimeWatermarkGeneratorBuilder를 사용하는 대표적인 예시는 아래와 같습니다.
NonKeyedPartitionStream stream = ...;
EventTimeWatermarkGeneratorBuilder<POJO> builder = EventTimeExtension
.newWatermarkGeneratorBuilder(pojo -> pojo.getTimeStamp()) // set event time extractor
.withIdleness(Duration.ofSeconds(10)) // set input idle timeout
.withMaxOutOfOrderTime(Duration.ofSeconds(30)) // set max out-of-order time
.periodicWatermark(Duration.ofMillis(200)); // set periodic watermark generation interval
stream.process(builder.buildAsProcessFunction())
.process(...);
주의: 타임스탬프와 이벤트 시간 워터마크 모두 1970-01-01T00:00:00Z의 Java 에포크 이후의 밀리초로 지정됩니다.
고정 지연 시간 (Fixed Amount of Lateness)
스트림에서 만날 수 있는 최대 지연을 미리 아는 시나리오가 있습니다. 예: 테스트를 위해 타임스탬프가 고정된 기간 내에 분포된 요소를 포함하는 사용자 지정 소스를 만드는 경우.
사용자는 EventTimeWatermarkGeneratorBuilder#withMaxOutOfOrderTime으로 이러한 경우를 처리할 수 있습니다. 즉 주어진 윈도우에 대한 최종 결과를 계산할 때 요소가 무시되기 전에 늦을 수 있는 최대 시간입니다. 지연은 t_w - t의 결과에 해당하며, 여기서 t는 요소의 (이벤트 시간) 타임스탬프이고 t_w는 이전 이벤트 시간 워터마크의 그것입니다. lateness > 0이면 요소는 늦은 것으로 간주되며 기본적으로 해당 윈도우에 대한 작업 결과를 계산할 때 무시됩니다.
Idle 입력 / 소스 처리
입력 split/partition/shard 중 하나가 한동안 이벤트를 운반하지 않으면 EventTimeWatermarkGeneratorBuilder도 워터마크의 기반이 될 새 정보를 얻지 못합니다. 우리는 이를 idle 입력 또는 idle 소스라고 부릅니다. 이는 일부 파티션이 여전히 이벤트를 운반할 수 있기 때문에 문제가 됩니다. 그 경우 이벤트 시간 워터마크는 모든 서로 다른 병렬 이벤트 시간 워터마크의 최소값으로 계산되므로 보류됩니다.
이를 처리하려면 EventTimeWatermarkGeneratorBuilder#withIdleness로 idle 타임아웃을 구성하여 idle을 감지하고 입력을 idle로 표시할 수 있습니다.
입력이 idle로 유지되면 EventTimeWatermarkGeneratorBuilder는 입력이 비활성임을 나타내는 idle 상태 워터마크를 방출합니다. 결과적으로 다운스트림 ProcessFunction 인스턴스는 이벤트 시간 워터마크를 결합할 때 이 입력을 무시합니다.
사용자 지정 ProcessFunction으로 워터마크 생성
사용자는 ProcessFunction을 사용자 지정하여 이벤트 시간 관련 워터마크를 만들고 보낼 수 있습니다. 이를 위해 다음 단계를 따라야 합니다.
- 사용자 지정
ProcessFunction에서EventTimeWatermarkDeclaration를 선언하고, 입력 idle 지원이 필요하면IdleStatusWatermarkDeclaration도 선언합니다. EventTimeWatermarkDeclaration를 사용해 이벤트 시간 관련 워터마크를 만듭니다.- (선택)
IdleStatusWatermarkDeclaration를 사용해 idle 상태 워터마크를 만듭니다. WatermarkManager를 통해 워터마크를 보냅니다.
다음은 이 접근 방식의 예시입니다.
public static class CustomProcessFunction
implements OneInputStreamProcessFunction<Integer, Integer> {
@Override
public Collection<? extends WatermarkDeclaration> watermarkDeclarations() {
return Set.of(
EventTimeExtension.EVENT_TIME_WATERMARK_DECLARATION,
EventTimeExtension.IDLE_STATUS_WATERMARK_DECLARATION
);
}
@Override
public void processRecord(Integer record, Collector<Integer> output, PartitionedContext ctx)
throws Exception {
// do something as needed
long eventTime = ... // get event time from record
// generate event time watermark and send to downstrea
LongWatermark eventTimeWatermark = EventTimeExtension.EVENT_TIME_WATERMARK_DECLARATION.newWatermark(eventTime);
ctx.getNonPartitionedContext()
.getWatermarkManager()
.emitWatermark(eventTimeWatermark);
}
}
주의: 타임스탬프와 이벤트 시간 워터마크 모두 1970-01-01T00:00:00Z의 Java 에포크 이후의 밀리초로 지정됩니다.
이벤트 시간 관련 워터마크 처리
이벤트 시간 관련 워터마크 생성에서 설명했듯이 Flink는 프로그래머가 자신의 타임스탬프를 할당하고 이벤트 시간 관련 워터마크를 방출할 수 있게 하는 추상화를 제공합니다.
일단 이벤트 시간 관련 워터마크가 생성되고 전파되면 Flink는 이를 사용해 프로그램의 이벤트 시간을 결정할 수 있습니다. 예: 이벤트 시간 워터마크가 Window의 종료 시간을 초과할 때 윈도우를 트리거.
또한 사용자는 이벤트 시간 관련 워터마크를 사용해 자신의 비즈니스 로직을 구현할 수 있습니다. 현재 Flink는 이벤트 시간 관련 워터마크를 사용하는 두 가지 방법을 제공합니다. 하나는 내장된 EventTimeProcessFunction을 사용하는 것이고, 다른 하나는 가장 기본적인 ProcessFunction을 사용해 이벤트 시간 관련 워터마크를 처리하는 것입니다.
첫 번째 방법은 사용자가 이벤트 타이머를 빠르게 등록/등록 해제할 수 있게 하지만, 두 번째 방법은 사용자가 이 기능을 직접 구현해야 한다는 점에 유의해야 합니다. 따라서 프로그램에서 이벤트 시간을 사용하거나 인지해야 할 때는 간단하고 효율적이므로 첫 번째 방법을 권장합니다.
각 접근 방식은 아래에 자세히 설명되어 있으며, 사용자는 요구에 맞게 하나를 선택해야 합니다.
EventTimeProcessFunction으로 이벤트 시간 관련 워터마크 처리
EventTimeProcessFunction은 Flink가 제공하며 이벤트 시간을 사용해야 하는 ProcessFunction의 래퍼입니다.
EventTimeProcessFunction을 사용할 때 사용자의 사용자 지정 ProcessFunction은 OneInputStreamProcessFunction / TwoInputBroadcastStreamProcessFunction / TwoInputNonBroadcastStreamProcessFunction / TwoOutputStreamProcessFunction 인터페이스 대신 OneInputEventTimeStreamProcessFunction / TwoInputBroadcastEventTimeStreamProcessFunction / TwoInputNonBroadcastEventTimeStreamProcessFunction / TwoOutputEventTimeStreamProcessFunction 인터페이스를 구현해야 합니다.
다음은 EventTimeProcessFunction과 OneInputEventTimeStreamProcessFunction의 인터페이스입니다.
@Experimental
public interface EventTimeProcessFunction extends ProcessFunction {
/**
* Initialize the {@link EventTimeProcessFunction} with an instance of {@link EventTimeManager}.
* Note that this method should be invoked before the open method.
*/
void initEventTimeProcessFunction(EventTimeManager eventTimeManager);
}
@Experimental
public interface OneInputEventTimeStreamProcessFunction<IN, OUT>
extends EventTimeProcessFunction, OneInputStreamProcessFunction<IN, OUT> {
/**
* The {@code #onEventTimeWatermark} method signifies that the EventTimeProcessFunction has
* received an EventTimeWatermark. Other types of watermarks will be processed by the {@code
* ProcessFunction#onWatermark} method.
*/
default void onEventTimeWatermark(
long watermarkTimestamp, Collector<OUT> output, NonPartitionedContext<OUT> ctx)
throws Exception {}
/**
* Invoked when an event-time timer fires. Note that it is only used in {@link
* KeyedPartitionStream}.
*/
default void onEventTimer(long timestamp, Collector<OUT> output, PartitionedContext<OUT> ctx) {}
}
주목해야 할 세 가지 메서드가 있습니다.
-
initEventTimeProcessFunction
- 이 메서드는
EventTimeProcessFunction이EventTimeManager의 인스턴스를 얻을 수 있게 합니다. 사용자는 이 인스턴스를 사용해 현재 이벤트 시간에 접근하고 필요에 따라 이벤트 타이머를 만들거나 삭제할 수 있습니다. 이벤트 타이머는 Keyed Partition Stream에서만 사용된다는 점에 유의하세요.
- 이 메서드는
-
onEventTimeWatermark
- 이 메서드는
EventTimeProcessFunction이 이벤트 시간 워터마크를 받았음을 나타냅니다.EventTimeProcessFunction에서 이벤트 시간 워터마크는EventTimeProcessFunction#onEventTimeWatermark로 처리되고, 다른 유형의 워터마크는EventTimeProcessFunction#onWatermark로 처리된다는 점에 유의하는 것이 중요합니다.
- 이 메서드는
-
onEventTimer
- 이 콜백 메서드는 이벤트 타이머에 의해 트리거됩니다. 이 메서드 안에서 사용자는 이벤트 타이머와 연관된 키와 이벤트 시간에 접근하고, 필요한 계산을 수행하며, 결과를 출력할 수 있습니다.
EventTimeProcessFunction을 구현하는 예시입니다.
class CustomEventTimeProcessFunction
implements OneInputEventTimeStreamProcessFunction<InputPojo, OutputPojo> {
private EventTimeManager eventTimeManager;
@Override
public void initEventTimeProcessFunction(EventTimeManager eventTimeManager) {
// get event time manager instance
this.eventTimeManager = eventTimeManager;
}
@Override
public void processRecord(
InputPojo record,
Collector<OutputPojo> output,
PartitionedContext<OutputPojo> ctx)
throws Exception {
...
// register event timer
eventTimeManager.registerTimer(targetTimestamp);
}
@Override
public void onEventTimeWatermark(
long watermarkTimestamp,
Collector<OutputPojo> output,
NonPartitionedContext<OutputPojo> ctx)
throws Exception {
// sense event time watermark arrival
}
@Override
public void onEventTimer(
long timestamp,
Collector<OutputPojo> output,
PartitionedContext<OutputPojo> ctx) {
// write your event timer callback here
}
}
EventTimeProcessFunction을 구현한 후 사용자는 EventTimeUtils#wrapProcessFunction을 사용해 사용자 지정 함수를 감싸야 합니다. 이 단계는 EventTimeManager 인스턴스를 포함한 필요한 구성 요소를 제공하고 타이머 및 기타 기능에 필요한 내장 상태를 선언하므로 필수적입니다.
EventTimeProcessFunction을 감싸는 예시입니다.
NonKeyedPartitionStream stream = ...;
stream.keyBy(x -> x.getKey())
.process(EventTimeExtension.wrapProcessFunction(new CustomEventTimeProcessFunction()))
.process(...);
사용자 지정 ProcessFunction으로 이벤트 시간 관련 워터마크 처리
마찬가지로 사용자는 EventTimeProcessFunction 대신 ProcessFunction을 구현하여 이벤트 시간 관련 워터마크를 처리할 수 있습니다.
이 접근 방식에서 사용자는 ProcessFunction#onWatermark 메서드 내에서 현재 받은 워터마크가 이벤트 시간 워터마크인지 idle 상태 워터마크인지 평가하고 그에 따라 적절한 처리 로직을 실행해야 합니다.
예시는 아래에 제공됩니다.
public static class CustomProcessFunction
implements OneInputStreamProcessFunction<Integer, Integer> {
@Override
public WatermarkHandlingResult onWatermark(
Watermark watermark,
Collector<Integer> output,
NonPartitionedContext<Integer> ctx) throws Exception {
if (EventTimeExtension.isEventTimeWatermark(watermark)) {
// do something as needed
...
return WatermarkHandlingResult.PEEK;
} else if (EventTimeExtension.isIdleStatusWatermark(watermark)) {
// do something as needed
...
return WatermarkHandlingResult.PEEK;
} else {
// do something as needed
...
}
}
}
중요한 점은 ProcessFunction#onWatermark가 이벤트 시간 관련 워터마크를 처리할 때 WatermarkHandlingResult#PEEK를 반환해야 한다는 것입니다. 이는 Flink 프레임워크가 워터마크 정의에 따라 처리 로직을 선택할 것임을 나타냅니다. 이벤트 시간 워터마크와 idle 상태 워터마크의 경우 Flink 프레임워크는 워터마크를 다운스트림으로 전달합니다.
반대로 ProcessFunction#onWatermark가 WatermarkHandlingResult#POP를 반환하면 워터마크는 Flink 프레임워크에 의해 다운스트림으로 전송되지 않습니다. 사용자는 이로 인해 워터마크가 손실될 수 있거나 워터마크를 수동으로 보내야 할 수 있다는 점을 알아야 합니다.
예시
Flink에서 이벤트 타이머 서비스를 사용하는 예시입니다.
NonKeyedPartitionStream<POJO> source = ...;
source.process(
EventTimeExtension.<POJO>newWatermarkGeneratorBuilder(pojo -> pojo.getTimestamp())
.periodicWatermark(Duration.ofMillis(200))
.buildAsProcessFunction()
)
.keyBy(pojo -> pojo.getKey())
.process(EventTimeExtension.wrapProcessFunction(
new OneInputEventTimeStreamProcessFunction<POJO, String>() {
private EventTimeManager eventTimeManager;
@Override
public void initEventTimeProcessFunction(EventTimeManager eventTimeManager) {
// get event time manager instance
this.eventTimeManager = eventTimeManager;
}
@Override
public void processRecord(
POJO record,
Collector<String> output,
PartitionedContext<String> ctx) throws Exception {
...
// register event timer
eventTimeManager.registerTimer(targetTimestamp);
}
@Override
public void onEventTimer(
long timestamp,
Collector<String> output,
PartitionedContext<String> ctx) {
// write your event timer callback here
}
}
)
)
.toSink(...);