REPLACE USING 플로우를 사용한 부분 스냅샷 교체

Important

이 기능은 베타 버전으로 제공됩니다.

REPLACE USING 플로우는 대상 테이블을 스트리밍 소스와 동기화하여 지정된 키 열과 일치하는 모든 행을 교체하고 나머지 데이터는 변경하지 않습니다.

SEQUENCE BY 열은 업데이트를 정렬하여, 업데이트가 순서 없이 도착하더라도 결과가 정확하도록 합니다. 각 키마다 가장 높은 시퀀스가 승리하며, 더 낮은 시퀀스 행은 이미 목표 내에 있는 더 높은 시퀀스를 덮어쓰지 않습니다. 동일한 키와 같은 시퀀스를 공유하는 행은 교체가 아니라 덧붙여집니다.

REPLACE USING 작동 원리

두 지역별 클릭 및 전환 이벤트를 포함하는 이벤트 테이블을 생각해 보자. 순서는 seq다음과 같다:

region_id device_type event_type seq
1 iOS click 1
1 Android 변환 1
2 iOS click 1
2 데스크톱 click 1

플로우는 REPLACE USING (region_id) SEQUENCE BY seq 영역 1과 3에 대한 이러한 업데이트를 받습니다. 지역 2에는 업데이트가 없습니다:

region_id device_type event_type seq
1 iOS click 2
1 Android 변환 2
1 데스크톱 click 2
3 iOS click 1
3 데스크톱 click 2

목표는 다음과 같습니다:

region_id device_type event_type seq 결과
1 iOS click 2 대체된 이유는 순차 2가 순차 1보다 크기 때문입니다
1 Android 변환 2 대체된 이유는 순차 2가 순차 1보다 크기 때문입니다
1 데스크톱 click 2 대체된 이유는 순차 2가 순차 1보다 크기 때문입니다
2 iOS click 1 이번 업데이트에는 키가 없기 때문에 손대지 않은 상태입니다
2 데스크톱 click 1 이번 업데이트에는 키가 없기 때문에 손대지 않은 상태입니다
3 데스크톱 click 2 추가. 영역 3의 seq 1 행은 추가되지 않는데, 키에 대해 가장 높은 시퀀스만 적용되기 때문입니다.

요구 사항

REPLACE USING 흐름에는 다음과 같은 요구 사항이 있습니다:

  • 데이터브릭스 런타임 18.2 이상, 클래식 또는 서버리스 컴퓨트에서 실행되는 흐름을 대체합니다. Databricks는 Unity Catalog를 추천합니다.
  • 원본은 스트리밍 원본이어야 합니다. REPLACE USING은 비스트리밍 소스를 허용하지 않습니다.
  • 최소 하나의 키 열과 정확히 한 개의 SEQUENCE BY 열을 지정해야 합니다.

REPLACE USING 플로우를 사용해야 하는 경우

Lakeflow 파이프라인은 기존 행을 덮어쓰는 세 가지 흐름을 제공합니다. 소스가 어떻게 생겼는지, 그리고 교체할 행을 어떻게 식별하는지에 따라 선택하세요:

  • 소스가 열을 기준으로 키가 지정된 일련의 부분 스냅샷인 경우 REPLACE USING을 사용하세요. REPLACE USING은 입력 데이터와 일치하는 데이터만 덮어쓰고, 나머지 데이터는 그대로 둡니다. 기본 키가 필요하지 않습니다.
  • 소스 소스가 명시적인 삽입, 업데이트, 삭제 작업이 있는 변경 데이터 캡처(CDC) 피드이거나 천천히 차원 변경(SCD) 타입 2 이력이 필요할 때는 AUTO CDC를 사용하세요. AUTO CDC도 진정한 기본 키가 필요합니다. AUTO CDC API: 파이프라인을 사용하여 변경 데이터 캡처 간소화를 참조하세요.
  • 소스가 스냅샷이고, 조건식으로 선택한 대상 테이블 범위(예: 최근 7일)를 배치 작업으로 다시 계산해 덮어쓰려는 경우 REPLACE WHERE를 사용하세요. 기본 키가 필요하지 않습니다. REPLACE WHERE 흐름을 사용한 Batch 처리를 참조하세요.

REPLACE USING 흐름을 만드세요

SQL 또는 Python에서 REPLACE USING 플로우를 정의하세요.

SQL

FLOW REPLACE USING와 함께 CREATE STREAMING TABLE 절을 인라인으로 사용하세요:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

또는 긴 형식 CREATE FLOW 구문을 사용합니다.

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

메모

BY NAME 는 SQL에 필요합니다. 위치가 아니라 이름을 기준으로 열을 매칭합니다.

Python

테이블과 플로우를 @dp.table함께 선언합니다:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

또는 @dp.replace_flow를 사용하여 기존 스트리밍 테이블을 대상으로 지정할 수도 있습니다:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using 는 키 열의 목록입니다. sequence_by는 열 이름 또는 Column 식이며, replace_using가 설정될 때마다 필요합니다.

시퀀싱 및 순서가 뒤바뀐 데이터

