Engineering Note

늦게 온 이벤트가 최신 상태를 덮어쓰지 않게: Fluss Versioned Merge Engine

재시도와 지연 도착으로 과거 이벤트가 최신 상태를 덮는 문제를 Apache Fluss Versioned Merge Engine으로 막는 방법을 Flink SQL 예제와 함께 설명합니다.

2026년 8월 24일 · Pletor Engineering apache-flussflinkstreamingdata-consistencyevent-processing
청록색 최신 데이터 레코드는 중앙 상태 저장소에 기록되고, 늦게 온 코랄색 과거 레코드는 비교 계층에서 걸러지는 등각형 일러스트
Versioned Merge Engine은 늦게 도착한 과거 상태가 현재 상태를 덮지 않게 합니다.

주문 상태 PAID를 기록한 뒤, 재시도나 지연 때문에 더 오래된 CONFIRMED 이벤트가 도착했다고 해 보겠습니다. 기본 키 테이블이 도착한 순서대로 행을 갱신하면, 원본 시스템에서는 이미 결제가 끝난 주문이 다시 확인 단계로 돌아간 것처럼 보일 수 있습니다.

Apache Fluss의 Versioned Merge Engine은 이 문제에서 현재 상태를 고르는 기준을 도착 순서에서 버전 값으로 바꿉니다. 같은 기본 키의 행이 여러 번 들어오면 가장 큰 버전 또는 시각을 가진 행만 남깁니다. 늦게 온 과거 이벤트가 최신 상태를 덮지 못하게 하는 기능입니다.

먼저 도착한 버전 20의 PAID 상태와 나중에 도착한 버전 19의 CONFIRMED 상태를 비교해, Versioned Merge Engine이 버전 20을 현재 상태로 유지하는 구조도
Versioned Merge Engine은 도착 순서가 아니라 원본이 부여한 버전으로 현재 상태를 고릅니다.

이 글은 Flink SQL로 최소 예제를 실행해 동작을 확인하고, 어떤 값에 버전을 부여해야 하는지와 이 기능이 해결하지 않는 범위를 함께 정리합니다. 핵심은 단순합니다. 최신 상태를 저장하려면 “마지막으로 도착한 값”과 “원본에서 가장 최신인 값”을 구분해야 합니다.

문제는 늦은 이벤트가 아니라 상태를 고르는 기준이다

이벤트가 네트워크를 거쳐 도착하는 순서는 원본 시스템의 변경 순서와 다를 수 있습니다. 작업 재시도, 장애 뒤 재처리, 여러 입력 경로의 지연은 모두 이전 변경을 뒤늦게 보낼 수 있습니다.

예를 들어 주문 서비스가 다음 순서로 상태를 만들었다고 하겠습니다.

원본 시스템의 변경 순서

source_version=19  CONFIRMED
source_version=20  PAID

그런데 Fluss 테이블에는 PAID가 먼저 쓰이고, 재시도된 CONFIRMED가 나중에 쓰일 수 있습니다. 기본 병합 방식인 LastRow는 뒤에 기록된 행을 최신 행으로 취급합니다. 원본의 변경 순서와 테이블의 쓰기 순서가 같다는 보장이 있을 때는 자연스러운 선택입니다. 하지만 이 예제처럼 순서가 달라질 수 있다면, 그 기준만으로는 현재 상태를 안전하게 표현할 수 없습니다.

Versioned Merge Engine은 기본 키별로 지정한 열을 비교합니다. 더 큰 값이 이미 저장되어 있다면 작은 값으로 들어온 행은 현재 상태를 바꾸지 않습니다. 공식 문서는 이 방식이 순서가 바뀐 데이터와 최종 일관성이 필요한 데이터를 병합하는 데 쓰인다고 설명합니다. Versioned Merge Engine 문서에서 지원하는 버전 열 형식과 제약도 확인할 수 있습니다.

LastRow, FirstRow, Versioned는 무엇이 다른가

Fluss 기본 키 테이블은 병합 엔진으로 같은 키의 여러 쓰기를 어떻게 하나의 현재 행으로 만들지 정합니다. 중요한 것은 엔진 이름보다 어떤 순서를 신뢰할 것인지입니다.

병합 엔진같은 키의 여러 행에서 남기는 값어울리는 상황
LastRow테이블에 마지막으로 기록된 행쓰기 도착 순서가 곧 업무상 최신 순서일 때
FirstRow테이블에 처음 기록된 행같은 사실의 중복 기록을 처음 값 하나로 막을 때
Versioned지정한 버전 열이 가장 큰 행재시도·재처리·지연 도착으로 쓰기 순서가 바뀔 수 있을 때

