Praca z historią tabel

W przypadku tabel Apache Iceberg i Delta Lake każda operacja modyfikująca tabelę tworzy nową wersję tabeli. Informacje historyczne służą do audytowania operacji, przywracania tabeli do wcześniejszego stanu lub wykonywania zapytań względem tabeli w określonym momencie za pomocą funkcji podróży w czasie.

Note

Nie używaj historii tabeli jako długoterminowego rozwiązania do archiwizacji danych. Użyj tylko ostatnich 7 dni dla operacji podróży czasowych, chyba że ustawiono konfiguracje przechowywania danych i dzienników na większą wartość.

Pobieranie historii tabel

Uruchom polecenie , DESCRIBE HISTORY aby pobrać informacje, w tym informacje o operacjach, użytkowniku i znaczniku czasu dla każdego zapisu w tabeli. Operacje są zwracane w odwrotnej kolejności chronologicznej.

Informacje o kolumnach zwracanych przez DESCRIBE HISTORY, wartościach w kolumnie operationParameters oraz metrykach dla poszczególnych operacji w kolumnie operationMetrics znajdziesz w artykule Schemat historii tabeli i metryki operacji.

Przechowywanie historii tabeli jest określane przez ustawienie tabeli logRetentionDuration, które domyślnie wynosi 30 dni.

Note

Przechodzenie w czasie i historia tabeli są kontrolowane przez różne progi przechowywania. Zobacz Podróż czasowa.

DESCRIBE HISTORY table_name       -- get the full history of the table

DESCRIBE HISTORY table_name LIMIT 1  -- get the last operation only

Aby uzyskać szczegółowe informacje o składni spark SQL, zobacz DESCRIBE HISTORY.

Szczegółowe informacje o składni języka Scala, Java i Python można znaleźć w dokumentacji interfejsu API usługi Delta Lake.

Eksplorator wykazu wyświetla wizualnie historię tabel na karcie Historia .

Identyfikowanie typu OPTIMIZE operacji

Automatyczne kompaktowanie, klastrowanie cieczy i porządkowanie Z są wyświetlane w historii tabeli jako OPTIMIZE operacje. Aby określić, który z nich został uruchomiony, sprawdź kolumnę operationParameters.

Aby sklasyfikować każdą OPTIMIZE operację w historii tabeli, uruchom następujące polecenie:

SELECT
  version,
  timestamp,
  CASE
    WHEN operationParameters.clusterBy IS NOT NULL AND operationParameters.clusterBy <> '[]' THEN 'Liquid clustering'
    WHEN operationParameters.zOrderBy IS NOT NULL AND operationParameters.zOrderBy <> '[]' THEN 'Z-ordering'
    WHEN operationParameters.auto = 'true' THEN 'Auto compaction'
    ELSE 'Manual OPTIMIZE'
  END AS optimize_type,
  operationParameters.auto AS is_auto_compaction,
  operationParameters.clusterBy AS cluster_by,
  operationParameters.zOrderBy AS z_order_by,
  operationMetrics.numRemovedFiles AS files_compacted,
  operationMetrics.numAddedFiles AS files_added,
  operationMetrics.numRemovedBytes AS bytes_removed,
  operationMetrics.numAddedBytes AS bytes_added
FROM (DESCRIBE HISTORY table_name)
WHERE operation = 'OPTIMIZE'
ORDER BY version DESC;

W poniższych sekcjach opisano szczegółowo każdą operationParameters wartość. Definicje kluczy operationMetrics wybieranych przez poprzednie zapytanie można znaleźć w artykule Metryki operacji.

Automatyczne kompaktowanie

Automatyczna kompakcja ustawia parametr auto na true. Azure Databricks wyzwala automatyczne kompaktowanie automatycznie po zapisie. Gdy auto ma wartość false, użytkownik lub zaplanowane zadanie uruchomił polecenie OPTIMIZE.

Na przykład operacja automatycznego kompaktowania pokazuje następujące elementy:

operationParameters: {
  "auto": "true"
}

Aby uzyskać więcej informacji na temat automatycznego kompaktowania, zobacz Automatyczne kompaktowanie.

Klastrowanie cieczy

