조인
조인 (Joining)
Join은 두 데이터 스트림을 공통 key로 요소를 매칭해 합치고, 매칭된 요소에 대해 계산을 수행하는 데 사용돼요. DataStream API V2에서 Join 연산을 자세히 소개해요.
출처: 문서
본문
참고: DataStream API V2는 기존 DataStream API를 점진적으로 대체하기 위한 새로운 API 집합이에요. 현재 실험 단계이며 프로덕션에는 완전히 사용할 수 없어요.
Join은 두 데이터 스트림을 공통 key를 기준으로 양쪽 스트림의 요소를 매칭해 합치고, 매칭된 요소에 대해 계산을 수행하는 데 사용돼요. 이 절에서는 DataStream의 Join 연산을 자세히 소개해요.
현재 DataStream은 non-window INNER join만 지원한다는 점에 유의하세요. interval join, lookup join, window join 같은 다른 유형의 조인은 나중에 지원될 거예요.
비윈도우 조인 (Non-Window Join)
Join을 사용하려면 사용자는 다음을 지정해야 해요.
- 조인될 좌측 및 우측 스트림
- 조인 key를 결정하는 두 스트림의 KeySelector
- 매칭된 데이터에 적용되는 처리 로직인 JoinFunction
먼저 Join을 수행하는 API를 소개하고, 그다음 JoinFunction을 소개할게요.
Join을 수행하는 API
Join 연산을 수행하는 두 가지 접근 방식이 있어요.
BuiltinFuncs.join을 사용해 두 스트림을 연결하고 Join을 실행- JoinFunction을 ProcessFunction으로 변환하고
KeyedPartitionStream#connectAndProcess를 활용해 Join을 실행
첫 번째 접근 방식에서 사용자는 입력 데이터 스트림의 유형에 따라 두 가지 옵션 중에 선택할 수 있어요.
- 두 입력 데이터 스트림이 모두 이미 Keyed Partition Stream이라면, 두 Keyed Partition Stream과 JoinFunction을 직접 결합해 조인된 데이터 스트림을 만들 수 있어요.
NonKeyedPartitionStream joinedStream = BuiltinFuncs.join(
keyedStream1,
keyedStream2,
new CustomJoinFunction()
);
- 두 입력 데이터 스트림이 모두 Non-Keyed Partition Stream이라면, 사용자는 해당 Join key KeySelector와 JoinFunction을 포함해 두 Non-Keyed Partition Stream을 조인된 데이터 스트림으로 변환해야 해요.
NonKeyedPartitionStream joinedStream = BuiltinFuncs.join(
stream1,
new CustomJoinKeySelector1(),
stream2,
new CustomJoinKeySelector2(),
new CustomJoinFunction()
);
두 번째 접근 방식에서 사용자는 JoinFunction을 ProcessFunction으로 변환할 수 있으며, 그 후 DataStream API가 처리할 수 있어요. 이 변환의 예시와 변환된 ProcessFunction을 사용하는 방법은 아래와 같아요.
TwoInputNonBroadcastStreamProcessFunction wrappedJoinFunction = BuiltinFuncs.join(new CustomJoinFunction());
NonKeyedPartitionStream joinedStream = keyedStream1.connectAndProcess(
keyedstream2,
wrappedJoinFunction
);
두 번째 접근 방식에서는 두 입력 데이터 스트림이 Keyed Partition Stream이어야 한다는 점을 유의하세요.
JoinFunction
JoinFunction은 매칭된 데이터를 어떻게 계산할지 설명하는 데 사용되는 인터페이스예요. processRecord 메서드 하나만 있어요. 사용자는 processRecord에서 매칭된 요소를 얻어 계산을 수행한 뒤 계산 결과를 출력할 수 있어요.
아래는 JoinFunction을 사용해 학생 개인 정보와 시험 점수를 연결하는 예시예요.
class JoinStudentInformationAndScore
implements JoinFunction<StudentInfo, ExamScore, EnrichedStudentExamScore> {
@Override
public void processRecord(
StudentInfo studentInfo,
ExamScore examScore,
Collector<EnrichedStudentExamScore> output,
RuntimeContext ctx)
throws Exception {
// do some calculation logic and emit joined result
EnrichedStudentExamScore studentExamScore =
new EnrichedStudentExamScore(studentInfo.getId(), studentInfo.getName(), examScore.getScore());
output.collect(studentExamScore);
}
}