조인

조인 (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);
    }
}

더 알아보기 (Learn more)