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.
W inżynierii danych backfilling odnosi się do procesu wstecznego przetwarzania danych historycznych za pośrednictwem potoku danych, który został zaprojektowany do przetwarzania danych bieżących lub strumieniowych.
Zazwyczaj jest to oddzielny przepływ wysyłający dane do istniejących tabel. Na poniższej ilustracji przedstawiono przepływ wypełniania wysyłający dane historyczne do brązowych tabel w potoku.
Niektóre scenariusze, które mogą wymagać wypełnienia danych:
- Przetwarzanie danych historycznych ze starszego systemu w celu wytrenowania modelu uczenia maszynowego (ML) lub utworzenia panelu analizy trendów historycznych.
- Ponowne przetwarzanie podzbioru danych z powodu problemu z jakością danych z nadrzędnymi źródłami danych.
- Wymagania biznesowe uległy zmianie i trzeba uzupełnić dane wstecz dla innego okresu czasowego, który nie był objęty początkowym kanałem przetwarzania danych.
- Logika biznesowa została zmieniona i trzeba ponownie przetworzyć zarówno dane historyczne, jak i bieżące.
Przepływ uzupełniania, którego używasz, zależy od tabeli docelowej i danych źródłowych. Dla powoli zmieniającego AUTO CDC się wymiaru (SCD) typu 1 obiektu docelowego z autorytatywną migawką użyj jednorazowego AUTO CDC FROM SNAPSHOT przepływu. Dla migracji SCD, która odtwarza historyczne zmiany, użyj jednorazowego AUTO CDC przepływu.
W przypadku uzupełniania tabeli tylko przez dopisywanie: Użyj specjalistycznego przepływu dopisywania za pomocą opcji ONCE, aby uzupełnić tabelę strumieniową tylko przez dopisywanie. Aby uzyskać więcej informacji na temat opcji, zobacz append_flow lub ONCE).
Zagadnienia dotyczące uzupełniania danych historycznych do tabeli strumieniowej
- Zazwyczaj dołącz dane do brązowej tabeli przesyłania strumieniowego. Warstwy srebra i złota pobierają nowe dane z warstwy brązu.
- Upewnij się, że twoje przetwarzanie danych może zarządzać zduplikowanymi danymi bez błędów w przypadku wielokrotnego dodawania tych samych danych.
- Upewnij się, że schemat danych historycznych jest zgodny z bieżącym schematem danych.
- Weź pod uwagę rozmiar wolumenu danych oraz wymaganą umowę SLA dotyczącą przetwarzania, a następnie odpowiednio skonfiguruj klaster i rozmiary wsadów.
Przykład: Dodaj backfill do istniejącego pipeline
Zakładając, że masz przepływ danych, który pozyskuje nieprzetworzone dane rejestracji wydarzeń ze źródła w chmurze, począwszy od 01 stycznia 2025 r. Później zdajesz sobie sprawę, że chcesz uzupełnić wstecznie poprzednie trzy lata danych historycznych dla docelowych przypadków użycia raportowania i analizy. Wszystkie dane są w jednej lokalizacji, partycjonowane według roku, miesiąca i dnia w formacie JSON.
Potok początkowy
** Oto kod początkowego potoku, który przyrostowo wczytuje nieprzetworzone dane rejestracji zdarzeń z magazynu w chmurze.
Python
from pyspark import pipelines as dp
source_root_path = spark.conf.get("registration_events_source_root_path")
begin_year = spark.conf.get("begin_year")
incremental_load_path = f"{source_root_path}/*/*/*"
# create a streaming table and the default flow to ingest streaming events
@dp.table(name="registration_events_raw", comment="Raw registration events")
def ingest():
return (
spark
.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.maxFilesPerTrigger", 100)
.option("cloudFiles.schemaEvolutionMode", "addNewColumns")
.option("modifiedAfter", "2025-01-01T00:00:00.000+00:00")
.load(incremental_load_path)
.where(f"year(timestamp) >= {begin_year}") # safeguard to not process data before begin_year
)
SQL
-- create a streaming table and the default flow to ingest streaming events
CREATE OR REFRESH STREAMING LIVE TABLE registration_events_raw AS
SELECT * FROM read_files(
"/Volumes/gc/demo/apps_raw/event_registration/*/*/*",
format => "json",
inferColumnTypes => true,
maxFilesPerTrigger => 100,
schemaEvolutionMode => "addNewColumns",
modifiedAfter => "2024-12-31T23:59:59.999+00:00"
)
WHERE year(timestamp) >= '2025'; -- safeguard to not process data before begin_year
W tej sekcji używamy opcji modifiedAfter Auto Loader, aby upewnić się, że nie przetwarzamy wszystkich danych ze ścieżki przechowywania w chmurze. Przetwarzanie przyrostowe jest odcięte na tej granicy.
Wskazówka
Inne źródła danych, takie jak Kafka, Kinesis i Azure Event Hubs, mają równoważne opcje czytnika, aby osiągnąć to samo zachowanie.
Wypełnianie danych z poprzednich 3 lat
Teraz chcesz dodać jeden lub więcej przepływów, aby uzupełnić poprzednie dane. W tym przykładzie wykonaj następujące czynności:
- Użyj
append onceprzepływu. Wykonuje to jednorazowe wypełnianie bez kontynuowania działania po tym pierwszym wypełnianiu. Kod pozostaje w pipeline, a jeśli pipeline zostanie kiedykolwiek w pełni odświeżony, uzupełnianie braków jest uruchamiane ponownie. - Utwórz trzy przepływy uzupełniania, każdy na jeden rok (w tym przypadku dane są dzielone w ścieżce według roku). W przypadku języka Python parametryzujemy tworzenie przepływów, ale w języku SQL powtarzamy kod trzy razy, raz dla każdego przepływu.
Jeśli pracujesz nad własnym projektem i nie korzystasz z obliczeń serverless, możesz zaktualizować maksymalną liczbę pracowników dla potoku. Zwiększenie maksymalnej liczby pracowników gwarantuje, że masz zasoby do przetwarzania danych historycznych przy jednoczesnym przetwarzaniu bieżących danych przesyłanych strumieniowo zgodnie z oczekiwanym poziomem usługi (SLA).
Wskazówka
Jeśli używasz bezserwerowych obliczeń z rozszerzonym skalowaniem automatycznym (ustawieniem domyślnym), klaster automatycznie zwiększa rozmiar po wzroście obciążenia.
Python
from pyspark import pipelines as dp
source_root_path = spark.conf.get("registration_events_source_root_path")
begin_year = spark.conf.get("begin_year")
backfill_years = spark.conf.get("backfill_years") # e.g. "2024,2023,2022"
incremental_load_path = f"{source_root_path}/*/*/*"
# meta programming to create append once flow for a given year (called later)
def setup_backfill_flow(year):
backfill_path = f"{source_root_path}/year={year}/*/*"
@dp.append_flow(
target="registration_events_raw",
once=True,
name=f"flow_registration_events_raw_backfill_{year}",
comment=f"Backfill {year} Raw registration events")
def backfill():
return (
spark
.read
.format("json")
.option("inferSchema", "true")
.load(backfill_path)
)
# create the streaming table
dp.create_streaming_table(name="registration_events_raw", comment="Raw registration events")
# append the original incremental, streaming flow
@dp.append_flow(
target="registration_events_raw",
name="flow_registration_events_raw_incremental",
comment="Raw registration events")
def ingest():
return (
spark
.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.maxFilesPerTrigger", 100)
.option("cloudFiles.schemaEvolutionMode", "addNewColumns")
.option("modifiedAfter", "2024-12-31T23:59:59.999+00:00")
.load(incremental_load_path)
.where(f"year(timestamp) >= {begin_year}")
)
# parallelize one time multi years backfill for faster processing
# split backfill_years into array
for year in backfill_years.split(","):
setup_backfill_flow(year) # call the previously defined append_flow for each year
SQL
-- create the streaming table
CREATE OR REFRESH STREAMING TABLE registration_events_raw;
-- append the original incremental, streaming flow
CREATE FLOW
registration_events_raw_incremental
AS INSERT INTO
registration_events_raw BY NAME
SELECT * FROM STREAM read_files(
"/Volumes/gc/demo/apps_raw/event_registration/*/*/*",
format => "json",
inferColumnTypes => true,
maxFilesPerTrigger => 100,
schemaEvolutionMode => "addNewColumns",
modifiedAfter => "2024-12-31T23:59:59.999+00:00"
)
WHERE year(timestamp) >= '2025';
-- one time backfill 2024
CREATE FLOW
registration_events_raw_backfill_2024
AS INSERT INTO ONCE
registration_events_raw BY NAME
SELECT * FROM read_files(
"/Volumes/gc/demo/apps_raw/event_registration/year=2024/*/*",
format => "json",
inferColumnTypes => true
);
-- one time backfill 2023
CREATE FLOW
registration_events_raw_backfill_2023
AS INSERT INTO ONCE
registration_events_raw BY NAME
SELECT * FROM read_files(
"/Volumes/gc/demo/apps_raw/event_registration/year=2023/*/*",
format => "json",
inferColumnTypes => true
);
-- one time backfill 2022
CREATE FLOW
registration_events_raw_backfill_2022
AS INSERT INTO ONCE
registration_events_raw BY NAME
SELECT * FROM read_files(
"/Volumes/gc/demo/apps_raw/event_registration/year=2022/*/*",
format => "json",
inferColumnTypes => true
);
Ta implementacja wyróżnia kilka ważnych wzorców.
Separacja obaw
- Przetwarzanie przyrostowe jest niezależne od operacji uzupełniania braków.
- Każdy przepływ ma własne ustawienia konfiguracji i optymalizacji.
- Istnieje wyraźne rozróżnienie między operacjami przyrostowymi a uzupełniania.
Kontrolowane wykonywanie
- Użycie opcji
ONCEgwarantuje, że każde wypełnienie jest uruchamiane dokładnie raz. - Przepływ wypełniania pozostaje na grafie potoku, ale staje się bezczynny po zakończeniu. Jest on gotowy do użycia podczas pełnego odświeżania, automatycznie.
- Istnieje wyraźny ślad audytu operacji uzupełniania w definicji potoku.
Optymalizacja przetwarzania
- Duże wypełnienie można podzielić na wiele mniejszych wypełnień w celu szybszego przetwarzania oraz kontroli nad przetwarzaniem.
- Użycie rozszerzonego skalowania automatycznego dynamicznie skaluje rozmiar klastra na podstawie bieżącego obciążenia klastra.
Ewolucja schematu
- Korzystanie z
schemaEvolutionMode="addNewColumns"obsługuje zmiany schematu w sposób bezproblemowy. - Masz jednolite wnioskowanie schematu dla danych historycznych i bieżących.
- Istnieje bezpieczna obsługa nowych kolumn w nowszych danych.
Dodaj uzupełnienie do tabeli AUTO CDC SCD Typ 1
Użyj jednorazowego AUTO CDC FROM SNAPSHOT przepływu, aby dodać autorytatywną migawkę do celu SCD typu 1, który również otrzymuje kanał z bieżącym rejestrem zmian (CDC). Wersja migawkowa oraz kolumna sekwencjonowania CDC tworzą jedną domenę porządkową. Nowsze zdarzenie CDC ma pierwszeństwo przed starszym snapshotem, podczas gdy nowsze snapshot ma pierwszeństwo przed starszym zdarzeniem CDC.
Requirements
Przed dodaniem zasypu upewnij się, że przepusty spełniają następujące wymagania:
- Obiekt docelowy wykorzystuje SCD Type 1.
- Obiekt docelowy ma dokładnie jeden
AUTO CDC FROM SNAPSHOTprzepływ oraz co najmniej jeden przepływ o unikatowej nazwieAUTO CDC. - Wszystkie przepływy używają tej samej liczby kluczy w tej samej kolejności. Nazwy kluczy snapshot-flow są porównywane bez rozróżniania wielkości liter z nazwami kluczy
AUTO CDC. Wiele przepływówAUTO CDCmusi używać identycznych nazw kluczy i wielkości liter. - Wersja migawkowa oraz każda kolumna sekwencjonowania CDC mają dokładnie ten sam typ danych.
- Przepływ
AUTO CDC FROM SNAPSHOTnie definiuje oczekiwań. - Przepływy
AUTO CDCnie korzystają zIGNORE NULL UPDATES. W Pythonie nie ustawiajignore_null_updates,ignore_null_updates_column_list, aniignore_null_updates_except_column_list. - Potok przetwarzania korzysta z trybu wyzwalanego. Ten wzorzec nie obsługuje potoków ciągłych.
Oba typy przepływów mogą korzystać z interfejsu potoku SQL lub Python. Możesz mieszać przepływy SQL i Python w tym samym obiekcie docelowym.
Migawka musi reprezentować pełny stan źródła w jego wersji. Jeśli klucz docelowy nie występuje w snapshotie, AUTO CDC FROM SNAPSHOT traktuje brak jako usunięcie w wersji snapshot. Wydarzenie CDC z nowszą wersją zachowuje lub przywraca klucz.
Dodaj uzupełnienie
Aby dodać jednorazowe uzupełnienie snapshot i kontynuować przetwarzanie zdarzeń CDC, wykonaj następujące kroki:
- Zachowaj istniejącą tabelę docelową i jej
AUTO CDCbieżące przepływy w definicji pipeline. - Zdefiniuj migawkę referencyjną i jej wersję. W przypadku callbacku w Python pierwsze wywołanie musi zwrócić migawkę i wersję. Wróć
Nonedopiero po przetworzeniu co najmniej jednej migawki. - Dodaj jeden
AUTO CDC FROM SNAPSHOTprzepływ w językuonce=TruePython lubONCESQL. Aby uzupełnić istniejące miejsce docelowe za pomocą SQL, dołącz zapytanieWITH VERSION. Przepływ snapshot SQL bezWITH VERSIONobsługuje jedynie początkowe ładowanie do pustego celu. - Uruchom wyzwalaną aktualizację pipeline, aby przetworzyć backfill i trwające zdarzenia CDC.
Poniższy przykład zaczyna się od istniejącego potoku Python, który stopniowo przetwarza zmiany z customers_cdc. Załóżmy, że już uruchomiłeś ten pipeline i customers wypełniłeś element docelowy:
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def customers_cdc():
return (
spark.readStream.table("main.bronze.customers_cdc")
.withColumn("change_timestamp", col("change_timestamp").cast("timestamp"))
)
dp.create_streaming_table("customers")
dp.create_auto_cdc_flow(
name="customers_incremental_cdc",
target="customers",
source="customers_cdc",
keys=["customer_id"],
sequence_by=col("change_timestamp"),
apply_as_deletes=expr("operation = 'DELETE'"),
except_column_list=["operation", "change_timestamp"],
stored_as_scd_type=1,
)
Aby uzupełnić ten istniejący obiekt docelowy o stan elementu customers_snapshot z dnia 1 stycznia 2025 r., dodaj następujący kod do tej samej definicji potoku. Zachowaj istniejącą tabelę docelową i AUTO CDC przepływ:
from datetime import datetime, timezone
from typing import Optional, Tuple
from pyspark.sql import DataFrame
backfill_version = datetime(2025, 1, 1, tzinfo=timezone.utc)
def backfill_snapshot_and_version(
latest_snapshot_version: Optional[datetime],
) -> Optional[Tuple[DataFrame, datetime]]:
if latest_snapshot_version is None:
return (spark.read.table("main.legacy.customers_snapshot"), backfill_version)
return None
dp.create_auto_cdc_from_snapshot_flow(
target="customers",
source=backfill_snapshot_and_version,
keys=["customer_id"],
stored_as_scd_type=1,
once=True,
)
Callback musi zwrócić migawkę i wersję przy pierwszym wywołaniu. Jeśli zwróci None przed przetworzeniem jakiejkolwiek migawki, aktualizacja pipeline kończy się niepowodzeniem. Po przetworzeniu migawki zwrócenie None sygnalizuje, że nie są dostępne dodatkowe migawki.
Wersja snapshotowa to typ Python datetime, co odpowiada typowi TIMESTAMP Spark SQL. Istniejący AUTO CDC przepływ rzutuje change_timestamp na TIMESTAMP, tak aby oba typy sekwencjonowania były dokładnie dopasowane. Przykład wykorzystuje język Python dla obu przepływów, ale możesz zdefiniować dowolny z tych przepływów w SQL i mieszać przepływy SQL i Python w tym samym obiekcie docelowym. Aby uzyskać składnię SQL, w tym wymagane WITH VERSION zapytanie dla niepustego celu, zobacz CREATE FLOW (potoki).
Po pomyślnym zatwierdzeniu przepływu migawki, kolejne aktualizacje przyrostowe pomijają go, podczas gdy AUTO CDC przepływ kontynuuje przetwarzanie nowych zdarzeń.
Important
Pełne odświeżenie obiektu docelowego powtarza jednorazowy proces snapshot. Trzymaj migawkę dostępną i upewnij się, że nadal odzwierciedla zamierzony stan, zanim wykonasz pełne odświeżenie.
Ten jednolity wzorzec backfill nie obsługuje obiektów docelowych typu SCD 2 ani dwuczasowych.
Przykład: Uzupełnienie celu SCD podczas migracji
Typowym scenariuszem migracji jest tabela powoli zmieniających się wymiarów (SCD), która już istnieje w systemie dziedziczonym z wieloletnią historią, ale której oryginalny feed zmian nie jest już dostępny. Ponieważ oryginalne zdarzenia zmian zniknęły, zamiast tego jednorazowo odtwarzasz historię tabeli legacy do nowego obiektu docelowego AUTO CDC, a następnie dołączasz świeży feed CDC na przyszłość. Więcej informacji o typach SCD i AUTO CDC znajdziesz w The AUTO CDC APIs: Simplify change data capture with pipelines.
Ten wzorzec to jednorazowy AUTO CDC przepływ do tej samej tabeli strumieniowej, do której jest kierowany bieżący AUTO CDC przepływ. Cel typu AUTO CDC akceptuje tylko przepływy typu AUTO CDC, więc ziarno również musi być przepływem typu AUTO CDC. Prosty INSERT INTO ONCE przepływ dodawania do tej samej tabeli nie przechodzi walidacji:
-
Utwórz docelową tabelę strumieniową, do której zapisuje twój
AUTO CDCprzepływ. -
Jednorazowo zainicjuj historię legacy za pomocą
AUTO CDC ONCEflow, który odczytuje tabelę SCD jako strumień, sekwencjonowany według kolumny początku obowiązywania w systemie legacy. Powtarzaj rzędy dziedzictwa jako wydarzenia zmiany, zamiast je samodzielnie kształtować.AUTO CDCtworzy kolumny historii__START_ATi__END_ATdla obiektu docelowego SCD typu 2, więc nie zapisuj bezpośrednio w tych kolumnach. -
Podłącz bieżący
AUTO CDCprzepływ odczytujący najnowszy strumień zmian.AUTO CDCrozstrzyga kolejność dla każdego klucza, więc cutover musi być obowiązywał dla każdego klucza biznesowego osobno: pierwsza aktywna zmiana każdego klucza musi następować po ostatniej zasiedowanej zmianie dla tego samego klucza. Wartość sekwencji, która jest jedynie późniejsza od globalnego starszego maksimum, może nadal być przestarzała dla konkretnego klucza, a wtedy pierwsza rzeczywista zmiana tego klucza zostaje zignorowana lub umieszczona w niewłaściwej kolejności.
Poniższy kod tworzy tabelę streamingową, która wykorzystuje powyższe kroki:
CREATE OR REFRESH STREAMING TABLE customers_history;
-- One-time seed: replay the legacy history as change events
CREATE FLOW customers_history_seed
AS AUTO CDC ONCE INTO customers_history
FROM stream(legacy.customers_scd2)
KEYS (customer_id)
SEQUENCE BY valid_from
STORED AS SCD TYPE 2;
-- Ongoing live CDC into the same target
CREATE FLOW customers_history_cdc
AS AUTO CDC INTO customers_history
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY change_timestamp
STORED AS SCD TYPE 2;
Oba przepływy muszą być zgodne pod względem kluczy, typu SCD oraz typu danych kolumny sekwencjonującej. W poprzednim przykładzie oba przepływy są porządkowane według znacznika czasu, przy czym do oddzielenia wstępnie załadowanej historii od strumienia danych na żywo używany jest jeden moment przełączenia. Jeśli starsza tabela używa do sekwencjonowania wartości innego typu niż strumień danych na żywo, rzutuj jedną z nich tak, aby typy były zgodne.
Ten sam schemat działa w przypadku celu SCD Type 1: zmień STORED AS SCD TYPE 2 na STORED AS SCD TYPE 1 w obu przepływach, a tabela docelowa zachowuje tylko bieżący wiersz dla każdego klucza. Zanim zdecydujesz się na którykolwiek z kształtów, zweryfikowaj na próbce kluczy, że pierwsza aktywna zmiana dla zasianego klucza generuje dokładnie jedną nową wersję i poprawnie zamyka poprzednią. Na tym etapie zwykle pojawia się przerwa w sekwencji dla poszczególnych klawiszy.