Częściowa wymiana snapshotów przepływami REPLACE USING

Ważna

Ta funkcja jest dostępna w wersji beta.

Przepływ REPLACE USING utrzymuje tabelę docelową zsynchronizowaną ze źródłem danych strumieniowych: zastępuje wszystkie wiersze zgodne z określonymi kolumnami kluczowymi i pozostawia wszystkie pozostałe dane bez zmian.

Kolumna SEQUENCE BY sortuje aktualizacje tak, aby wynik był poprawny nawet wtedy, gdy aktualizacje docierają w niewłaściwej kolejności. Dla każdego klucza decyduje najwyższy numer sekwencji, a wiersz o niższym numerze sekwencji nigdy nie nadpisuje wiersza o wyższym numerze, który już znajduje się w tabeli docelowej. Wiersze o tym samym kluczu i tej samej sekwencji są dodawane, a nie zastępowane.

Jak działa REPLACE USING

Rozważmy tabelę zdarzeń, która zawiera zdarzenia kliknięć i konwersji dla dwóch regionów, ułożone w sekwencję wzorem seq:

region_id typ_urządzenia event_type seq
1 iOS click 1
1 Android konwersja 1
2 iOS click 1
2 biurko click 1

REPLACE USING (region_id) SEQUENCE BY seq przepływ otrzymuje te aktualizacje dla regionów 1 i 3. Region 2 nie ma żadnych aktualizacji:

region_id typ_urządzenia event_type seq
1 iOS click 2
1 Android konwersja 2
1 biurko click 2
3 iOS click 1
3 biurko click 2

Celem jest:

region_id typ_urządzenia event_type seq Outcome
1 iOS click 2 Zastąpiony, ponieważ numer sekwencyjny 2 jest większy niż numer sekwencyjny 1
1 Android konwersja 2 Zastąpiono, ponieważ sekwencja 2 jest większa niż sekwencja 1
1 biurko click 2 Zastąpiony, bo sekw. 2 jest większy niż sekw. 1
2 iOS click 1 Niezmienione, ponieważ klucz nie jest obecny w tej aktualizacji
2 biurko click 1 Bez zmian, ponieważ tego klucza nie ma w tej aktualizacji
3 biurko click 2 Dodano. Wiersz z sekwencją 1 dla regionu 3 nie zostaje dodany, ponieważ dla klucza stosowana jest tylko najwyższa sekwencja.

Requirements

Przepływy REPLACE USING mają następujące wymagania:

  • Przepływy REPLACE USING działają w środowisku Databricks Runtime 18.2 i nowszym, na klasycznych lub bezserwerowych zasobach obliczeniowych. Databricks poleca Unity Catalog.
  • Źródło musi być źródłem przesyłania strumieniowego. REPLACE USING odrzuca źródło niebędące streamingem.
  • Musisz określić co najmniej jedną kolumnę klucza i dokładnie jedną kolumnę SEQUENCE BY .

Kiedy używać przepływów REPLACE USING

Rurociągi Lakeflow oferują trzy przepływy, które nadpisują istniejące rzędy. Wybierz na podstawie wyglądu źródła i sposobu identyfikacji wierszy do zastąpienia:

  • Używaj REPLACE USING, gdy źródłem jest seria częściowych migawek identyfikowanych według kolumny. REPLACE USING nadpisuje tylko dane, które mają zgodność z danymi wejściowymi, pozostawiając wszystkie pozostałe dane nietknięte. Nie wymaga klucza głównego.
  • Używaj AUTO CDC, gdy źródłem jest strumień przechwytywania zmian danych (CDC) z jawnymi operacjami wstawiania, aktualizacji i usuwania lub gdy potrzebujesz historii wymiaru wolnozmiennego typu 2 (SCD). AUTO CDC wymaga również prawdziwego klucza głównego. Zobacz Interfejsy API AUTO CDC: upraszczają przechwytywanie zmian danych za pomocą potoków.
  • Użyj REPLACE WHERE, gdy źródłem jest migawka i chcesz ponownie obliczyć i nadpisać zakres w tabeli docelowej określony predykatem, na przykład z ostatnich 7 dni, w ramach operacji wsadowej. Nie wymaga klucza głównego. Zobacz Przetwarzanie wsadowe z użyciem przepływów REPLACEWHERE.

Utwórz przepływ REPLACE USING

Zdefiniuj przepływy REPLACE USING w SQL lub Pythonie.

SQL

Użyj klauzuli FLOW REPLACE USING w tekście wraz z 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);

Alternatywnie należy użyć składni długiej: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);

Note

BY NAME jest wymagany w języku SQL. Dopasowuje kolumny według nazw, a nie pozycji.

Python

Zadeklaruj tabelę i przepływ razem z @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")

Alternatywnie wybierz istniejącą tabelę strumieniową za pomocą @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 to lista kluczowych kolumn. sequence_by jest nazwą kolumny lub wyrażeniem Column , i jest wymagana za każdym razem, gdy replace_using jest ustawiona.

Ustalanie kolejności i dane poza kolejnością