Mechanizm Liquid Clustering wypełnia parametr clusterBy nazwami kolumn klastrowania. Pusta clusterBy tablica ([]) wskazuje tylko kompaktowanie plików.

Na przykład operacja, która grupowała dane według kolumn date i region, pokazuje następujące informacje:

operationParameters: {
  "clusterBy": "[\"date\",\"region\"]"
}

Aby uzyskać więcej informacji na temat klastrowania płynnego, zobacz Używanie klastrowania płynnego dla tabel.

Kolejność Z

Porządkowanie Z wypełnia parametr zOrderBy nazwami kolumn porządku Z. Pusta zOrderBy tablica ([]) wskazuje, że operacja nie zastosowała porządkowania Z.

Na przykład operacja, która zastosowała kolejność Z w date kolumnie, pokazuje następujące elementy:

operationParameters: {
  "zOrderBy": "[\"date\"]"
}

Zakres operacji

Parametr predicate wskazuje, czy operacja została uruchomiona w pełnej tabeli, czy tylko jej części:

  • Pusta predicate tablica ([]) oznacza, że operacja została uruchomiona w całej tabeli.
  • Wypełniona predicate tablica oznacza, że docelowe OPTIMIZE table_name WHERE <partition_predicate> polecenie było uruchamiane tylko na partycjach, które są zgodne z predykatem.

Na przykład operacja ukierunkowana na partycje pasujące do year = 2024 wyświetla następujące informacje:

operationParameters: {
  "predicate": "[\"'year = 2024\"]"
}

Podróż czasowa

Podróż czasowa obsługuje wykonywanie zapytań dotyczących poprzednich wersji tabeli na podstawie sygnatury czasowej lub wersji tabeli (zarejestrowanej w dzienniku transakcji). Możesz użyć podróży w czasie dla aplikacji, takich jak:

  • Ponowne tworzenie analiz, raportów lub danych wyjściowych, takich jak dane wyjściowe modelu uczenia maszynowego. Może to być przydatne w przypadku debugowania lub inspekcji, zwłaszcza w branżach regulowanych.
  • Pisanie złożonych zapytań czasowych.
  • Naprawianie błędów w danych.
  • Zapewnienie izolacji migawek dla zestawu zapytań dotyczących szybko zmieniających się tabel.

Note

W środowisku Databricks Runtime 18.0 lub nowszym zapytania dotyczące podróży w czasie są blokowane, jeśli żądają wersji starszej deletedFileRetentionDuration niż właściwość tabeli (domyślnie 7 dni). W przypadku tabel zarządzanych przez katalog Unity, dotyczy to Databricks Runtime 12.2 lub nowszego.

Składnia podróży czasowej

Tworzysz zapytanie do tabeli z funkcją podróży w czasie, dodając klauzulę po specyfikacji nazwy tabeli.

  • timestamp_expression może być jednym z:
    • '2018-10-18T22:15:12.013Z', czyli ciąg, który można przekształcić na znacznik czasu
    • cast('2018-10-18 13:36:32 CEST' as timestamp)
    • '2018-10-18', czyli ciąg daty
    • current_timestamp() - interval 12 hours
    • date_sub(current_date(), 1)
    • Dowolne inne wyrażenie, które można rzutować na znacznik czasu
  • version to długa wartość, którą można uzyskać z danych wyjściowych elementu DESCRIBE HISTORY table_spec.

Ani timestamp_expression ani version nie mogą być podzapytaniami.

Akceptowane są tylko ciągi daty lub znacznika czasu. Przykład: "2019-01-01" i "2019-01-01T00:00:00.000Z". Zobacz następujący kod jako przykład składni.

SQL

SELECT * FROM people10m TIMESTAMP AS OF '2018-10-18T22:15:12.013Z';
SELECT * FROM people10m VERSION AS OF 123;

Python

df1 = spark.read.option("timestampAsOf", "2019-01-01").table("people10m")
df2 = spark.read.option("versionAsOf", 123).table("people10m")

Można również użyć składni @, aby określić znacznik czasu lub wersję jako część nazwy tabeli. Znacznik czasu musi być w yyyyMMddHHmmssSSS formacie. Możesz określić wersję za pomocą polecenia @v. Zobacz następujący kod jako przykład składni.

