스트리밍 분석
스트리밍 분석 (Streaming Analytics)
이 페이지는 스트리밍 분석에서 핵심이 되는 이벤트 시간(Event Time), 워터마크(Watermark), 윈도우(Window) 개념을 다뤄요. 무한 스트림에서 신뢰할 수 있는 분석 결과를 얻기 위해 시간을 다루는 방법을 배워요.
출처: 문서
본문
이벤트 시간과 워터마크 (Event Time and Watermarks)
소개 (Introduction)
Flink는 세 가지 서로 다른 시간 개념을 명시적으로 지원해요.
- event time (이벤트 시간): 이벤트가 발생한 시간으로, 이벤트를 생성(또는 저장)하는 장치가 기록한 시간이에요
- ingestion time (수집 시간): Flink가 이벤트를 수집하는 순간 기록하는 타임스탬프예요
- processing time (처리 시간): 파이프라인의 특정 연산자가 이벤트를 처리하는 시각이에요
재현 가능한 결과를 얻으려면, 예를 들어 특정 날짜의 첫 거래 시간 동안 주식이 도달한 최대 가격을 계산할 때는 event time을 사용해야 해요. 이렇게 하면 결과가 계산이 수행된 시점에 의존하지 않아요. 이런 종류의 실시간 애플리케이션은 때때로 processing time으로 수행되기도 하지만, 그 경우 결과는 그 시간 동안 실제로 발생한 이벤트가 아니라 우연히 그 시간에 처리된 이벤트에 의해 결정돼요. processing time에 기반한 분석을 계산하면 불일치가 발생하고, 과거 데이터를 재분석하거나 새 구현을 테스트하기 어려워져요.
이벤트 시간 다루기 (Working with Event Time)
event time을 사용하려면 Flink가 이벤트 시간의 진행 상황을 추적하는 데 사용할 Timestamp Extractor와 Watermark Generator도 함께 제공해야 해요. 이에 대해서는 아래의 "워터마크 다루기" 절에서 다룰 거예요. 하지만 먼저 워터마크가 무엇인지 설명해야겠어요.
워터마크 (Watermarks)
워터마크가 왜 필요한지, 어떻게 동작하는지 보여주는 간단한 예를 살펴볼게요. 이 예에는 어느 정도 순서가 뒤섞여 도착하는 타임스탬프가 있는 이벤트 스트림이 있어요. 표시된 숫자는 이벤트가 실제로 발생한 시각을 나타내는 타임스탬프예요. 첫 번째로 도착한 이벤트는 시간 4에 발생했고, 그 다음에는 더 이른 시간 2에 발생한 이벤트가 오는 식이에요.
이제 스트림 정렬기(stream sorter)를 만든다고 상상해 보세요. 이것은 스트림의 각 이벤트가 도착할 때 처리하고, 같은 이벤트를 포함하되 타임스탬프 순서대로 정렬된 새 스트림을 방출하는 애플리케이션이에요.
몇 가지 관찰:
- 스트림 정렬기가 가장 먼저 보는 요소는 4지만, 정렬된 스트림의 첫 요소로 즉시 방출할 수는 없어요. 순서가 뒤섞여 도착했을 수 있고, 더 이른 이벤트가 아직 도착하지 않았을 수 있기 때문이에요. 사실 이 스트림의 미래에 대한 신적인 지식을 가진 것처럼, 스트림 정렬기가 결과를 만들기 전에 적어도 2가 도착할 때까지 기다려야 함을 알 수 있어요. 어느 정도의 버퍼링과 지연이 필요해요.
- 이 작업을 잘못하면 영원히 기다리게 될 수 있어요. 먼저 정렬기는 시간 4의 이벤트를 보고, 그 다음 시간 2의 이벤트를 봤어요. 타임스탬프가 2보다 작은 이벤트가 도착할까요? 아마도요. 아닐 수도 있어요. 영원히 기다려도 1을 못 볼 수 있어요. 결국 용기를 내고 2를 정렬된 스트림의 시작으로 방출해야 해요.
- 그러면 필요한 것은, 어떤 주어진 타임스탬프가 있는 이벤트에 대해 언제 더 이른 이벤트의 도착을 기다리는 것을 멈출지 정의하는 일종의 정책이에요.
이것이 바로 워터마크가 하는 일이에요 — 워터마크는 언제 더 이른 이벤트를 기다리는 것을 멈출지 정의해요.
Flink에서 event time 처리는 watermark generator에 의존하는데, 이 생성기는 **워터마크(watermark)**라고 하는 특별한 타임스탬프가 있는 요소를 스트림에 삽입해요. 시간 t에 대한 워터마크는 스트림이 (아마도) 시간 t까지 완전하다는 주장(assertion)이에요.
이 스트림 정렬기는 언제 기다리는 것을 멈추고 2를 밀어내어 정렬된 스트림을 시작해야 할까요? 타임스탬프가 2 이상인 워터마크가 도착할 때예요.
- 워터마크를 어떻게 생성할지 결정하는 다양한 정책을 상상할 수 있어요. 각 이벤트는 어느 정도 지연 후에 도착하고, 이런 지연은 다양해서 일부 이벤트는 다른 것보다 더 지연돼요. 한 가지 간단한 접근 방식은 이러한 지연이 어떤 최대 지연에 의해 제한된다고 가정하는 것이에요. Flink는 이 전략을 bounded-out-of-orderness 워터마킹이라고 불러요. 더 복잡한 워터마킹 접근 방식을 상상하는 것은 쉽지만, 대부분의 애플리케이션에서는 고정된 지연이 충분히 잘 동작해요.
지연 vs 완전성 (Latency vs. Completeness)
워터마크에 대한 또 다른 생각은, 워터마크가 스트리밍 애플리케이션의 개발자인 당신에게 **지연(latency)**과 완전성(completeness) 사이의 트레이드오프를 제어할 수 있게 해준다는 것이에요. 결과를 만들기 전에 입력에 대한 완전한 지식을 가질 수 있는 여유가 있는 배치 처리와 달리, 스트리밍에서는 결국 더 많은 입력을 보기 위해 기다리는 것을 멈추고 어떤 결과를 만들어야 해요.
짧은 bounded delay로 워터마킹을 공격적으로 구성해, 입력에 대한 상당히 불완전한 지식으로 결과를 만들 위험을 감수할 수 있어요 — 즉, 빠르게 만들어낸 어쩌면 틀린 결과일 수 있어요. 또는 더 오래 기다려 입력 스트림에 대한 더 완전한 지식을 활용하는 결과를 만들 수도 있어요.
초기 결과를 빠르게 만들고, 추가(지연) 데이터가 처리될 때 그 결과에 대한 업데이트를 공급하는 하이브리드 솔루션을 구현하는 것도 가능해요. 이는 일부 애플리케이션에 좋은 접근 방식이에요.
지연 (Lateness)
지연(Lateness)은 워터마크를 기준으로 정의돼요. Watermark(t)는 스트림이 시간 t까지 완전함을 주장해요. 이 워터마크를 따르면서 타임스탬프가 ≤ t인 모든 이벤트는 **지연(late)**된 것이에요.
워터마크 다루기 (Working with Watermarks)
event-time 기반 이벤트 처리를 수행하려면 Flink가 각 이벤트와 연관된 시간을 알아야 하고, 스트림에 워터마크가 포함되어야 해요.
실습에서 사용하는 Taxi 데이터 소스는 이러한 세부 사항을 처리해줘요. 하지만 직접 만든 애플리케이션에서는 이를 직접 처리해야 하는데, 보통 이벤트에서 타임스탬프를 추출하고 필요할 때 워터마크를 생성하는 클래스를 구현해서 해요. 가장 쉬운 방법은 WatermarkStrategy를 사용하는 것이에요.
DataStream<Event> stream = ...;
WatermarkStrategy<Event> strategy = WatermarkStrategy
.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((event, timestamp) -> event.timestamp);
DataStream<Event> withTimestampsAndWatermarks =
stream.assignTimestampsAndWatermarks(strategy);
윈도우 (Windows)
Flink는 매우 표현력이 풍부한 윈도우 의미론을 갖추고 있어요. 이 절에서 배우게 될 내용은 다음과 같아요.
- 윈도우가 어떻게 무한 스트림에서 집계를 계산하는 데 사용되는지
- Flink가 지원하는 윈도우 유형
- 윈도우 집계를 가진 DataStream 프로그램을 구현하는 방법
소개 (Introduction)
스트림 처리를 할 때 다음과 같은 질문에 답하기 위해 스트림의 유한 부분 집합에 대해 집계된 분석을 계산하고 싶은 것은 자연스러운 일이에요.
- 분당 페이지 조회 수
- 주당 사용자별 세션 수
- 센서별 분당 최대 온도
Flink로 윈도우 분석을 계산하는 것은 두 가지 주요 추상화에 의존해요. 이벤트를 윈도우에 할당하는 Window Assigners(필요에 따라 새 윈도우 객체를 생성)와, 윈도우에 할당된 이벤트에 적용되는 Window Functions이에요.
Flink의 윈도우 API에는 윈도우 함수를 언제 호출할지 결정하는 Triggers와 윈도우에 수집된 요소를 제거할 수 있는 Evictors라는 개념도 있어요.
기본 형태로 keyed 스트림에 윈도우를 적용하는 방법은 다음과 같아요.
stream
.keyBy(<key selector>)
.window(<window assigner>)
.reduce|aggregate|process(<window function>);
keyed가 아닌 스트림에도 윈도우를 사용할 수 있지만, 이 경우 처리가 병렬로 수행되지 않는다는 점을 명심하세요.
stream
.windowAll(<window assigner>)
.reduce|aggregate|process(<window function>);
Window Assigners
Flink에는 여러 내장 윈도우 할당자 유형이 있으며, 아래에 설명돼요. 이러한 윈도우 할당자가 어떤 용도로 쓰일 수 있는지, 어떻게 지정하는지의 몇 가지 예시가 있어요.
- 분당 페이지 조회 수:
TumblingEventTimeWindows.of(Duration.ofMinutes(1)) - 10초마다 계산되는 분당 페이지 조회 수:
SlidingEventTimeWindows.of(Duration.ofMinutes(1), Duration.ofSeconds(10)) - 세션 간 간격이 최소 30분인 세션별 페이지 조회 수:
EventTimeSessionWindows.withGap(Duration.ofMinutes(30))
Durations(기간)는 Duration.ofMillis(n), Duration.ofSeconds(n), Duration.ofMinutes(n), Duration.ofHours(n), Duration.ofDays(n) 중 하나로 지정할 수 있어요.
시간 기반 윈도우 할당자(세션 윈도우 포함)는 event time 버전과 processing time 버전이 모두 있어요. 이 두 유형의 시간 윈도우 사이에는 중요한 트레이드오프가 있어요. processing time 윈도우를 사용할 때는 다음 제한을 받아들여야 해요.
- 과거 데이터를 올바르게 처리할 수 없음
- 순서가 뒤섞인 데이터를 올바르게 처리할 수 없음
- 결과가 비결정적
하지만 더 낮은 지연 시간이라는 장점이 있어요.
count 기반 윈도우를 사용할 때는, 이러한 윈도우가 배치가 완료될 때까지 fire되지 않는다는 점을 명심하세요. 부분 윈도우를 처리하기 위해 타임아웃을 지정할 방법은 없지만, custom Trigger로 그 동작을 직접 구현할 수는 있어요.
global window 할당자는 (같은 key를 가진) 모든 이벤트를 같은 global window에 할당해요. 이는 custom Trigger로 직접 custom 윈도우를 할 때만 유용해요. 유용해 보일 수 있는 많은 경우에서, 다른 절에서 설명하는 것처럼 ProcessFunction을 사용하는 것이 더 나을 거예요.
Window Functions
윈도우의 내용을 처리하는 기본 옵션은 세 가지예요.
ProcessWindowFunction을 사용해 배치로 처리 (윈도우 내용이 담긴Iterable이 전달됨)- 각 이벤트가 윈도우에 할당될 때 호출되는
ReduceFunction또는AggregateFunction으로 증분 처리 - 두 가지의 조합:
ReduceFunction또는AggregateFunction의 사전 집계 결과를, 윈도우가 트리거될 때ProcessWindowFunction에 제공
다음은 방법 1과 3의 예시예요. 각 구현은 1분 event time 윈도우에서 각 센서의 최대값을 찾고, (key, end-of-window-timestamp, max_value)를 포함하는 Tuple 스트림을 생성해요.
ProcessWindowFunction 예시
DataStream<SensorReading> input = ...;
input
.keyBy(x -> x.key)
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(1)))
.process(new MyWastefulMax());
public static class MyWastefulMax extends ProcessWindowFunction<
SensorReading, // input type
Tuple3<String, Long, Integer>, // output type
String, // key type
TimeWindow> { // window type
@Override
public void process(
String key,
Context context,
Iterable<SensorReading> events,
Collector<Tuple3<String, Long, Integer>> out) {
int max = 0;
for (SensorReading event : events) {
max = Math.max(event.value, max);
}
out.collect(Tuple3.of(key, context.window().getEnd(), max));
}
}
이 구현에서 주의할 점이 몇 가지 있어요.
- 윈도우에 할당된 모든 이벤트는 윈도우가 트리거될 때까지 keyed Flink state에 버퍼링되어야 해요. 이는 상당히 비용이 들 수 있어요.
- 우리의
ProcessWindowFunction은 윈도우에 대한 정보를 포함하는Context객체를 전달받아요. 그 인터페이스는 다음과 같아요.
public abstract class Context implements java.io.Serializable {
public abstract W window();
public abstract long currentProcessingTime();
public abstract long currentWatermark();
public abstract KeyedStateStore windowState();
public abstract KeyedStateStore globalState();
public abstract <X> void output(OutputTag<X> outputTag, X value);
}
windowState와 globalState는 그 key의 모든 윈도우에 대해 key별, window별, 또는 전역 key 정보를 저장할 수 있는 곳이에요. 예를 들어 현재 윈도우에 대해 무언가를 기록하고, 이후 윈도우를 처리할 때 사용하고 싶다면 유용할 수 있어요.
증분 집계 예시 (Incremental Aggregation Example)
DataStream<SensorReading> input = ...;
input
.keyBy(x -> x.key)
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(1)))
.reduce(new MyReducingMax(), new MyWindowFunction());
private static class MyReducingMax implements ReduceFunction<SensorReading> {
public SensorReading reduce(SensorReading r1, SensorReading r2) {
return r1.value() > r2.value() ? r1 : r2;
}
}
private static class MyWindowFunction extends ProcessWindowFunction<
SensorReading, Tuple3<String, Long, SensorReading>, String, TimeWindow> {
@Override
public void process(
String key,
Context context,
Iterable<SensorReading> maxReading,
Collector<Tuple3<String, Long, SensorReading>> out) {
SensorReading max = maxReading.iterator().next();
out.collect(Tuple3.of(key, context.window().getEnd(), max));
}
}
Iterable이 정확히 하나의 reading을 포함한다는 점에 주목하세요 — MyReducingMax가 계산한 사전 집계된 최대값이에요.
지연 이벤트 (Late Events)
기본적으로 event time 윈도우를 사용할 때 지연 이벤트는 버려져요. 이에 대한 더 많은 제어를 제공하는 윈도우 API의 두 가지 선택적 부분이 있어요.
Side Outputs라는 메커니즘을 사용해, 버려질 이벤트를 대체 출력 스트림으로 수집하도록 마련할 수 있어요. 예시는 다음과 같아요.
OutputTag<Event> lateTag = new OutputTag<Event>("late"){};
SingleOutputStreamOperator<Event> result = stream
.keyBy(...)
.window(...)
.sideOutputLateData(lateTag)
.process(...);
DataStream<Event> lateStream = result.getSideOutput(lateTag);
또한 지연 이벤트가 계속해서 적절한 윈도우(상태가 유지된)에 할당되는 허용 지연(allowed lateness) 간격을 지정할 수도 있어요. 기본적으로 각 지연 이벤트는 윈도우 함수가 다시 호출되게 해요 (때로 late firing이라고 함). 기본 허용 지연은 0이에요. 즉, 워터마크보다 뒤에 있는 요소는 버려지거나(또는 side output으로 보내짐) 돼요. 예를 들어,
stream
.keyBy(...)
.window(...)
.allowedLateness(Duration.ofSeconds(10))
.process(...);
허용 지연이 0보다 크면, 너무 늦어 버려질 이벤트만 (구성된 경우) side output으로 보내져요.
놀라움 (Surprises)
Flink 윈도우 API의 일부 측면은 예상한 대로 동작하지 않을 수 있어요. flink-user 메일링 리스트 등의 자주 묻는 질문에 기반해, 윈도우에 대해 당신을 놀라게 할 수 있는 몇 가지 사실이 있어요.
슬라이딩 윈도우는 복사본을 만든다 (Sliding Windows Make Copies)
슬라이딩 윈도우 할당자는 많은 윈도우 객체를 만들 수 있고, 각 이벤트를 관련된 모든 윈도우로 복사해요. 예를 들어 24시간 길이의 슬라이딩 윈도우가 15분마다 있다면, 각 이벤트는 4 * 24 = 96개의 윈도우로 복사돼요.
시간 윈도우는 Epoch에 정렬된다 (Time Windows are Aligned to the Epoch)
1시간 길이의 processing-time 윈도우를 사용하고 12:05에 애플리케이션을 시작한다고 해서 첫 윈도우가 1:05에 닫히는 것은 아니에요. 첫 윈도우는 55분 길이가 되고 1:00에 닫혀요.
단, tumbling 및 sliding 윈도우 할당자는 윈도우의 정렬을 변경하는 데 사용할 수 있는 선택적 offset 매개변수를 받아요. 자세한 내용은 Tumbling Windows와 Sliding Windows를 참고하세요.
윈도우는 윈도우를 따라갈 수 있다 (Windows Can Follow Windows)
예를 들어 이렇게 하는 것이 동작해요.
stream
.keyBy(t -> t.key)
.window(<window assigner>)
.reduce(<reduce function>)
.windowAll(<same window assigner>)
.reduce(<same reduce function>);
Flink 런타임이 이 병렬 사전 집계를 (ReduceFunction 또는 AggregateFunction을 사용한다면) 알아서 해줄 것이라고 기대할 수도 있지만, 그렇지 않아요.
이것이 동작하는 이유는 시간 윈도우가 생성하는 이벤트가 윈도우의 끝 시각을 기준으로 타임스탬프가 할당되기 때문이에요. 예를 들어 1시간 윈도우가 생성한 모든 이벤트는 한 시간의 끝을 표시하는 타임스탬프를 가져요. 이 이벤트를 소비하는 이후의 윈도우는 이전 윈도우와 같거나 그 배수인 지속 시간을 가져야 해요.
빈 TimeWindow에는 결과가 없다 (No Results for Empty TimeWindows)
윈도우는 이벤트가 할당될 때만 생성돼요. 따라서 주어진 시간 프레임에 이벤트가 없으면 결과가 보고되지 않아요.
지연 이벤트는 지연 병합을 일으킬 수 있다 (Late Events Can Cause Late Merges)
세션 윈도우는 병합될 수 있는 윈도우의 추상화에 기반해요. 각 요소는 처음에 새 윈도우에 할당되고, 그 후 두 윈도우 사이의 간격이 충분히 작을 때마다 윈도우가 병합돼요. 이런 식으로 지연 이벤트가 이전에 분리되어 있던 두 세션을 가르는 간격을 연결해, 지연 병합(late merge)을 만들어낼 수 있어요.
Hands-on
이 절과 함께하는 실습은 Hourly Tips Exercise예요.
추가 읽기 (Further Reading)
- Timely Stream Processing
- Windows