Praca z danymi nieustrukturyzowanymi w dużych ilościach

Na tej stronie przedstawiono, jak przechowywać, zadawać zapytania i przetwarzać nieustrukturyzowane pliki danych przy użyciu woluminów Unity Catalog. Dowiesz się, jak przekazywać pliki, wysyłać zapytania o metadane, przetwarzać pliki za pomocą funkcji sztucznej inteligencji, stosować kontrolę dostępu i udostępniać woluminy innym organizacjom. Jeśli to możliwe, dołączono instrukcje dotyczące pracy z tym samouczkiem przy użyciu interfejsu użytkownika Eksploratora wykazu. Jeśli nie jest wyświetlana opcja Eksplorator wykazu , użyj podanych poleceń języka Python lub SQL.

Aby zapoznać się z pełnym omówieniem możliwości woluminów i przypadków użycia, odwołaj się do Co to są woluminy Unity Catalog?.

Uwaga / Notatka

Ten samouczek wykorzystuje funkcje AI do przetwarzania plików według ścieżek. Dostępny w wersji beta, typ pozwala FILE przechowywać odwołania do plików i metadane jako wartości kolumn w tabeli. Zobacz typ pliku i dane niestrukturalne.

Wymagania

  • Obszar roboczy usługi Azure Databricks z włączonym Unity Catalog.
  • CREATE CATALOG uprawnienia do magazynu metadanych. Zobacz Tworzenie katalogów. Jeśli nie możesz utworzyć wykazu, poproś administratora o dostęp lub użyj istniejącego katalogu, w którym masz CREATE SCHEMA uprawnienia.
  • Databricks Runtime 14.3 LTS lub nowsze.
  • W przypadku funkcji sztucznej inteligencji: obszar roboczy w obsługiwanym regionie.
  • Dla OpenSharing: uprawnienia CREATE SHARE i CREATE RECIPIENT do magazynu metadanych. Zobacz Bezpieczne udostępnianie danych i zasobów sztucznej inteligencji.

Krok 1. Tworzenie woluminu

Utwórz katalog, schemat i wolumin do przechowywania plików. Aby uzyskać szczegółowe instrukcje dotyczące zarządzania woluminami, zobacz Tworzenie i zarządzanie woluminami Unity Catalog.

Krok 1.1. Tworzenie wykazu i schematu

SQL

-- Create a catalog
CREATE CATALOG IF NOT EXISTS unstructured_data_lab;
USE CATALOG unstructured_data_lab;

-- Create a schema
CREATE SCHEMA IF NOT EXISTS raw;
USE SCHEMA raw;

Python

spark.sql("CREATE CATALOG IF NOT EXISTS unstructured_data_lab")
spark.sql("USE CATALOG unstructured_data_lab")
spark.sql("CREATE SCHEMA IF NOT EXISTS raw")
spark.sql("USE SCHEMA raw")

Eksplorator wykazu

  1. Kliknij ikonę Dane.Wykaz na pasku bocznym.
  2. Kliknij Utwórz> katalog.
  3. Wprowadź unstructured_data_lab jako nazwę katalogu.
  4. Kliknij pozycję Utwórz.
  5. Kliknij pozycję Wyświetl wykaz.

Na stronie wykazu:

  1. Kliknij pozycję Utwórz schemat.
  2. Wprowadź raw jako nazwę schematu.
  3. Kliknij pozycję Utwórz.

Krok 1.2. Tworzenie woluminu zarządzanego

SQL

CREATE VOLUME IF NOT EXISTS files_volume
COMMENT 'Volume for storing unstructured data files';

Python

spark.sql("""
    CREATE VOLUME IF NOT EXISTS files_volume
    COMMENT 'Volume for storing unstructured data files'
""")

Eksplorator wykazu

Na stronie schematu:

  1. Kliknij Utwórz>wolumin.
  2. Wprowadź files_volume jako nazwę woluminu.
  3. Sprawdź, czy wybrano wolumin zarządzany .
  4. Kliknij pozycję Utwórz.

Krok 2. Przekazywanie plików

Prześlij pliki na wolumin. Aby uzyskać kompleksowe przykłady zarządzania plikami, zobacz Praca z plikami w woluminach Unity Catalog.

Krok 2.1. Przekazywanie plików

Możesz użyć przykładów z databricks-datasets w tym samouczku lub przesłać własne pliki, korzystając z interfejsu użytkownika Eksploratora katalogów.

Uwaga / Notatka

