Samouczek: COPY INTO z usługą Spark SQL

Usługa Databricks zaleca użycie polecenia COPY INTO do przyrostowego i zbiorczego ładowania danych dla źródeł danych zawierających tysiące plików.

W tym samouczku użyjesz polecenia COPY INTO, aby załadować dane JSON z woluminu Katalogu Unity do tabeli Delta w obszarze roboczym Azure Databricks. Przykładowy zestaw danych Wanderbricks jest używany jako źródło danych. Aby uzyskać bardziej zaawansowane przypadki użycia Auto Loader, zobacz Co to jest moduł automatycznego ładowania?.

Wymagania

Krok 1. Konfigurowanie środowiska

Kod w tym samouczku używa woluminu Unity Catalog do przechowywania plików źródłowych JSON. Zastąp <catalog> katalogiem, w którym masz uprawnienia CREATE SCHEMA i uprawnienia CREATE VOLUME. Jeśli nie możesz uruchomić kodu, skontaktuj się z administratorem obszaru roboczego.

Utwórz notes i dołącz go do zasobu obliczeniowego. Następnie uruchom następujący kod, aby skonfigurować schemat i wolumin na potrzeby tego samouczka.

Python

# Set parameters and reset demo environment

catalog = "<catalog>"

username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first()[0]
schema = f"copyinto_{username}_db"
volume = "copy_into_source"
source = f"/Volumes/{catalog}/{schema}/{volume}"

spark.sql(f"SET c.catalog={catalog}")
spark.sql(f"SET c.schema={schema}")
spark.sql(f"SET c.volume={volume}")

spark.sql(f"DROP SCHEMA IF EXISTS {catalog}.{schema} CASCADE")
spark.sql(f"CREATE SCHEMA {catalog}.{schema}")
spark.sql(f"CREATE VOLUME {catalog}.{schema}.{volume}")

SQL

-- Reset demo environment

DROP SCHEMA IF EXISTS <catalog>.copy_into_tutorial CASCADE;
CREATE SCHEMA <catalog>.copy_into_tutorial;
CREATE VOLUME <catalog>.copy_into_tutorial.copy_into_source;

Krok 2. Zapisywanie przykładowych danych na woluminie w formacie JSON

Polecenie COPY INTO ładuje dane ze źródeł opartych na plikach. Odczyt z przykładowej tabeli Wanderbricksbookings i zapis partii rekordów jako plików JSON do woluminu, symulując nadchodzące dane z systemu zewnętrznego.

Python

# Write a batch of Wanderbricks bookings data as JSON to the volume

bookings = spark.read.table("samples.wanderbricks.bookings")
batch_1 = bookings.orderBy("booking_id").limit(20)
batch_1.write.mode("append").json(f"{source}/bookings")

SQL

Do zapisywania plików na woluminie wymagany jest język Python. W rzeczywistym przepływie pracy te dane zostaną dostarczone z systemu zewnętrznego.

%python
# Write a batch of Wanderbricks bookings data as JSON to the volume

bookings = spark.read.table("samples.wanderbricks.bookings")
batch_1 = bookings.orderBy("booking_id").limit(20)
batch_1.write.mode("append").json("/Volumes/<catalog>/copy_into_tutorial/copy_into_source/bookings")

Krok 3: Ładuj dane JSON idempotentnie przy użyciu COPY INTO

Utwórz docelową tabelę delty przed użyciem polecenia COPY INTO. Nie musisz podawać niczego innego niż nazwa tabeli w instrukcji CREATE TABLE . Ponieważ ta akcja jest idempotentna, usługa Databricks ładuje dane tylko raz, nawet jeśli kod jest uruchamiany wiele razy.

Python

# Create target table and load data

spark.sql(f"CREATE TABLE IF NOT EXISTS {catalog}.{schema}.bookings_target")

spark.sql(f"""
  COPY INTO {catalog}.{schema}.bookings_target
  FROM '/Volumes/{catalog}/{schema}/{volume}/bookings'
  FILEFORMAT = JSON
  FORMAT_OPTIONS ('mergeSchema' = 'true')
  COPY_OPTIONS ('mergeSchema' = 'true')
""")

SQL

-- Create target table and load data

CREATE TABLE IF NOT EXISTS <catalog>.copy_into_tutorial.bookings_target;

COPY INTO <catalog>.copy_into_tutorial.bookings_target
FROM '/Volumes/<catalog>/copy_into_tutorial/copy_into_source/bookings'
FILEFORMAT = JSON
FORMAT_OPTIONS ('mergeSchema' = 'true')
COPY_OPTIONS ('mergeSchema' = 'true')

Krok 4. Podgląd zawartości tabeli

Sprawdź, czy tabela zawiera 20 wierszy z pierwszej partii danych rezerwacji usługi Wanderbricks i czy schemat został poprawnie wywnioskowany z plików źródłowych JSON.

Python

# Review loaded data

display(spark.sql(f"SELECT * FROM {catalog}.{schema}.bookings_target"))

SQL

-- Review loaded data

SELECT * FROM <catalog>.copy_into_tutorial.bookings_target

Krok 5. Ładowanie większej ilości danych i podgląd wyników

Dodatkowe dane pochodzące z systemu zewnętrznego można symulować, zapisując kolejną partię rekordów i uruchamiając COPY INTO ponownie. Uruchom następujący kod, aby napisać drugą partię danych.

Python

# Write another batch of Wanderbricks bookings data as JSON

bookings = spark.read.table("samples.wanderbricks.bookings")
batch_2 = bookings.orderBy(bookings.booking_id.desc()).limit(20)
batch_2.write.mode("append").json(f"{source}/bookings")

SQL

Do zapisywania plików na woluminie wymagany jest język Python. W rzeczywistym przepływie pracy te dane zostaną dostarczone z systemu zewnętrznego.

%python
# Write another batch of Wanderbricks bookings data as JSON

bookings = spark.read.table("samples.wanderbricks.bookings")
batch_2 = bookings.orderBy(bookings.booking_id.desc()).limit(20)
batch_2.write.mode("append").json("/Volumes/<catalog>/copy_into_tutorial/copy_into_source/bookings")

Następnie ponownie uruchom COPY INTO polecenie z kroku 3 i wyświetl podgląd tabeli, aby potwierdzić nowe rekordy. Ładowane są tylko nowe pliki.

Python

# Confirm new data was loaded

display(spark.sql(f"SELECT COUNT(*) AS total_rows FROM {catalog}.{schema}.bookings_target"))

SQL

-- Confirm new data was loaded

SELECT COUNT(*) AS total_rows FROM <catalog>.copy_into_tutorial.bookings_target

Krok 6. Samouczek dotyczący czyszczenia

Po ukończeniu tego samouczka możesz wyczyścić skojarzone zasoby, jeśli nie chcesz ich przechowywać. Upuść schemat, tabele i wolumin i usuń wszystkie dane.

Python

# Drop schema and all associated objects

spark.sql(f"DROP SCHEMA IF EXISTS {catalog}.{schema} CASCADE")

SQL

-- Drop schema and all associated objects

DROP SCHEMA IF EXISTS <catalog>.copy_into_tutorial CASCADE;

Dodatkowe zasoby