Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Ważna
Ta funkcja jest dostępna w wersji beta.
Przepływy REPLACE USING w potokach Lakeflow utrzymują tabelę w aktualnym stanie na podstawie strumienia częściowych obrazów migawkowych. Zastępują wszystkie wiersze odpowiadające określonym kolumnom klucza i pozostawiają resztę tabeli bez zmian. Nadają się do źródeł, które okresowo wysyłają pełny zestaw wierszy przypisanych do danego klucza, na przykład plików wczytywanych ponownie, gdy ich zawartość ulega zmianie.
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ąpiony, ponieważ numer sekwencyjny 2 jest większy niż numer sekwencyjny 1 |
| 1 | biurko | click | 2 | Zastąpiony, ponieważ numer sekwencyjny 2 jest większy niż numer sekwencyjny 1 |
| 2 | iOS | click | 1 | Niezmienione, ponieważ klucz nie jest obecny w tej aktualizacji |
| 2 | biurko | click | 1 | Niezmienione, ponieważ klucz nie jest obecny 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. Kolumny kluczy nie mogą być powtarzane, a ich typy muszą być sortowalne. Typy atomowe, takie jak liczby całkowite, ciągi i daty, mogą być kluczami.MAPIVARIANTnie mogą być kluczami.
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.
Note
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.
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.
Obsługiwane operacje
W zapytaniu przepływu obsługiwane są następujące operacje:
- Rzutowanie kolumnowe: wybierz lub przesuń kolejność podzbioru kolumn.
- Filtry za pomocą
WHERE. - Wyrażenia skalarne, takie jak
CAST, arytmetyka orazCASE. - Deduplikacja za pomocą
SELECT DISTINCT. - Granice wierszy z
LIMIT. - Funkcje generatora, takie jak
EXPLODEiPOSEXPLODE. - Sumy dwóch źródeł strumieniowych za pomocą
UNION ALL. - Połączenia strumienia z danymi statycznymi: wewnętrzne i lewostronne zewnętrzne, ze strumieniem po lewej stronie.
- Złączenia wewnętrzne strumień–strumień.
Następujące operacje wymagają dodatkowej konfiguracji:
- Okna czasowe, takie jak
window(ts, '5 minutes'), wymagają watermarku. - Złączenia zewnętrzne typu strumień-strumień wymagają znacznika wodnego oraz warunku dotyczącego zakresu czasu.
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.
Następujące operacje nie są obsługiwane w zapytaniu przepływu:
- Agregacje, takie jak
SUM,COUNT, iGROUP BY. - Funkcje okien inne niż okna oparte na czasie, takie jak
ROW_NUMBER() OVER (...), nawet przy użyciu znaku wodnego. - Sortowanie za pomocą
ORDER BY. - Operacje zbiorów, takie jak
INTERSECTiEXCEPT. - Związki zawodowe, które mieszają źródło streamingu z źródłem niestreamującym.
- Odczyt z źródła niebędącego strumieniem, takiego jak
spark.range(). - Pełne i prawe zewnętrzne złączenia strumienia z danymi statycznymi.
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")