Polecenia języka Python umożliwiają kopiowanie plików z databricks-datasets do woluminu, nawet jeśli nie znasz języka Python. Aby uzyskać instrukcje dotyczące uruchamiania poleceń w notesach, zobacz Zarządzanie notesami usługi Databricks .

Python

# Upload a single image file
dbutils.fs.cp(
    "dbfs:/databricks-datasets/flower_photos/roses/10090824183_d02c613f10_m.jpg",
    "/Volumes/unstructured_data_lab/raw/files_volume/rose.jpg"
)

# Upload a single PDF file
dbutils.fs.cp(
    "dbfs:/databricks-datasets/COVID/CORD-19/2020-03-13/COVID.DATA.LIC.AGMT.pdf",
    "/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf"
)

# Upload a directory
local_dir = "dbfs:/databricks-datasets/samples/data/mllib"
volume_path = "/Volumes/unstructured_data_lab/raw/files_volume/sample_files"

for file_info in dbutils.fs.ls(local_dir):
    source = file_info.path
    dest = f"{volume_path}/{file_info.name}"
    dbutils.fs.cp(source, dest, recurse=True)
    print(f"Uploaded: {file_info.name}")

Eksplorator wykazu

Kod Python na karcie Python wysyła dwa pliki (JPG i PDF) oraz katalog, który zawiera pliki .txt i .csv. Aby przekazać pliki przy użyciu Eksploratora katalogu:

  1. Na stronie woluminu kliknij opcję Prześlij do tego woluminu.
  2. W oknie dialogowym Przekazywanie plików w obszarze Pliki kliknij przycisk przeglądaj lub przeciągnij i upuść pliki do strefy upuszczania.
  3. W obszarze Wolumin docelowy sprawdź, czy wolumin utworzony w poprzednim kroku został wybrany.

Krok 2.2. Weryfikacja przesyłania

SQL

LIST '/Volumes/unstructured_data_lab/raw/files_volume/';

Python

files = dbutils.fs.ls("/Volumes/unstructured_data_lab/raw/files_volume/")
for f in files:
    print(f"{f.name}\t{f.size} bytes")

Eksplorator wykazu

Po przekazaniu plików są one wyświetlane na stronie woluminu. Kliknij nazwę pliku, aby wyświetlić podgląd, lub kliknij katalog, aby wyświetlić poszczególne pliki.

Alternatywa: użyj polecenia magic %fs

Użyj polecenia magicznego %fs.

%fs ls /Volumes/unstructured_data_lab/raw/files_volume/

Krok 3. Wykonywanie zapytań dotyczących metadanych pliku

Wykonywanie zapytań o informacje o pliku, aby zrozumieć, co znajduje się w woluminie. Aby uzyskać więcej wzorców wykonywania zapytań, zobacz Wyświetlanie listy i wykonywanie zapytań dotyczących plików w woluminach przy użyciu języka SQL.

Krok 3.1. Wyświetlanie metadanych pliku

SQL

SELECT
  path,
  _metadata.file_name,
  _metadata.file_size,
  _metadata.file_modification_time
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile'
);

Python

df = (
    spark.read
    .format("binaryFile")
    .option("recursiveFileLookup", "true")
    .load("/Volumes/unstructured_data_lab/raw/files_volume/")
)

df.select("path", "modificationTime", "length").show(truncate=False)

Eksplorator wykazu

Strona woluminu w Eksploratorze wykazu zawiera nazwę każdego pliku (w tym rozszerzenie), rozmiar i datę ostatniej modyfikacji .

Krok 4. Wykonywanie zapytań i przetwarzanie plików

Korzystanie z funkcji sztucznej inteligencji usługi Azure Databricks w celu wyodrębniania zawartości z dokumentów i analizowania obrazów. Aby zapoznać się z pełnym omówieniem funkcji sztucznej inteligencji, zobacz Wzbogacanie danych przy użyciu funkcji sztucznej inteligencji.

Uwaga / Notatka

Funkcje sztucznej inteligencji wymagają obszaru roboczego w obsługiwanym regionie. Zobacz Wzbogacanie danych przy użyciu funkcji sztucznej inteligencji.

Jeśli nie masz dostępu do funkcji sztucznej inteligencji, zamiast tego użyj standardowych bibliotek języka Python. Rozwiń poniższe sekcje alternatywne, aby uzyskać przykłady.

Krok 4.1. Analizowanie dokumentów

SQL