FirstRow는 처음 본 이벤트를 유지하는 중복 제거에, Versioned는 여러 상태 후보 중 원본 기준으로 더 새로운 상태를 고르는 데 초점이 있습니다. 둘 다 “중복”이라는 단어와 함께 설명되기 쉽지만, 판단 기준이 다릅니다. FirstRow는 최초 도착 여부를, Versioned는 버전 값을 봅니다. FirstRow Merge Engine 문서도 FirstRow가 첫 행만 유지하며 UPDATE, DELETE, 부분 갱신을 지원하지 않는다고 명시합니다.

아래 예제는 Fluss catalog를 선택한 Flink SQL 세션에서 실행합니다. order_id가 기본 키이고, source_version이 원본 주문 서비스가 부여한 단조 증가 버전이라고 가정합니다.

SET 'execution.runtime-mode' = 'batch';
SET 'table.dml-sync' = 'true';

CREATE TABLE order_status (
  order_id STRING NOT NULL,
  status STRING,
  source_version BIGINT NOT NULL,
  source_updated_at TIMESTAMP(3),
  PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
  'bucket.num' = '1',
  'table.merge-engine' = 'versioned',
  'table.merge-engine.versioned.ver-column' = 'source_version'
);

table.dml-synctrue로 두면 각 INSERT 작업이 끝난 뒤 다음 명령을 실행합니다. 그래서 바로 뒤의 조회가 아직 끝나지 않은 쓰기 작업을 앞지르지 않습니다. source_updated_at은 결과를 읽을 때 원본 변경 시각을 함께 보기 위해 넣은 열입니다. 병합 판단에는 source_version만 사용합니다. Versioned Merge Engine은 INT, BIGINT, TIMESTAMP, TIMESTAMP_LTZ 계열을 버전 열로 사용할 수 있습니다. 다만 주문 상태처럼 변경 순서가 중요하다면, 시간보다 원본이 보장하는 증가 버전이 더 분명한 기준인 경우가 많습니다.

먼저 최신 상태를 기록합니다.

INSERT INTO order_status VALUES (
  'o-1042', 'PAID', 20, TIMESTAMP '2026-08-14 09:00:00'
);

이제 재시도되었거나 늦게 도착한 이전 상태를 기록합니다. 테이블 입장에서는 이 행이 더 나중에 들어오지만, 원본 버전은 더 작습니다.

INSERT INTO order_status VALUES (
  'o-1042', 'CONFIRMED', 19, TIMESTAMP '2026-08-14 08:59:00'
);

SELECT order_id, status, source_version
FROM order_status
WHERE order_id = 'o-1042';

조회 결과는 다음과 같습니다.

+----------+--------+----------------+
| order_id | status | source_version |
+----------+--------+----------------+
| o-1042   | PAID   |             20 |
+----------+--------+----------------+

CONFIRMED 행을 없애거나 별도 오류로 처리한 것이 아닙니다. Versioned Merge Engine이 이 기본 키의 현재 상태에는 버전 20이 더 새롭다고 판단했기 때문에, 버전 19가 현재 행을 바꾸지 않은 것입니다.

버전 열은 이벤트 ID와 다르다

event_idsource_version을 같은 용도로 쓰면 설계가 흐려집니다.

  • event_id는 “이 이벤트를 전에 처리했는가”를 판단하는 식별자입니다. 같은 사실을 재전송했는지 확인하는 데 적합합니다.
  • source_version은 “서로 다른 두 상태 중 어느 쪽이 더 나중의 상태인가”를 판단하는 순서 값입니다.

한 주문에 CONFIRMEDPAID가 각각 한 번씩만 왔다고 해도, 둘은 중복 이벤트가 아니라 서로 다른 상태 변경입니다. 이 경우 이벤트 ID로는 어느 상태를 현재 값으로 남길지 정할 수 없습니다. 반대로 버전만으로는 같은 이벤트가 여러 번 전송됐는지 추적하기 어렵습니다. 필요하다면 두 값을 함께 보관해야 합니다.