SQL

-- Timestamp version
SELECT * FROM people10m@20190101000000000
-- Version number
SELECT * FROM people10m@v123

Python

# Timestamp version
spark.read.table("people10m@20190101000000000")
# Version number
spark.read.table("people10m@v123")

Konfigurowanie przechowywania danych dla zapytań dotyczących podróży w czasie

Aby wykonać zapytanie dotyczące poprzedniej wersji tabeli, należy zachować zarówno dziennik, jak i pliki danych dla tej wersji:

  • Pliki danych są usuwane, gdy VACUUM zostanie uruchomiony na tabeli.
  • Pliki dziennika są usuwane automatycznie po wersjach tabeli punktów kontrolnych.

Aby zwiększyć próg przechowywania danych dla tabel, należy skonfigurować następujące właściwości tabeli, zastępując element <format> albo delta :iceberg

  • <format>.logRetentionDuration = "interval <interval>": określa, jak długo jest przechowywana historia tabeli. Wartość domyślna to interval 30 days.
    • W środowisku Databricks Runtime w wersji 18.0 lub nowszej logRetentionDuration musi być większe lub równe deletedFileRetentionDuration. W przypadku tabel zarządzanych przez katalog Unity, dotyczy to Databricks Runtime 12.2 lub nowszego.
  • <format>.deletedFileRetentionDuration = "interval <interval>": określa wartość progową VACUUM używaną do usuwania plików danych, do których nie odwołuje się już bieżąca wersja tabeli. Wartość domyślna to interval 7 days.

Na przykład, aby uzyskać dostęp do danych historycznych z ostatnich 30 dni, ustaw delta.deletedFileRetentionDuration = "interval 30 days", co odpowiada ustawieniu domyślnemu dla delta.logRetentionDuration.

Important

Zwiększenie progu przechowywania danych może spowodować zwiększenie kosztów magazynowania, ponieważ utrzymuje się więcej plików danych.

Właściwości tabeli można określić podczas tworzenia tabeli lub ustawić za pomocą instrukcji ALTER TABLE . Zobacz Informacje o właściwościach tabeli.

Przykłady podróży czasowych

Aby naprawić skutki przypadkowego usunięcia danych z tabeli użytkownika 111:

INSERT INTO my_table
  SELECT * FROM my_table TIMESTAMP AS OF date_sub(current_date(), 1)
  WHERE userId = 111

Aby naprawić przypadkowe nieprawidłowe aktualizacje tabeli:

MERGE INTO my_table target
  USING my_table TIMESTAMP AS OF date_sub(current_date(), 1) source
  ON source.userId = target.userId
  WHEN MATCHED THEN UPDATE SET *

Aby wysłać zapytanie dotyczące liczby nowych klientów dodanych w ciągu ostatniego tygodnia:

SELECT
(
  SELECT count(distinct userId)
  FROM my_table
)
-
(
  SELECT count(distinct userId)
  FROM my_table TIMESTAMP AS OF date_sub(current_date(), 7)
) AS new_customers

Punkty kontrolne dziennika transakcji

Dziennik transakcji rejestruje wersje tabeli jako pliki JSON w katalogu dziennika transakcji wraz z danymi tabeli.

Aby zoptymalizować wykonywanie zapytań dotyczących punktów kontrolnych, wersje tabel są agregowane do plików punktów kontrolnych Parquet, co zwiększa wydajność, uniemożliwiając odczytywanie wszystkich wersji historii tabel w formacie JSON. Użytkownicy nie muszą bezpośrednio korzystać z punktów kontrolnych.

Usługa Azure Databricks optymalizuje częstotliwość tworzenia punktów kontrolnych pod kątem rozmiaru i obciążenia danych. Częstotliwość punktów kontrolnych może ulec zmianie bez powiadomienia.

Przywracanie tabeli do wcześniejszego stanu

RESTORE Użyj polecenia , aby przywrócić tabelę do poprzedniej wersji lub znacznika czasu, w tym w następujących scenariuszach:

  • Możesz przywrócić już przywróconą tabelę.
  • Możesz przywrócić sklonowaną tabelę.