Kolumna SEQUENCE BY sprawia, że wynik jest niezależny od kolejności, w jakiej napływają aktualizacje. Wiersz jest przypisywany do klucza tylko wtedy, gdy jego sekwencja jest większa niż sekwencja już zapisana dla tego klucza, więc późny lub powtarzany wiersz starszy od aktualnej wartości jest ignorowany. Klucze nieobecne w aktualizacji pozostają nietknięte.

Postępuj zgodnie z poniższymi zasadami, aby zastępowanie działało w przewidywalny sposób:

Practice Powód
Użyj sekwencji, która rośnie ściśle dla każdej wersji klucza, na przykład znacznika czasu, numeru wersji lub przesunięcia w logu. Dwa wiersze z tym samym kluczem i tą samą sekwencją są zachowywane, co skutkuje podwójnymi wierszami dla tego klucza.
Użyj sekwencji niezerowej. Sekwencja zerowa może prowadzić do nieokreślonego zachowania.

Expectations

ZASTĄPNIJ UŻYWANIE przepływów wspierających oczekiwania. warn i fail zachowuje się tak jak w innych przepływach: nadal narusza wiersze, warn rejestruje naruszenie i fail zatrzymuje aktualizację. Zobacz Zarządzanie jakością danych przy użyciu oczekiwań dotyczących przepływu danych.

Asercja drop traktuje naruszający wiersz tak, jakby źródło nigdy go nie wygenerowało. Usunięty wiersz nie zastępuje, nie usuwa ani nie modyfikuje pasujących kluczy w docelowej tabeli:

  • Dropping następuje przed deduplikacją, więc flow zachowuje najnowszą poprawną wersję klucza.
  • Jeśli każdy wiersz wejściowy dla klucza zostanie usunięty, istniejące wiersze klucza pozostają nietknięte.
  • Ponieważ porzucony wiersz nie wyznacza poziomu sekwencji, późniejsza poprawna aktualizacja nadal się pojawia, nawet jeśli jej sekwencja jest niższa niż porzucona linia.

Ograniczenia

Przepływy „Zastąp przy użyciu” mają następujące ograniczenia:

  • REPLACE USING obsługuje pojedynczy przepływ dla każdej tabeli docelowej. Nie jest obsługiwane łączenie REPLACE USING z innym typem przepływu na tym samym celu.
  • Tabela docelowa musi zostać utworzona w potoku.
  • Źródło musi być źródłem przesyłania strumieniowego.
  • Musisz określić co najmniej jedną kolumnę klucza i kolumnę SEQUENCE BY . Kolumny klucza nie mogą być powtarzane, a typ każdej kolumny klucza musi być sortowalny. Typy atomowe, takie jak liczby całkowite, ciągi i daty, mogą być kluczami, podczas gdy MAP i VARIANT nie.
  • W przypadku autonomicznych tabel strumieniowych zobacz Zastosowanie częściowego zastępowania migawek za pomocą przepływów REPLACE USING, aby zapoznać się z różnicami w składni.

Examples

Poniższe przykłady odczytują dane z samples.wanderbricks.booking_updates, przykładowej tabeli zmian stanu rezerwacji, która jest dostępna w każdej przestrzeni roboczej z obsługą Unity Catalog. Każda rezerwacja pojawia się ponownie przy każdej zmianie, więc booking_id pojawia się ponownie z nowym booking_update_id. Zobacz zbiór danych Wanderbricks.

Przykład 1: Zachowaj najnowszy rekord dla każdego klucza

Zachowaj tylko aktualny stan każdej rezerwacji. Przepływ używa booking_id jako klucza i numeruje sekwencję według booking_update_id, więc najnowsza aktualizacja danej rezerwacji zastępuje wcześniejsze. Zamiast tego używaj AUTO CDC, gdy źródłem jest feed zmian z wyraźnymi operacjami wstawiania, aktualizacji i usuwania.

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")

Ten przykład jest porządkowany według booking_update_id, a nie według znacznika czasu updated_at, ponieważ kilka aktualizacji tej samej rezerwacji może mieć ten sam znacznik czasu. Wiersze zrównane w sekwencji są dodawane, a nie zastępowane, co pozostawiałoby więcej niż jeden wiersz dla tych rezerwacji.

Przykład 2: Klucz złożony z więcej niż jednej kolumny

Gdy rekord zostanie zidentyfikowany przez kombinację kolumn, wypisz je wszystkie w .REPLACE USING Tutaj każda rezerwacja jest identyfikowana przez (property_id, booking_id), więc przepływ zachowuje aktualny stan każdej rezerwacji na obiekt. Jeśli kolumna kluczowa może mieć wartość null, REPLACE USING traktuje wartość null jako odpowiadającą wartości null, zamiast pomijać wiersz.

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")

Przykład 3: Usuń nieprawidłowe rekordy z oczekiwaniem

Dodaj oczekiwanie, aby nieprawidłowe wiersze nie trafiały do miejsca docelowego. Porzucony wiersz traktowany jest tak, jakby źródło go nigdy nie wygenerowało: nie zastępuje ani nie usuwa dopasowanego klucza, a przepływ wraca do najnowszego poprawnego wiersza dla tego klucza. Ten przepływ odrzuca aktualizacje, które nie mają dodatniego 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")