SELECT
  path AS file_path,
  ai_parse_document(content, map('version', '2.0')) AS parsed_content
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile',
  fileNamePattern => '*.pdf'
);

Python

result_df = spark.sql("""
    SELECT
      path AS file_path,
      ai_parse_document(content, map('version', '2.0')) AS parsed_content
    FROM read_files(
      '/Volumes/unstructured_data_lab/raw/files_volume/',
      format => 'binaryFile',
      fileNamePattern => '*.pdf'
    )
""")
display(result_df)
Alternatywa: Analizowanie plików PDF bez funkcji sztucznej inteligencji

Jeśli funkcje sztucznej inteligencji nie są dostępne w Twoim regionie, użyj bibliotek języka Python:

%pip install PyPDF2==3.0.1

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
from PyPDF2 import PdfReader
import io

@udf(returnType=StringType())
def extract_pdf_text(content):
    if content is None:
        return None
    try:
        reader = PdfReader(io.BytesIO(content))
        return "\n".join(page.extract_text() or "" for page in reader.pages)
    except Exception as e:
        return f"Error: {str(e)}"

df = spark.read.format("binaryFile") \
    .option("pathGlobFilter", "*.pdf") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/")

result_df = df.withColumn("text_content", extract_pdf_text("content"))
display(result_df.select("path", "text_content"))

Krok 4.2. Analizowanie obrazów

SQL

SELECT
  path,
  ai_query(
    'databricks-llama-4-maverick',
    'Describe this image in one sentence:',
    files => content
  ) AS description
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile',
  fileNamePattern => '*.{jpg,jpeg,png}'
)
WHERE _metadata.file_size < 5000000;

Python

result_df = spark.sql("""
    SELECT
      path,
      ai_query(
        'databricks-llama-4-maverick',
        'Describe this image in one sentence:',
        files => content
      ) AS description
    FROM read_files(
      '/Volumes/unstructured_data_lab/raw/files_volume/',
      format => 'binaryFile',
      fileNamePattern => '*.{jpg,jpeg,png}'
    )
    WHERE _metadata.file_size < 5000000
""")
display(result_df)
Alternatywa: wyodrębnianie metadanych obrazu bez funkcji sztucznej inteligencji

Aby wyodrębnić metadane obrazu bez funkcji sztucznej inteligencji:

%pip install pillow==10.4.0

from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from PIL import Image
import io

image_schema = StructType([
    StructField("width", IntegerType()),
    StructField("height", IntegerType()),
    StructField("format", StringType())
])

@udf(returnType=image_schema)
def get_image_info(content):
    if content is None:
        return None
    try:
        img = Image.open(io.BytesIO(content))
        return {"width": img.width, "height": img.height, "format": img.format}
    except:
        return None

df = spark.read.format("binaryFile") \
    .option("pathGlobFilter", "*.{jpg,jpeg,png}") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/")

result_df = df.withColumn("image_info", get_image_info("content"))
display(result_df.select("path", "image_info.*"))

Krok 4.3. Filtrowanie i analizowanie według nazwy pliku

Ten przykład filtruje pliki obrazów z podciągem "rose" w nazwie pliku.

SQL

SELECT
  path AS file_path,
  ai_query(
    'databricks-llama-4-maverick',
    'Describe this image in one sentence:',
    files => content
  ) AS description
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile',
  fileNamePattern => '*.{jpg,jpeg,png}'
)
WHERE _metadata.file_name ILIKE '%rose%';

Python

result_df = spark.sql("""
    SELECT
      path AS file_path,
      ai_query(
        'databricks-llama-4-maverick',
        'Describe this image in one sentence:',
        files => content
      ) AS description
    FROM read_files(
      '/Volumes/unstructured_data_lab/raw/files_volume/',
      format => 'binaryFile',
      fileNamePattern => '*.{jpg,jpeg,png}'
    )
    WHERE _metadata.file_name ILIKE '%rose%'
""")
display(result_df)

Krok 4.4. Łączenie plików z tabelami ustrukturyzowanymi

W tym przykładzie użyto numerów wierszy do parowania plików z przejazdami taksówek w celach demonstracyjnych. W środowisku produkcyjnym dołącz do istotnych kluczy biznesowych.

SQL

