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.
Próby i ponowne uruchomienia są nieuniknione w każdym prawdziwym pipeline, dlatego ta strona wyjaśnia gwarancje przetwarzania, które dają pipeline Lakeflow, oraz jak bezpiecznie uruchomić elementy, które piszesz, do ponownego uruchamiania.
Overview
Dwie powiązane właściwości decydują o tym, czy ponowne uruchomienie pociągu jest bezpieczne:
- Idempotencja oznacza, że potok generuje ten sam wynik niezależnie od tego, ile razy uruchomisz go dla tych samych danych wejściowych. Ponowne uruchomienie po awarii, dwukrotne uzupełnienie zakresu dat lub ręczne ponowne uruchomienie zadania nigdy nie tworzy powielonych wierszy ani nie uszkadza stanu.
- Gwarancja przetwarzania opisuje, ile razy każdy rekord wpływa na wynik. Przynajmniej raz przetwarzanie gwarantuje, że każdy rekord zostanie przetworzony, ale niepowodzenie i ponowna próba mogą przetworzyć niektóre rekordy więcej niż raz, co grozi powstaniem duplikatów. Przetwarzanie dokładnie raz gwarantuje, że każdy rekord wpływa na wynik tak, jak gdyby został przetworzony dokładnie jeden raz, nawet w przypadku ponownych prób, bez duplikatów i bez braków.
Potoki Lakeflow są domyślnie idempotentne w zakresie elementów, którymi zarządzają, i zapewniają przetwarzanie dokładnie jeden raz we własnych tabelach zarządzanych. Ważne jest, aby zrozumieć, gdzie te gwarancje przestają być automatyczne, aby móc dodać odpowiednie zabezpieczenia na obrzeżach pipeline'u.
Jak to działa
Potoki Lakeflow zapewniają przetwarzanie dokładnie raz i idempotentność w zarządzanych przez nie przepływach, a także udostępniają narzędzia, które pomagają zachować idempotentność tworzonej przez Ciebie logiki.
Dokładnie jednorazowe przetwarzanie dla tabel zarządzanych
W tabelach zarządzanych domyślnie otrzymujesz przetwarzanie dokładnie raz. Tabele strumieniowe wykorzystują punkty kontrolne Structured Streaming w połączeniu z transakcyjnymi zapisami Delta Lake: każda mikropartia zatwierdza jednocześnie swoje offsety źródłowe i dane wyjściowe, więc ponowiona mikropartia po awarii albo kończy się pełnym powodzeniem, albo jest w całości wycofywana i ponawiana — nigdy nie dochodzi do jej częściowego zastosowania dwa razy. Dotyczy to pozyskiwania plików za pomocą Auto Loader, odczytów z Apache Kafka, Kinesis i Azure Event Hubs oraz operacji AUTO CDC upsert, bez konieczności pisania kodu.
Jeśli źródło przynajmniej raz wysyła ten sam rekord wielokrotnie, pipeline przetwarza je jako unikalne rekordy i zapisuje je wszystkie do Twojej tabeli. Usunięcie tych duplikatów to twoja odpowiedzialność. Zobacz Deduplikuj przynajmniej raz źródła.
Idempotentność odczytów również wynika z tych samych punktów kontrolnych. Auto Loader i punkty kontrolne tabel strumieniowych gwarantują, że każdy plik źródłowy lub offset jest przetwarzany tylko raz na potrzeby śledzenia stanu, dzięki czemu po awarii ponowne uruchomienie potoku jest wznawiane od punktu kontrolnego, zamiast prowadzić do ponownego przetwarzania lub pomijania danych. Osiąga się to dzięki użyciu tabel strumieniowych opartych na spark.readStream zamiast ręcznie napisanych pętli wsadowych. Zobacz Tablice streamingowe.
Użyj AUTO CDC zamiast ręcznie pisanego MERGE
AUTO CDC INTO jest z natury idempotentny względem swojego keys i sequence_by. Dwukrotne zastosowanie tego samego rekordu zmian lub zastosowanie rekordów w niewłaściwej kolejności skutkuje tym samym stanem końcowym, ponieważ potok używa kolumny z numerem sekwencyjnym do określenia, czy przychodzący wiersz jest faktycznie nowszy niż ten już zapisany:
CREATE FLOW customers_cdc_flow AS AUTO CDC INTO customers_silver
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY sequence_num
STORED AS SCD TYPE 1;
Jeśli napiszesz własną logikę upsert poza AUTO CDC (rzadko, ale czasem jest to konieczne przy złożonych warunkach scalania), oprzyj ją na stabilnym kluczu biznesowym i zadbaj o to, aby można ją było bezpiecznie zastosować dwukrotnie, na przykład MERGE ... WHEN MATCHED oparty na order_id zamiast ślepego INSERT. Więcej informacji można znaleźć w interfejsach API AUTO CDC: upraszczanie przechwytywania zmian danych za pomocą potoków danych.
Zadbaj, aby własne transformacje były idempotentne
Aby zachować idempotentność logiki przy ponownym wykonywaniu operacji zapisu, stosuj się do tych dwóch zasad:
- Unikaj niedeterministycznych przekształceń w widokach materializowanych. Ponieważ widok materializowany może być przeliczany od podstaw lub przyrostowo, unikaj funkcji, których wynik zależy od momentu ich uruchomienia, a nie od tego, jakie są dane wejściowe. Na przykład nie używaj
current_timestamp()do obliczania wartości biznesowej, która powinna pozostać stała po zapisaniu; bierz znacznik czasu ze zdarzenia źródłowego lub przekaż go jako parametr, aby ponowne obliczenia dawały identyczny wynik. - Zaprojektuj pełne odświeżenia dla bezpieczeństwa. Pełne odświeżenie usuwa tabelę i oblicza ją ponownie od zera, co jest bezpieczne tylko wtedy, gdy każde źródło wejściowe nadal może wygenerować pełną historię. Jeśli źródło udostępnia tylko przesuwne okno zmian, pełne odświeżenie podrzędnej tabeli
AUTO CDCmoże po cichu spowodować utratę historii, dlatego projektując retencję danych źródła i tematu, trzeba mieć to na uwadze.
Złap dokładnie raz na krawędziach
Miejscem, w którym semantyka exactly-once przestaje być automatyczna, są granice obszarów bezpośrednio kontrolowanych przez potok przetwarzania, takie jak zapisy do systemów zewnętrznych. Gdy zapisujesz dane do systemu zewnętrznego, zadbaj o to, aby sam zapis był idempotentny, na przykład stosując operację upsert według klucza po stronie systemu docelowego, ponieważ ponowione przetwarzanie mikropartii mogłoby spowodować dwukrotny zapis tej samej partii. Poniższy sink zapisuje każdą partycję partii z wykonawców i używa klucza idempotentnego, aby powtórzona partia nie zapisywała się podwójnie:
from pyspark import pipelines as dp
@dp.foreach_batch_sink(name="orders_to_external_api")
def write_orders_to_api(batch_df, batch_id):
def write_partition(rows):
# Open one client per partition.
for row in rows:
# Use an idempotency key (order_id) so a retried batch doesn't double-write.
upsert_to_external_system(key=row.order_id, payload=row.asDict())
batch_df.select("order_id", "amount").foreachPartition(write_partition)
Więcej informacji na temat zapisywania do systemów zewnętrznych znajdziesz w Sinks in Lakeflow pipelines.
Deduplikuj źródła przynajmniej raz
Gdy źródło może dostarczyć rekord więcej niż raz, usuwaj duplikaty na dalszym etapie przetwarzania. Połącz znak wodny z dropDuplicatesWithinWatermark, który jest świadomy znaków wodnych i nie wymaga stanu nieograniczonego do wykrywania duplikatów. Deduplikuj kolumny, które jednoznacznie identyfikują zdarzenie. Tożsamość może obejmować kilka kolumn, gdy żadna pojedyncza kolumna nie jest unikalna sama w sobie. W poniższym przykładzie numer sekwencyjny kliknięcia jest unikalny tylko w swojej sesji, więc dwie kolumny razem identyfikują zdarzenie:
from pyspark import pipelines as dp
@dp.table(name="clicks_deduped")
def clicks_deduped():
return (
spark.readStream.table("clicks_bronze")
.withWatermark("click_ts", "5 minutes")
.dropDuplicatesWithinWatermark(["session_id", "click_seq_num"])
)
Wybierz te kolumny z kontraktu unikalności danych źródłowych, a nie na podstawie tego, co wydaje się unikalne w danych przykładowych. Kolumny, które mogą legalnie się powtarzać, odrzucają rzeczywiste wydarzenia, gdy traktujesz je jako tożsamość. Typowym przykładem jest sytuacja, gdy użytkownik kliknie tę samą reklamę dwa razy: deduplikacja według użytkownika i reklamy po cichu pomija drugie kliknięcie.
AUTO CDC semantyka operacji upsert oparta na kluczach naturalnie również scala duplikaty, więc kierowanie danych dostarczanych co najmniej raz przez przepływ AUTO CDC oparty na stabilnym kluczu biznesowym to kolejny sposób na osiągnięcie stanu exactly-once.
Limitations
Przetwarzanie dokładnie raz ma zastosowanie do zarządzanych przepływów Delta-to-Delta. Traktuj następujące krawędzie jako przynajmniej raz i dodaj tam jawną deduplikację lub logikę idempotentnego zapisu:
-
foreach_batch_sinkoraz niestandardowe zapisy zewnętrzne. Spark gwarantuje, że dla partii zostanie podjęta próba co najmniej raz, ale ponowna próba przetworzenia partii po częściowym zapisie może spowodować, że niektóre wiersze będą widoczne dwukrotnie w systemie zewnętrznym. Spraw, aby operacja zapisu do systemu zewnętrznego była idempotentna, na przykład przez zastosowanie operacji upsert względem klucza naturalnego albo zapisanie identyfikatora partii, na podstawie którego odbiorca może usuwać duplikaty. - Kafka jak zlew. Tematy Kafki nie obsługują transakcyjnych zapisów w modelu exactly-once w taki sposób, jak Delta, więc ponowiona mikropartia zapisująca do Kafki może powodować powstawanie zduplikowanych komunikatów. Jeśli odbiorcy po stronie downstream są wrażliwi na duplikaty, należy usuwać duplikaty po stronie odbiorcy, na przykład na podstawie identyfikatora zdarzenia.
- Niestandardowe źródła danych Python używane jako źródła. To, czy odczyty są przetwarzane dokładnie raz, zależy od tego, czy implementacja źródła poprawnie raportuje offsety i wznawia działanie od nich. Jeśli nie śledzi offsetów, traktuj to jako mechanizm co najmniej raz i deduplikuj na dalszym etapie przetwarzania za pomocą
dropDuplicatesna podstawie identyfikatora zdarzenia lub polegając na semantyce upsert opartej na kluczu wAUTO CDC.
Co do zasady, jeśli cały Twój potok danych jest typu Delta-to-Delta (tabele strumieniowe i widoki zmaterializowane odczytują i zapisują tabele Delta za pośrednictwem zarządzanych przepływów), to masz już semantykę exactly-once. W momencie, gdy dodasz foreach_batch_sink, nie-delta sink lub niezweryfikowane niestandardowe źródło, traktuj tę konkretną krawędź przynajmniej raz i dodaj tam logikę idempotentnego zapisu lub deduplikacji.