Weź pod uwagę następujące wymagania:

  • Aby przywrócić tabelę, musisz mieć MODIFY uprawnienia do tabeli.
  • Po usunięciu plików danych ręcznie lub za pomocą VACUUM nie można przywrócić tabeli do starszej wersji, która odwołuje się do tych plików. Przywracanie do tej wersji częściowo jest nadal możliwe, jeśli spark.sql.files.ignoreMissingFiles jest ustawione na true.
  • Aby przywrócić według znacznika czasu, użyj formatów yyyy-MM-dd HH:mm:ss lub yyyy-MM-dd.
RESTORE TABLE target_table TO VERSION AS OF <version>;
RESTORE TABLE target_table TO TIMESTAMP AS OF <timestamp>;

Aby uzyskać szczegółowe informacje o składni, zobacz RESTORE.

Zachowanie transmisji strumieniowej

Przywracanie to operacja zmiany danych i może spowodować zduplikowanie danych dla obciążeń podrzędnych. Wpisy dziennika dodane przez RESTORE polecenie zawierają wartość dataChange ustawioną na true.

W przypadku obciążeń podrzędnych, takich jak zadanie przesyłania strumieniowego ze strukturą , które przetwarza aktualizacje tabeli, wpisy dziennika zmian danych dodane przez operację przywracania są uznawane za nowe aktualizacje danych, a ich przetwarzanie może spowodować zduplikowanie danych.

Przykład:

Wersja tabelaryczna Operation Aktualizacje dzienników Rekordy w aktualizacjach dziennika zmian danych
0 INSERT AddFile(/path/to/file-1, dataChange = true) (imię = Viktor, wiek = 29), (imię = George, wiek = 55)
1 INSERT AddFile(/path/to/file-2, dataChange = true) (imię = George, wiek = 39)
2 OPTIMIZE AddFile(/path/to/file-3, dataChange = false), RemoveFile(/path/to/file-1), RemoveFile(/path/to/file-2) Brak rekordów. OPTIMIZE kompaktowanie nie zmienia danych w tabeli.
3 RESTORE(version=1) RemoveFile(/path/to/file-3), AddFile(/path/to/file-1, dataChange = true), AddFile(/path/to/file-2, dataChange = true) (imię = Viktor, wiek = 29), (imię = George, wiek = 55), (imię = George, wiek = 39)

W poprzednim przykładzie polecenie RESTORE powoduje wyświetlenie zmian, które były wcześniej widoczne podczas odczytu tabeli w wersjach 0 i 1. Jeśli zapytanie przesyłane strumieniowo ponownie odczytuje tę tabelę, te pliki są traktowane jako nowo dodane dane i są przetwarzane ponownie.

Przywracanie metryk

Po zakończeniu RESTORE raportuje następujące metryki jako ramkę danych z jednym wierszem:

  • table_size_after_restore: rozmiar tabeli po przywróceniu.

  • num_of_files_after_restore: liczba plików w tabeli po przywróceniu.

  • num_removed_files: liczba plików usuniętych (logicznie usuniętych) z tabeli.

  • num_restored_files: liczba przywróconych plików z powodu wycofywania.

  • removed_files_size: całkowity rozmiar w bajtach plików usuniętych z tabeli.

  • restored_files_size: całkowity rozmiar w bajtach przywróconych plików.

    Przykład przywracania metryk

Znajdowanie ostatniej wersji zatwierdzenia

Aby uzyskać numer wersji ostatniego zatwierdzenia zapisanego przez bieżący SparkSession we wszystkich wątkach i wszystkich tabelach, wykonaj zapytanie o konfigurację SQL spark.databricks.<format>.lastCommitVersionInSession. Zastąp ciąg <format> ciągiem delta lub iceberg, w zależności od formatu tabeli.

Przykład:

SQL

SET spark.databricks.delta.lastCommitVersionInSession

Python

spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")

Scala

spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")

Jeśli nie dokonano żadnych zatwierdzeń przez SparkSession, zapytanie dotyczące klucza zwraca pustą wartość.

Note

Jeśli współdzielisz ten sam SparkSession między wieloma wątkami, to tak, jakby współdzielić zmienną między wieloma wątkami. Mogą wystąpić warunki wyścigu dotyczące współbieżnych aktualizacji wartości konfiguracji.