Używanie Auto Loader z Unity Catalog

Moduł Auto Loader może bezpiecznie pozyskiwać dane z lokalizacji zewnętrznych skonfigurowanych z Unity Catalogiem. Aby dowiedzieć się więcej na temat bezpiecznego łączenia magazynu z Unity Catalog, zobacz część Połącz się z magazynem obiektów w chmurze przy użyciu Unity Catalog. Automatyczne ładowanie opiera się na Strukturalnym Przesyłaniu Strumieniowym na potrzeby przetwarzania przyrostowego; aby uzyskać zalecenia i ograniczenia, zobacz Używanie Unity Catalog ze Strukturalnym Przesyłaniem Strumieniowym.

Uwaga

W środowisku Databricks Runtime 11.3 LTS lub nowszym można używać Auto Loader w trybach dostępu standardowego lub dedykowanego (dawniej tryb dostępu współdzielonego i tryb pojedynczego użytkownika).

Tryb wyświetlania listy katalogów jest domyślnie obsługiwany.

Określanie lokalizacji dla zasobów Auto Loader w Unity Catalog

Model zabezpieczeń Unity Catalog zakłada, że wszystkie lokalizacje magazynowania wymienione w obciążeniu będą zarządzane przez Unity Catalog. Databricks zaleca zawsze przechowywać informacje o punktach kontrolnych i ewolucji schematów w lokalizacjach przechowywania zarządzanych przez Unity Catalog. Unity Catalog nie umożliwia zagnieżdżania punktów kontrolnych ani wnioskowania schematu i plików ewolucji w katalogu tabeli.

Przetwarzanie danych z magazynu w chmurze przy użyciu Unity Catalog

W poniższych przykładach założono, że wykonujący użytkownik ma READ FILES uprawnienia do lokalizacji zewnętrznej, uprawnienia właściciela na tabelach docelowych oraz następujące konfiguracje i uprawnienia.

Uwaga

Usługa Azure Data Lake Storage jest jedynym typem magazynu usługi Azure obsługiwanym przez Unity Catalog.

Lokalizacja usługi Storage Dotacja
abfss://autoloader-source@<storage-account>.dfs.core.windows.net/json-data READ FILES
abfss://dev-bucket@<storage-account>.dfs.core.windows.net READ FILES, , WRITE FILESCREATE TABLE

Użyj Auto Loader do ładowania do tabeli zarządzanej przez Unity Catalog

W poniższych przykładach pokazano, jak używać Auto Loader do wprowadzania danych do tabeli zarządzanej w Unity Catalog.

Python

checkpoint_path = "abfss://dev-bucket@<storage-account>.dfs.core.windows.net/_checkpoint/dev_table"

(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .load("abfss://autoloader-source@<storage-account>.dfs.core.windows.net/json-data")
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .trigger(availableNow=True)
  .toTable("dev_catalog.dev_database.dev_table"))

SQL

CREATE OR REFRESH STREAMING TABLE dev_catalog.dev_database.dev_table
AS SELECT * FROM STREAM read_files(
  'abfss://autoloader-source@<storage-account>.dfs.core.windows.net/json-data',
  format => 'json'
);

Gdy używasz read_files w instrukcji CREATE STREAMING TABLE w potokach Lakeflow, punkty kontrolne i lokalizacje schematów są automatycznie zarządzane.

Załaduj do tabeli zewnętrznej Unity Catalog za pomocą Auto Loader

Aby przechowywać dane w określonej lokalizacji pamięci masowej, użyj zewnętrznej tabeli Unity Catalog zamiast tabeli zarządzanej. Na przykład użyj tabeli zewnętrznej, aby udostępnić dane klientom spoza usługi Databricks lub zarejestrować istniejące dane. W przypadku tabel zewnętrznych ustawiasz ścieżkę przechowywania. Zobacz Praca z tabelami zewnętrznymi.

Aby używać Auto Loader z tabelą zewnętrzną w Unity Catalog, najpierw zarejestruj tabelę za pomocą CREATE TABLE ... LOCATION, a następnie zapisuj do niej strumieniowo, używając jej nazwy. Lokalizacja tabeli musi znajdować się w lokalizacji zewnętrznej , w której masz CREATE EXTERNAL TABLE uprawnienia. Lokalizacja punktu kontrolnego musi również znajdować się w lokalizacji zewnętrznej zarządzanej przez Unity Catalog. Użyj oddzielnej ścieżki od danych tabeli.

checkpoint_path = "abfss://dev-bucket@<storage-account>.dfs.core.windows.net/_checkpoint/dev_table"
table_path = "abfss://dev-bucket@<storage-account>.dfs.core.windows.net/external/dev_table"

# One-time: register the external table in UC.
spark.sql(f"""
  CREATE TABLE IF NOT EXISTS dev_catalog.dev_database.dev_table
  USING DELTA
  LOCATION '{table_path}'
""")

(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .load("abfss://autoloader-source@<storage-account>.dfs.core.windows.net/json-data")
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .trigger(availableNow=True)
  .toTable("dev_catalog.dev_database.dev_table"))