중복 제거

중복 제거 (Deduplication)

중복 제거는 특정 컬럼 집합에 대해 중복된 행을 제거하여 첫 번째 행 또는 마지막 행만 유지하는 연산입니다. 어떤 경우에는 업스트림 ETL 작업이 end-to-end 정확히-한 번(exactly-once)이 아니어서 장애 복구 시 싱크에 중복 레코드가 생길 수 있습니다. 그러나 중복 레코드는 SUM, COUNT 같은 다운스트림 분석 작업의 정확성에 영향을 주므로, 추가 분석 전에 중복 제거가 필요합니다.

출처: 문서

본문

Flink는 Top-N 쿼리와 같은 방식으로 ROW_NUMBER()를 사용해 중복을 제거합니다. 이론적으로 중복 제거는 N이 1이고 처리 시간이나 이벤트 시간으로 순서를 정하는 Top-N의 특수한 경우입니다.

다음은 중복 제거 문의 문법입니다:

SELECT [column_list]
FROM table_name
QUALIFY ROW_NUMBER() OVER ([PARTITION BY col1[, col2...]] ORDER BY time_attr [asc|desc]) = 1

파라미터 설명:

  • ROW_NUMBER(): 각 행에 1부터 시작하는 고유한 순차 번호를 할당합니다.
  • PARTITION BY col1[, col2...]: 파티션 컬럼, 즉 중복 제거 키를 지정합니다.
  • ORDER BY time_attr [asc|desc]: 순서를 지정하는 컬럼으로, 반드시 시간 속성(time attribute)이어야 합니다. 현재 Flink는 처리 시간 속성과 이벤트 시간 속성을 지원합니다. ASC 정렬은 첫 번째 행을 유지하고, DESC 정렬은 마지막 행을 유지함을 뜻합니다.
  • WHERE rownum = 1: Flink가 이 쿼리를 중복 제거로 인식하려면 rownum = 1이 필요합니다.

참고: 위 패턴은 정확히 지켜야 합니다. 그렇지 않으면 옵티마이저가 쿼리를 변환하지 못합니다.

다음 예시들은 스트리밍 테이블에서 중복 제거를 사용하는 SQL 쿼리를 보여줍니다.

CREATE TABLE Orders (
  order_id  STRING,
  user        STRING,
  product     STRING,
  num         BIGINT,
  proctime AS PROCTIME()
) WITH (...);

-- remove duplicate rows on order_id and keep the first occurrence row,
-- because there shouldn't be two orders with the same order_id.
SELECT order_id, user, product, num
FROM (
  SELECT *,
    ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY proctime ASC) AS row_num
  FROM Orders)
WHERE row_num = 1

더 알아보기 (Learn more)