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 행은 추가되지 않는데, 키에 대해 가장 높은 시퀀스만 적용되기 때문입니다. |
요구 사항
REPLACEMENT USING FLOWS는 다음과 같은 요구사항을 가집니다:
- 데이터브릭스 런타임 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
기대치를 지원하는 흐름 사용 대체.
warn
fail 그리고 다른 흐름 warn 에서처럼 계속 위반 행을 기록 fail 하고 업데이트를 중단합니다.
파이프라인 기대를 사용하여 데이터 품질을 관리하기를 참조하세요.
기대는 drop 위반 논란을 출처가 만들어내지 않은 것처럼 취급합니다. 삭제된 행은 대상 테이블에서 일치하는 키를 대체하거나 삭제하거나 수정하지 않습니다:
- 드롭은 중복 제거 전에 이루어지므로, 플로우는 키의 최신 유효 버전을 유지합니다.
- 만약 키의 모든 들어오는 행이 삭제되면, 그 키의 기존 행은 그대로 남습니다.
- 삭제된 행은 시퀀스 플로어를 설정하지 않으므로, 나중에 유효한 업데이트가 삭제된 행보다 낮더라도 여전히 발생한다.
Limitations
REPLACEMENT USING 플로우에는 다음과 같은 제한이 있습니다:
- REPLACE USING 는 대상 테이블당 단일 플로우를 지원합니다. 같은 타겟에서 REPLACE USING 타입과 다른 플로우 유형을 결합하는 것은 지원되지 않습니다.
- 대상 테이블은 파이프라인 내에서 만들어야 합니다.
- 원본은 스트리밍 원본이어야 합니다.
- 최소 하나의 키 열과 하나의
SEQUENCE BY열을 지정해야 합니다. 키 열은 반복할 수 없으며, 각 키 열의 타입은 정렬 가능해야 합니다. 정수, 문자열, 날짜와 같은 원자 타입은 키가 될 수 있지만MAP, 와VARIANT는 그렇지 않습니다.
예제
다음 예시들은 모든 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")