Pobieraj pliki jako typ pliku

Important

Ta funkcja jest dostępna w wersji beta. Administratorzy obszaru roboczego mogą kontrolować dostęp do tej funkcji ze strony Podglądy . Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.

Typ FILE przechowuje i zapytuje odniesienia do nieustrukturyzowanych plików (dokumentów, obrazów i dźwięków) w tabelach. Ta strona pokazuje, jak odkrywać pliki, pobierać je jako FILE referencje oraz stopniowo pobierać nowe pliki w miarę ich pojawiania się.

Dla odniesienia do FILE tego typu, zobacz FILE typ. Aby uzyskać przegląd podejść do pozyskiwania danych niestrukturalnych, zobacz TYP pliku i dane niestrukturalne.

Uwaga / Notatka

FILE Kolumny nie mają określonego porządku. Nie można użyć kolumny FILE jako kolumny partycjonującej, klasterizacyjnej ani klucza Z-order. Aby uzyskać więcej informacji, zobacz Limity.

Tryby przechowywania

Referencja FILE może być przechowywana w jednym z dwóch trybów:

  • FILE EXTERNAL odnosi się do plików, które już istnieją w woluminie Unity Catalog. Databricks nie obsługuje przechowywania FILE EXTERNAL referencji dla plików przechowywanych poza woluminami.
  • FILE MANAGED przechowuje kopie plików w przechowywaniu zarządzanym przez Unity Catalog. Pliki ze źródeł spoza woluminów, takich jak SharePoint, Google Drive czy SFTP, muszą być pobierane i przechowywane jako FILE MANAGED.

Użyj list_files do odkrywania plików

Użyj list_files funkcji tabelowej o wartości tabelowej , aby odnaleźć pliki dostępne na danej ścieżce. Zwraca jeden wiersz na plik z , pathsize, modification_time, oraz odniesieniemFILE:

SELECT * FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');

Aby odkryć pliki w źródle wymagającym połączenia z Unity Catalog, takim jak SharePoint, Google Drive lub SFTP, dodaj parametr:connection

SELECT * FROM list_files('https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection');

list_files domyślnie odkrywa pliki rekurencyjne. Aby dowiedzieć się więcej, zobacz list_files funkcję tabelową.

Pobieranie plików jako odniesienia do plików

Wybierz podejście do pobierania plików w zależności od miejsca przechowywania plików. Aby odwołać się do plików już znajdujących się w woluminie Unity Catalog, użyj FILE EXTERNAL. Aby pobrać pliki z zewnętrznego źródła, skopiuj je do zarządzanej pamięci masowej jako FILE MANAGED.

Pobieranie plików woluminów jako PLIK ZEWNĘTRZNY

Aby pobrać pliki już istniejące w woluminie Unity Catalog, użyj CREATE TABLE AS SELECT instrukcji (CTAS) z .list_files Tworzy to tabelę z kolumną FILE EXTERNAL , która odwołuje się do każdego pliku w miejscu, bez kopiowania jego treści. Poniższy przykład tworzy tabelę documents z nazwą pliku, metadanymi oraz referencjami FILE dla każdego pliku:

CREATE TABLE documents AS
  SELECT _metadata.file_name, *
  FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');

Pobieranie zewnętrznych plików źródłowych jako FILE MANAGED

Aby generować FILE referencje dla plików w źródle takim jak SharePoint, Google Drive lub SFTP, najpierw pobierz pliki i przechowuj je jako FILE MANAGED. FILE EXTERNAL nie jest obsługiwany dla plików przechowywanych poza woluminami.

Poniższy przykład pobiera pliki z SharePoint do tabeliFILE MANAGED:

SQL

CREATE TABLE managed_documents (
  file_name STRING,
  path STRING,
  size BIGINT,
  modification_time TIMESTAMP,
  file FILE MANAGED
) USING DELTA
  TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');