가장 좋은 버전 열은 원본 시스템이 이미 보장하는 단조 증가 값입니다. 주문 행의 수정 번호, 변경 시퀀스, 데이터베이스 변경 로그의 위치처럼 변경 순서를 뜻하는 값이 여기에 해당합니다. 원본 변경 시각을 쓰는 방법도 있지만, 같은 시각의 두 변경을 어떻게 정렬할지와 시계 오차를 먼저 정해야 합니다. Versioned Merge Engine은 같은 버전의 행 가운데 나중에 쓴 행도 갱신 대상으로 보므로, 동률이 가능한 시간값만으로는 결정 규칙이 부족할 수 있습니다.

이 기능이 해결하지 않는 범위

Versioned Merge Engine은 기본 키 테이블 안에서 현재 행을 고르는 기능입니다. 다음 문제까지 대신 해결하지는 않습니다.

  • 모든 상태 변화를 보관하는 일: PAID만 남기는 현재 상태 테이블에서는 CONFIRMED가 언제 들어왔는지 조회할 수 없습니다. 상태 전이 감사나 재처리가 필요하면 원본 이벤트를 Log Table에도 추가 전용으로 기록해야 합니다.
  • 외부 시스템과의 원자적 처리: Fluss 테이블의 병합 결과가 외부 DB 갱신이나 결제 API 호출까지 원자적으로 묶어 주지는 않습니다. 그 경계에는 별도의 멱등성, 결과 조회, Outbox 같은 설계가 필요할 수 있습니다.
  • 부분 갱신과 삭제: Versioned Merge Engine은 SQL UPDATE, DELETE, 부분 갱신을 지원하지 않습니다. 필요한 갱신 상태를 모두 담은 행을 INSERT로 기록하는 모델에 맞춰야 합니다.

따라서 보통은 두 테이블을 나란히 둡니다.

order_status_events      Log Table
  모든 상태 변경을 추가 전용으로 보관한다.

order_status             Versioned Primary Key Table
  원본 버전이 가장 큰 행을 현재 상태로 제공한다.

이 구성이면 감사·재처리·상태 전이 분석은 Log Table에서 하고, 서비스 조회나 최신 상태 조인은 Versioned 기본 키 테이블에서 처리할 수 있습니다. Log Table과 Primary Key Table의 역할 차이는 Apache Fluss란: Kafka·Flink·Iceberg 사이의 스트리밍 레이크하우스에서 더 자세히 다뤘습니다.

도입 전에 확인할 질문

Versioned Merge Engine은 단순히 기본 키 테이블에 붙이는 안전장치가 아닙니다. 다음 질문에 “예”라고 답할 수 있을 때 의미가 큽니다.

  • 원본 시스템이 키별로 비교 가능한 버전 또는 신뢰할 수 있는 변경 시각을 제공하는가?
  • 재시도, 재처리, 여러 입력 경로 때문에 예전 이벤트가 늦게 도착할 수 있는가?
  • 소비자가 모든 중간 상태보다 최신 상태 조회를 더 자주 필요로 하는가?
  • 중간 상태도 필요하다면, 이를 별도 Log Table에 보관하고 있는가?

운영에서는 늦은 이벤트가 무조건 오류라는 전제도 버려야 합니다. 늦은 이벤트 비율, 버전 역전 횟수, 버전이 없는 입력 비율을 애플리케이션 또는 Flink 작업의 메트릭으로 관찰하면, 병합 엔진이 조용히 가리고 있는 입력 품질 문제를 발견할 수 있습니다. Konduo는 여러 IT 인프라의 상태, 메트릭·로그 근거, 알림 대응을 한 운영 흐름으로 연결하는 플랫폼입니다. Fluss를 포함한 데이터 경로를 운영할 때도 현재 상태만 보지 않고, 그 상태에 이르는 입력 흐름을 함께 관찰하는 관점이 필요합니다.

마무리

기본 키 테이블에 최신 상태를 저장한다고 해서 자동으로 원본의 최신 상태가 남는 것은 아닙니다. 마지막으로 도착한 행을 최신으로 볼 수 있는지, 아니면 원본의 버전으로 판단해야 하는지를 먼저 정해야 합니다.

재시도와 지연 도착이 현실적인 경로라면 Versioned Merge Engine은 현재 상태 테이블을 훨씬 명확하게 만듭니다. 다만 이것을 이벤트 이력, 외부 부수 효과, 전체 처리의 정확성을 보장하는 장치로 확대 해석해서는 안 됩니다. Log Table에는 사실을 남기고, Versioned 기본 키 테이블에는 원본 버전이 가장 큰 현재 상태를 남기는 역할 분리가 가장 안전한 출발점입니다.

함께 읽기 좋은 글