-- This example demonstrates joining file metadata with structured data
-- by pairing files with taxi trips using row numbers
WITH files_with_row AS (
  SELECT
    path,
    SPLIT(path, '/')[SIZE(SPLIT(path, '/')) - 1] AS file_name,
    length,
    ROW_NUMBER() OVER (ORDER BY path) AS file_row
  FROM read_files(
    '/Volumes/unstructured_data_lab/raw/files_volume/',
    format => 'binaryFile'
  )
),
trips_with_row AS (
  SELECT
    tpep_pickup_datetime,
    pickup_zip,
    dropoff_zip,
    fare_amount,
    ROW_NUMBER() OVER (ORDER BY tpep_pickup_datetime) AS trip_row
  FROM samples.nyctaxi.trips
  WHERE pickup_zip IS NOT NULL
  LIMIT 5
)
SELECT
  f.path,
  f.file_name,
  f.length,
  t.pickup_zip,
  t.dropoff_zip,
  t.fare_amount,
  t.tpep_pickup_datetime
FROM files_with_row f
INNER JOIN trips_with_row t ON f.file_row = t.trip_row;

Python

from pyspark.sql.functions import col, row_number, element_at, split
from pyspark.sql.window import Window

# Read files and add row numbers
files_df = spark.read.format("binaryFile") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/") \
    .withColumn("file_name", element_at(split(col("path"), "/"), -1))

files_with_row = files_df.alias("files") \
    .withColumn("file_row", row_number().over(Window.orderBy("path")))

# Get trips and add row numbers
trips_df = spark.table("samples.nyctaxi.trips") \
    .filter(col("pickup_zip").isNotNull()) \
    .limit(5)

trips_with_row = trips_df.alias("trips") \
    .withColumn("trip_row", row_number().over(Window.orderBy("tpep_pickup_datetime")))

# Join on row numbers
result_df = files_with_row \
    .join(trips_with_row, col("file_row") == col("trip_row"), "inner") \
    .select(
        "files.path",
        "files.file_name",
        "files.length",
        "trips.pickup_zip",
        "trips.dropoff_zip",
        "trips.fare_amount",
        "trips.tpep_pickup_datetime"
    )

display(result_df)

Krok 5. Stosowanie kontroli dostępu

Kontroluj, kto może odczytywać i zapisywać pliki w twoich woluminach. Aby dowiedzieć się więcej na temat zarządzania uprawnieniami w Unity Catalog, zobacz Zarządzanie uprawnieniami w Unity Catalog.

Krok 5.1. Udzielanie dostępu

SQL

-- Replace <user-or-group-name> with your workspace group or user name

-- Grant read access
GRANT READ VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
TO `<user-or-group-name>`;

-- Grant read and write access
GRANT READ VOLUME, WRITE VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
TO `<user-or-group-name>`;

-- Grant all privileges
GRANT ALL PRIVILEGES ON VOLUME unstructured_data_lab.raw.files_volume
TO `<user-or-group-name>`;

Python

# Replace <user-or-group-name> with your workspace group or user name
spark.sql("""
    GRANT READ VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
    TO `<user-or-group-name>`
""")

spark.sql("""
    GRANT READ VOLUME, WRITE VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
    TO `<user-or-group-name>`
""")

spark.sql("""
    GRANT ALL PRIVILEGES ON VOLUME unstructured_data_lab.raw.files_volume
    TO `<user-or-group-name>`
""")

Eksplorator wykazu

  1. Przejdź do karty Uprawnienia na stronie woluminu.
  2. Kliknij Grant.
  3. Wprowadź adres e-mail użytkownika lub nazwę grupy.
  4. Wybierz uprawnienia do udzielenia.
  5. Kliknij przycisk Potwierdź.

Krok 5.2. Wyświetlanie bieżących uprawnień

SQL

SHOW GRANTS ON VOLUME unstructured_data_lab.raw.files_volume;

Python

display(spark.sql("SHOW GRANTS ON VOLUME unstructured_data_lab.raw.files_volume"))

Eksplorator wykazu

Karta Uprawnienia na stronie woluminu pokazuje, którzy użytkownicy i grupy mają dostęp do woluminu.

Krok 6. Konfigurowanie pozyskiwania przyrostowego

Użyj modułu automatycznego ładującego, aby automatycznie przetwarzać nowe pliki w miarę ich przybycia do woluminu. Ten schemat jest przydatny w przypadku procesów roboczych ciągłego pobierania danych. Aby uzyskać więcej wzorców wczytywania, zobacz Typowe wzorce ładowania danych.

Krok 6.1. Tworzenie tabeli przesyłania strumieniowego

SQL