INSERT INTO managed_documents
  SELECT _metadata.file_name, *
  FROM read_files(
    'https://example.sharepoint.com/sites/my-site/',
    connection => 'my_sharepoint_connection',
    format => 'file');

Python

(spark.read.format("file")
  .option("databricks.connection", "my_sharepoint_connection")
  .load("https://example.sharepoint.com/sites/my-site/")
  .selectExpr("_metadata.file_name", "*")
  .writeTo("managed_documents").append())

Scala

spark.read.format("file")
  .option("databricks.connection", "my_sharepoint_connection")
  .load("https://example.sharepoint.com/sites/my-site/")
  .selectExpr("_metadata.file_name", "*")
  .writeTo("managed_documents").append()

Używaj potoków do stopniowego pobierania nowych plików

Aby pobierać nowe pliki w momencie ich pojawiania się, użyj tabeli strumieniowej w potoku Lakeflow, która odczytuje źródło za pomocą STREAM read_files(..., format => 'file'). Każda aktualizacja potoku przetwarza tylko pliki dodane po ostatniej aktualizacji. Zobacz read_files i iskrz potoki deklaratywne.

Aby stopniowo przesyłać pliki ze źródła takiego jak Google Drive:

  1. Ustaw kanał potoku na .PREVIEW Pobieranie FILE referencji w potoku wymaga kanału PREVIEW .

  2. Zdefiniuj tabelę strumieniową, która odczytuje źródło z , STREAM read_files(..., format => 'file')jak w następującym kodzie:

    SQL

    CREATE STREAMING TABLE streaming_documents (
      path STRING,
      size BIGINT,
      modification_time TIMESTAMP,
      file FILE MANAGED
    )
    TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/')
    AS SELECT *
      FROM STREAM read_files(
        'https://drive.google.com/drive/folders/my-folder-id',
        connection => 'my_gdrive_connection',
        format => 'file');
    

    Python

    from pyspark import pipelines as dp
    
    @dp.table(
      name="streaming_documents",
      schema="path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED",
      table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
    )
    def streaming_documents():
      return (
        spark.readStream.format("cloudFiles")
          .option("cloudFiles.format", "file")
          .option("databricks.connection", "my_gdrive_connection")
          .load("https://drive.google.com/drive/folders/my-folder-id")
      )
    

Stosuj aktualizacje i usuwanie z AUTO CDC

Pobieranie streamingu dodaje nowe pliki, ale nie rejestruje aktualizacji ani usuwań ze źródła. Aby zastosować te zmiany, przeczytaj źródłowy kanał zmiany za pomocą AUTO CDC.

Ostrzeżenie

Databricks zaleca, aby najpierw umieścić dane zmiany w tabeli zarządzanej, jak w poniższym przykładzie, a następnie zastosować je AUTO CDC do tej tabeli. Bezpośrednie zastosowanie AUTO CDC do STREAM read_files(..., readChangeFeed => true) ponownego odczytu źródła dla każdego przepływu w dalszej fazie przepływu, co może zwiększyć koszty przetwarzania.