이 열은 SEQUENCE BY 업데이트가 도착하는 순서와 무관하게 결과를 만듭니다. 행은 해당 키에 적용되는 시퀀스가 이미 저장된 시퀀스보다 클 때만 적용되므로, 현재 값보다 오래된 늦거나 재재생된 행은 무시됩니다. 업데이트에 없는 키는 그대로 남겨둡니다.

교체가 예측 가능하게 작동하도록 다음 절차를 따르세요:

연습 이유
타임스탬프, 버전 번호, 로그 오프셋 등 키 버전마다 엄격히 증가하는 순서를 사용하세요. 같은 키와 같은 시퀀스를 가진 두 행이 모두 유지되어, 해당 키에 중복된 행이 생깁니다.
null이 아닌 시퀀스를 사용하세요. 널 시퀀스는 정의되지 않은 동작을 초래할 수 있습니다.

Expectations

흐름을 사용해 지원 기대치를 대체합니다. warnfail은 다른 흐름에서와 마찬가지로 동작합니다. warn는 위반 행을 그대로 두고 위반을 기록하며, fail는 업데이트를 중단합니다. 파이프라인 기대를 사용하여 데이터 품질을 관리하기를 참조하세요.

drop 검증 조건은 조건을 위반한 행을 소스가 아예 생성하지 않은 것처럼 처리합니다. 삭제된 행은 대상 테이블에서 일치하는 키를 대체하거나 삭제하거나 수정하지 않습니다:

  • 드롭은 중복 제거 전에 이루어지므로, 플로우는 키의 최신 유효 버전을 유지합니다.
  • 만약 키의 모든 들어오는 행이 삭제되면, 그 키의 기존 행은 그대로 남습니다.
  • 삭제된 행은 시퀀스 플로어를 설정하지 않으므로, 나중에 유효한 업데이트가 삭제된 행보다 낮더라도 여전히 발생한다.

Limitations

REPLACE USING 흐름에는 다음과 같은 제한 사항이 있습니다:

  • REPLACE USING 는 대상 테이블당 단일 플로우를 지원합니다. 같은 타겟에서 REPLACE USING 타입과 다른 플로우 유형을 결합하는 것은 지원되지 않습니다.
  • 대상 테이블은 파이프라인 내에서 만들어야 합니다.
  • 원본은 스트리밍 원본이어야 합니다.
  • 최소 하나의 키 열과 하나의 SEQUENCE BY 열을 지정해야 합니다. 키 열은 반복할 수 없으며, 각 키 열의 타입은 정렬 가능해야 합니다. 정수, 문자열, 날짜와 같은 원자 타입은 키가 될 수 있지만 MAP , 와 VARIANT 는 그렇지 않습니다.
  • 독립 실행형 스트리밍 테이블의 경우 구문 차이는 REPLACE USING flows를 사용한 부분 스냅샷 대체 적용을 참조하세요.

예제

다음 예시들은 모든 Unity Catalog 지원 작업 공간에서 제공되는 예약 상태 변경 샘플 표에서 samples.wanderbricks.booking_updates가져옵니다. 각 예약은 변경될 때마다 한 번씩 표시되므로, booking_id가 새로운 booking_update_id와 함께 반복됩니다. Wanderbricks 데이터셋을 참조하세요.

예시 1: 각 키에 대한 최신 기록을 유지합니다

각 예약의 현재 상태만 유지하세요. 플로우는 booking_id를 기준으로 키를 지정하고 booking_update_id를 기준으로 시퀀스를 정하므로, 해당 예약의 가장 최근 업데이트가 이전 업데이트를 대체합니다. 소스가 명시적 삽입, 업데이트 및 삭제 작업이 명시된 변경 피드인 경우 대신 AUTO CDC를 사용하세요.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

이 예시에서는 updated_at 타임스탬프가 아니라 booking_update_id를 기준으로 순서를 정하는데, 동일한 예약에 대한 여러 업데이트가 같은 타임스탬프를 공유할 수 있기 때문입니다. 순서와 일치하는 행은 교체되지 않고 덧붙여져 예약을 위해 한 행 이상이 남게 됩니다.

예시 2: 여러 열에 키가 있음

레코드가 여러 열의 조합으로 식별되는 경우, REPLACE USING에 모두 나열하세요. 여기서 각 예약은 (property_id, booking_id)로 식별되므로, 플로우는 숙소별로 각 예약의 현재 상태를 유지합니다. 키 열이 null일 수 있는 경우, REPLACE USING은 행을 건너뛰지 않고 null을 null과 일치시킵니다.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

예시 3: 기대값을 사용하여 유효하지 않은 레코드 삭제

잘못된 행이 대상에 들어가지 않도록 하는 익스펙테이션을 추가하세요. 삭제된 행은 소스가 해당 키를 생성하지 않은 것처럼 처리합니다: 일치하는 키를 대체하거나 삭제하지 않고, 흐름은 해당 키에 대한 최신 유효 행으로 폴트됩니다. 이 흐름은 양수 total_amount가 없는 업데이트를 제외합니다.

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")