CREATE OR REFRESH STREAMING TABLE document_ingestion
SCHEDULE EVERY 1 HOUR
AS SELECT
  path,
  modificationTime,
  length,
  content,
  _metadata,
  current_timestamp() AS ingestion_time
FROM STREAM(read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/incoming/',
  format => 'binaryFile'
));

Python

from pyspark.sql.functions import current_timestamp, col

dbutils.fs.mkdirs("/Volumes/unstructured_data_lab/raw/files_volume/incoming/")

df = spark.readStream.format("cloudFiles") \
    .option("cloudFiles.format", "binaryFile") \
    .option("pathGlobFilter", "*.pdf") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/incoming/")

df_enriched = df \
    .withColumn("ingestion_time", current_timestamp()) \
    .withColumn("source_file", col("_metadata.file_path"))

query = df_enriched.writeStream \
    .option("checkpointLocation",
            "/Volumes/unstructured_data_lab/raw/files_volume/_checkpoints/docs") \
    .trigger(availableNow=True) \
    .toTable("document_ingestion")

query.awaitTermination()

Krok 7. Udostępnianie plików za pomocą funkcji OpenSharing

Bezpieczne udostępnianie woluminów użytkownikom w innych organizacjach przy użyciu funkcji OpenSharing. Przed udostępnieniem musisz utworzyć adresata. Odbiorca reprezentuje zewnętrzną organizację lub użytkownika, który może uzyskiwać dostęp do udostępnionych danych. Aby uzyskać informacje o konfiguracji adresatów, zobacz Tworzenie adresatów danych w usłudze OpenSharing (udostępnianie usługi Databricks-to-Databricks).

Krok 7.1: Utworzyć i skonfigurować udział

SQL

-- Create a share
CREATE SHARE IF NOT EXISTS unstructured_data_share
COMMENT 'Document files for partners';

-- Add the volume
ALTER SHARE unstructured_data_share
ADD VOLUME unstructured_data_lab.raw.files_volume;

-- Create a recipient
CREATE RECIPIENT IF NOT EXISTS <partner_org>
USING ID '<recipient-sharing-identifier>';

-- Grant access
GRANT SELECT ON SHARE unstructured_data_share
TO RECIPIENT <partner_org>;

Python

spark.sql("""
    CREATE SHARE IF NOT EXISTS unstructured_data_share
    COMMENT 'Document files for partners'
""")

spark.sql("""
    ALTER SHARE unstructured_data_share
    ADD VOLUME unstructured_data_lab.raw.files_volume
""")

spark.sql("""
    CREATE RECIPIENT IF NOT EXISTS <partner_org>
    USING ID '<recipient-sharing-identifier>'
""")

spark.sql("""
    GRANT SELECT ON SHARE unstructured_data_share
    TO RECIPIENT <partner_org>
""")

Krok 7.2. Uzyskiwanie dostępu do udostępnionych danych (jako adresata)

SQL

-- View available shares
SHOW SHARES IN PROVIDER <provider_name>;

-- Create a catalog from the share
CREATE CATALOG IF NOT EXISTS shared_documents
FROM SHARE <provider_name>.unstructured_data_share;

-- Query shared files
SELECT * EXCEPT (content), _metadata
FROM read_files(
  '/Volumes/shared_documents/raw/files_volume/',
  format => 'binaryFile'
)
LIMIT 10;

Python

spark.sql("SHOW SHARES IN PROVIDER <provider_name>").show()

spark.sql("""
    CREATE CATALOG IF NOT EXISTS shared_documents
    FROM SHARE <provider_name>.unstructured_data_share
""")

df = spark.read.format("binaryFile") \
    .load("/Volumes/shared_documents/raw/files_volume/")

df.select("path", "modificationTime", "length").show(10)

Krok 8. Czyszczenie plików

Usuń pliki, gdy nie są już potrzebne.

Python

# Delete a single file
dbutils.fs.rm("/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf")

# Delete a directory recursively
dbutils.fs.rm("/Volumes/unstructured_data_lab/raw/files_volume/sample_files/", recurse=True)

CLI

# Delete a single file
databricks fs rm dbfs:/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf

# Delete a directory recursively
databricks fs rm -r dbfs:/Volumes/unstructured_data_lab/raw/files_volume/sample_files/
Alternatywa: Użyj standardowego języka Python
import os
os.remove("/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf")

import shutil
shutil.rmtree("/Volumes/unstructured_data_lab/raw/files_volume/sample_files/")

Dodatkowe zasoby

Kontynuuj naukę o woluminach

Odwołania do funkcji SQL