조인
조인 (Joining)
데이터 스트림을 조인하는 다양한 방법을 다루는 문서예요. Flink의 DataStream API는 윈도우 조인(Window Join)과 인터벌 조인(Interval Join) 두 가지를 지원해요.
출처: 문서
본문
윈도우 조인 (Window Join)
윈도우 조인은 공통 키(common key)를 공유하고 같은 윈도우 안에 있는 두 스트림의 요소를 조인해요. 이 윈도우들은 윈도우 할당기(window assigner)로 정의할 수 있으며, 두 스트림 모두의 요소에 대해 평가돼요.
양쪽의 요소들은 사용자 정의 JoinFunction 또는 FlatJoinFunction으로 전달되어, 조인 조건을 만족하는 결과를 사용자가 방출(emit)할 수 있어요.
일반적인 사용 방식은 다음과 같이 요약할 수 있어요:
stream.join(otherStream)
.where(<KeySelector>)
.equalTo(<KeySelector>)
.window(<WindowAssigner>)
.apply(<JoinFunction>);
의미론(semantics)에 대한 몇 가지 참고 사항:
- 두 스트림 요소의 쌍별 조합(pairwise combination)을 만드는 것은 내부 조인(inner-join)처럼 동작해요. 즉 한 스트림의 요소는 조인할 다른 스트림의 대응 요소가 없으면 방출되지 않아요.
- 실제로 조인된 요소들은 타임스탬프로 각각의 윈도우 안에 여전히 있는 가장 큰 타임스탬프를 가지게 돼요. 예를 들어 경계가
[5, 10)인 윈도우는 조인된 요소가 타임스탬프 9를 가지게 해요.
다음 섹션에서는 몇 가지 예시 시나리오를 통해 서로 다른 종류의 윈도우 조인이 어떻게 동작하는지 개요를 제공할게요.
텀블링 윈도우 조인 (Tumbling Window Join)
텀블링 윈도우 조인을 수행하면, 공통 키와 공통 텀블링 윈도우를 가진 모든 요소가 쌍별 조합으로 조인되어 JoinFunction 또는 FlatJoinFunction으로 전달돼요. 이것은 내부 조인처럼 동작하므로, 자기 텀블링 윈도우 안에 다른 스트림의 요소가 없는 한 스트림의 요소는 방출되지 않아요!
그림에서 볼 수 있듯이, 크기가 2밀리초인 텀블링 윈도우를 정의하면 [0,1], [2,3], ... 형태의 윈도우가 만들어져요. 그림은 JoinFunction으로 전달될 각 윈도우의 모든 요소의 쌍별 조합을 보여줘요. 텀블링 윈도우 [6,7]에서는 주황색 요소 ⑥과 ⑦과 조인될 초록색 스트림의 요소가 없기 때문에 아무것도 방출되지 않는다는 점에 주의하세요.
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import java.time.Duration;
...
DataStream<Integer> orangeStream = ...;
DataStream<Integer> greenStream = ...;
orangeStream.join(greenStream)
.where(<KeySelector>)
.equalTo(<KeySelector>)
.window(TumblingEventTimeWindows.of(Duration.ofMillis(2)))
.apply (new JoinFunction<Integer, Integer, String> (){
@Override
public String join(Integer first, Integer second) {
return first + "," + second;
}
});
슬라이딩 윈도우 조인 (Sliding Window Join)
슬라이딩 윈도우 조인을 수행하면, 공통 키와 공통 슬라이딩 윈도우를 가진 모든 요소가 쌍별 조합으로 조인되어 JoinFunction 또는 FlatJoinFunction으로 전달돼요. 현재 슬라이딩 윈도우 안에 다른 스트림의 요소가 없는 한 스트림의 요소는 방출되지 않아요! 일부 요소는 한 슬라이딩 윈도우에서는 조인되지만 다른 윈도우에서는 조인되지 않을 수 있다는 점에 주의하세요!
이 예시에서는 크기가 2밀리초인 슬라이딩 윈도우를 사용하고 1밀리초씩 슬라이드해, [-1, 0],[0,1],[1,2],[2,3], … 형태의 슬라이딩 윈도우를 만들어요. x축 아래의 조인된 요소들은 각 슬라이딩 윈도우에 대해 JoinFunction으로 전달되는 것들이에요. 여기서 주황색 ②가 윈도우 [2,3]에서 초록색 ③과 조인되지만, 윈도우 [1,2]에서는 아무것과도 조인되지 않는 것을 볼 수 있어요.
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;
import java.time.Duration;
...
DataStream<Integer> orangeStream = ...;
DataStream<Integer> greenStream = ...;
orangeStream.join(greenStream)
.where(<KeySelector>)
.equalTo(<KeySelector>)
.window(SlidingEventTimeWindows.of(Duration.ofMillis(2) /* size */, Duration.ofMillis(1) /* slide */))
.apply (new JoinFunction<Integer, Integer, String> (){
@Override
public String join(Integer first, Integer second) {
return first + "," + second;
}
});
세션 윈도우 조인 (Session Window Join)
세션 윈도우 조인을 수행하면, *"결합(combined)"*될 때 세션 기준을 충족하는 같은 키의 모든 요소가 쌍별 조합으로 조인되어 JoinFunction 또는 FlatJoinFunction으로 전달돼요. 다시 말해 이것은 내부 조인을 수행하므로, 한 스트림의 요소만 포함하는 세션 윈도우가 있으면 어떤 출력도 방출되지 않아요!
여기서는 각 세션이 최소 1ms의 갭(gap)으로 나뉘는 세션 윈도우 조인을 정의해요. 세 개의 세션이 있고, 처음 두 세션에서는 두 스트림의 조인된 요소가 JoinFunction으로 전달돼요. 세 번째 세션에서는 초록색 스트림에 요소가 없으므로 ⑧과 ⑨는 조인되지 않아요!
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.windowing.assigners.EventTimeSessionWindows;
import java.time.Duration;
...
DataStream<Integer> orangeStream = ...;
DataStream<Integer> greenStream = ...;
orangeStream.join(greenStream)
.where(<KeySelector>)
.equalTo(<KeySelector>)
.window(EventTimeSessionWindows.withGap(Duration.ofMillis(1)))
.apply (new JoinFunction<Integer, Integer, String> (){
@Override
public String join(Integer first, Integer second) {
return first + "," + second;
}
});
인터벌 조인 (Interval Join)
인터벌 조인은 공통 키를 가진 두 스트림(여기서는 A와 B라고 부를게요)의 요소를 조인하는데, 이때 스트림 B의 요소 타임스탬프가 스트림 A 요소 타임스탬프에 대한 상대적 시간 인터벌 안에 있는 경우를 말해요.
이것은 더 공식적으로
b.timestamp ∈ [a.timestamp + lowerBound; a.timestamp + upperBound] 또는
a.timestamp + lowerBound <= b.timestamp <= a.timestamp + upperBound
로 표현할 수 있어요. 여기서 a와 b는 공통 키를 공유하는 A와 B의 요소예요. 하한(lower bound)이 항상 상한(upper bound)보다 작거나 같기만 하면, 하한과 상한 모두 음수 또는 양수일 수 있어요. 인터벌 조인은 현재 내부 조인만 수행해요.
요소의 쌍이 ProcessJoinFunction으로 전달되면, 두 요소 중 더 큰 타임스탬프가 할당돼요(ProcessJoinFunction.Context로 접근 가능).
인터벌 조인은 현재 이벤트 타임(event time)만 지원해요.
위 예시에서는 하한 -2밀리초, 상한 +1밀리초로 'orange'와 'green' 두 스트림을 조인해요. 기본적으로 이 경계들은 포함(inclusive)되지만, .lowerBoundExclusive()와 .upperBoundExclusive()를 적용해 동작을 바꿀 수 있어요.
더 공식적인 표기법을 다시 사용하면, 삼각형으로 표시된 것처럼
orangeElem.ts + lowerBound <= greenElem.ts <= orangeElem.ts + upperBound
로 번역돼요.
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction;
import java.time.Duration;
...
DataStream<Integer> orangeStream = ...;
DataStream<Integer> greenStream = ...;
orangeStream
.keyBy(<KeySelector>)
.intervalJoin(greenStream.keyBy(<KeySelector>))
.between(Duration.ofMillis(-2), Duration.ofMillis(1))
.process (new ProcessJoinFunction<Integer, Integer, String>(){
@Override
public void processElement(Integer left, Integer right, Context ctx, Collector<String> out) {
out.collect(left + "," + right);
}
});