Pobieraj feed zmian w dwóch krokach. Poniższy przykład pobiera feed zmian z SharePoint, a następnie stosuje go do docelowej tabeli streamingowej jako SCD typu 1:

  1. Zapisz dane zmiany do tabeli streamingowej z plikami zarządzanymi, jak w poniższym kodzie. Ustaw readChangeFeed => true na read_files zwracanie zwrotu zmian, który zawiera kolumny _file_id, _sequence, i _is_deleted metadanych.

    SQL

    CREATE OR REFRESH STREAMING TABLE documents_changes (
      _file_id STRING,
      _sequence BIGINT,
      _is_deleted BOOLEAN,
      path STRING,
      size BIGINT,
      modification_time TIMESTAMP,
      file FILE MANAGED
    )
    TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/')
    AS SELECT *
      FROM STREAM read_files(
        'https://example.sharepoint.com/sites/my-site/',
        connection => 'my_sharepoint_connection',
        format => 'file',
        readChangeFeed => true);
    

    Python

    from pyspark import pipelines as dp
    
    @dp.table(
      name="documents_changes",
      table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
    )
    def documents_changes():
      return (
        spark.readStream.format("cloudFiles")
          .option("cloudFiles.format", "file")
          .option("databricks.connection", "my_sharepoint_connection")
          .option("cloudFiles.readChangeFeed", "true")
          .load("https://example.sharepoint.com/sites/my-site/")
      )
    
  2. Użyj do AUTO CDC zastosowania zmian z tej tabeli do docelowej tabeli streamingowej, jak w poniższym kodzie. Użyj _file_id jako klucza, _sequence jako kolumny sekwencji oraz do identyfikacji _is_deleted usuniętych miejsc.

    SQL

    CREATE OR REFRESH STREAMING TABLE documents
      TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');
    
    CREATE FLOW documents_cdc AS AUTO CDC INTO
      documents
    FROM STREAM documents_changes
      KEYS (_file_id)
      APPLY AS DELETE WHEN _is_deleted = true
      SEQUENCE BY _sequence
      COLUMNS * EXCEPT (_is_deleted, _sequence)
      STORED AS SCD TYPE 1;
    

    Python

    from pyspark import pipelines as dp
    from pyspark.sql.functions import col, expr
    
    dp.create_streaming_table(
      name="documents",
      table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
    )
    
    dp.create_auto_cdc_flow(
      target = "documents",
      source = "documents_changes",
      keys = ["_file_id"],
      sequence_by = col("_sequence"),
      apply_as_deletes = expr("_is_deleted = true"),
      except_column_list = ["_is_deleted", "_sequence"],
      stored_as_scd_type = 1
    )
    

Konwertowanie danych binarnych w linii na odniesienia do PLIKU

Jeśli tabela już przechowuje zawartość pliku jako dane binarne w linii, użyj create_file funkcji do zapisu tych danych na pamięć i wygenerowania referencji FILE .

Poniższe przykłady wykorzystują tabelę generowaną przez użytkownika, raw_documents, z kolumną name i kolumną content zawierającą dane binarne.

Zapisz dane binarne do woluminu jako FILE EXTERNAL

Aby zapisać pliki do wolumenu katalogu Unity jako pliki zewnętrzne, przekaż a destination_path do create_file, zgodnie z następującym kodem:

SQL

CREATE TABLE documents (name STRING, file FILE EXTERNAL) USING DELTA;

INSERT INTO documents (name, file)
  SELECT
    name,
    create_file(
      content => content,
      destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name
    )
  FROM raw_documents;

Python

(spark.read.table("raw_documents")
  .selectExpr(
    "name",
    "create_file(content => content, destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name) AS file")
  .writeTo("documents").append())

Scala

spark.read.table("raw_documents")
  .selectExpr(
    "name",
    "create_file(content => content, destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name) AS file")
  .writeTo("documents").append()

Zapisz dane binarne do zarządzanej pamięci jako FILE MANAGED

Aby przechowywać pliki jako pliki zarządzane, wywołaj create_file tylko zawartość binarną. Gdy pominiesz destination_path, Unity Catalog przesyła zawartość do zarządzanej lokalizacji pamięci masowej:

SQL

CREATE TABLE managed_documents (name STRING, file FILE MANAGED) USING DELTA
  TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');

INSERT INTO managed_documents (name, file)
  SELECT name, create_file(content => content)
  FROM raw_documents;

Python

(spark.read.table("raw_documents")
  .selectExpr("name", "create_file(content => content) AS file")
  .writeTo("managed_documents").append())

Scala

spark.read.table("raw_documents")
  .selectExpr("name", "create_file(content => content) AS file")
  .writeTo("managed_documents").append()

Następne kroki