Dateien als FILE-Typ eintragen

Important

Dieses Feature befindet sich in der Betaversion. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.

Der FILE Typ speichert und fragt Verweise auf unstrukturierte Dateien (Dokumente, Bilder und Audio) in Tabellen ab. Diese Seite zeigt, wie man Dateien entdeckt, als FILE Referenzen einnimmt und neue Dateien schrittweise einführt, sobald sie eintreffen.

Für die Referenz zum Typ FILE siehe FILE Typ. Für einen Überblick über Ansätze zur Erfassung unstrukturierter Daten siehe FILE type and unstructured data.

Hinweis

FILE Spalten haben keine definierte Reihenfolge. Du kannst eine FILE Spalte nicht als Partitionierungsspalte, Clustering-Spalte oder Z-Ordnungsschlüssel verwenden. Weitere Informationen finden Sie unter Grenzwerte.

Speichermodi

Eine Referenz FILE kann in einem von zwei Modi gespeichert werden:

  • FILE EXTERNAL verweist auf Dateien, die bereits in einem Unity-Katalog-Volume existieren. Databricks unterstützt keine Speicherung FILE EXTERNAL von Referenzen für Dateien, die außerhalb der Volumes liegen.
  • FILE MANAGED Speichert Kopien von Dateien in vom Unity Catalog verwalteten Speicher. Dateien von Quellen außerhalb von Volumen, wie SharePoint, Google Drive oder SFTP, müssen als eingetragen und gespeichert werdenFILE MANAGED.

Verwenden list_files Sie zum Entdecken von Dateien

Verwenden Sie die list_files tabellenwertige Funktion tabellenwertige Funktion , um die auf einem Pfad verfügbaren Dateien zu entdecken. Sie gibt pro Datei eine Zeile mit ihrem path, size, , modification_timeund einer Referenz FILE zurück:

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

Um Dateien in einer Quellquelle, die eine Unity-Catalog-Verbindung benötigt, wie SharePoint, Google Drive oder SFTP, zu entdecken, fügen Sie den Parameter hinzuconnection:

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

list_files Entdeckt Dateien standardmäßig rekursiv. Um mehr zu erfahren, siehe list_files tabellenwertige Funktion.

Dateien als FILE-Referenzen eintragen

Wählen Sie eine Aufnahmemethode basierend darauf, wo Sie Ihre Dateien speichern. Um bereits in einem Unity-Catalog-Volume auf Dateien zuzugreifen, verwenden FILE EXTERNALSie . Um Dateien von einer externen Quelle zu übertragen, kopieren Sie sie in verwalteten Speicher als FILE MANAGED.

Importiere Volumendateien als EXTERNE Datei

Um bereits in einem Unity-Katalog-Volume vorhandene Dateien einzulesen, verwenden Sie eine CREATE TABLE AS SELECT (CTAS)-Anweisung mit list_files. Dies erstellt eine Tabelle mit einer Spalte FILE EXTERNAL , die jede Datei an Ort und Stelle referenziert, ohne deren Inhalt zu kopieren. Das folgende Beispiel erstellt eine Tabelle documents mit Dateinamen, Metadaten und einer FILE Referenz für jede Datei:

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

Externe Quellcodedateien als FILE MANAGED eintragen

Um Referenzen für Dateien in einer Quelle wie SharePoint, Google Drive oder SFTP zu erzeugenFILE, werden die Dateien zuerst eingetragen und als FILE MANAGEDabgelegt werden. FILE EXTERNAL wird für Dateien, die außerhalb der Volumes gespeichert sind, nicht unterstützt.

Das folgende Beispiel überträgt Dateien von SharePoint in eine FILE MANAGED Tabelle:

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()

Verwenden Sie Pipelines, um inkrementelle neue Dateien einzutragen

Um neue Dateien zu importieren, sobald sie eintreffen, verwenden Sie eine Streaming-Tabelle in einer Lakeflow-Pipeline, die die Quellcode mit STREAM read_files(..., format => 'file')liest. Jedes Pipeline-Update verarbeitet nur die nach dem letzten Update hinzugefügten Dateien. Siehe read_files und zünde deklarative Pipelines an.

Um Dateien schrittweise von einer Quelle wie Google Drive zu streamen:

  1. Setze den Kanal der Pipeline auf PREVIEW. Das Aufnehmen FILE von Referenzen in einer Pipeline erfordert den Kanal PREVIEW .

  2. Definiere eine Streaming-Tabelle, die die Quelle mit STREAM read_files(..., format => 'file')liest, wie im folgenden Code:

    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")
      )
    

Aktualisieren und Löschungen bei AUTO CDC anwenden

Ein Streaming-Ingest fügt neue Dateien hinzu, erfasst aber keine Updates oder Löschungen vom Quellcode. Um diese Änderungen anzuwenden, lies den Quell-Änderungsfeed mit AUTO CDC.

Warning

Databricks empfiehlt, die Änderungsdaten zuerst in einer verwalteten Tabelle zu landen, wie im folgenden Beispiel, und dann auf diese Tabelle anzuwenden AUTO CDC . Direkt anzuwenden AUTO CDC , um STREAM read_files(..., readChangeFeed => true) den Quellwechsel-Feed für jeden nachgelagerten Fluss erneut zu lesen, was die Verarbeitungskosten erhöhen könnte.

Nehmen Sie den Wechselfeed in zwei Schritten ein. Das folgende Beispiel speichert den Änderungsfeed von SharePoint und wendet ihn dann auf eine Ziel-Streaming-Tabelle als SCD-Typ 1 an:

  1. Schreiben Sie die Änderungsdaten in eine Streaming-Tabelle mit verwalteten Dateien, wie im folgenden Code. Aktiviert readChangeFeed => trueread_files , um den Änderungsfeed zurückzugeben, der die _file_idSpalten , _sequence, und _is_deleted Metadaten enthält.

    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. Verwenden Sie AUTO CDC , um die Änderungen aus dieser Tabelle auf eine Ziel-Streaming-Tabelle anzuwenden, wie im folgenden Code. Verwenden Sie _file_id sie als Schlüssel, _sequence als Sequenzspalte und _is_deleted zur Identifizierung von Löschungen.

    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
    )
    

Inline-Binärdaten in FILE-Referenzen umwandeln

Wenn eine Tabelle bereits Dateiinhalte als Inline-Binärdaten speichert, verwenden create_file Sie Function, um diese Daten in den Speicher zu schreiben und eine Referenz FILE zu erstellen.

Die folgenden Beispiele verwenden eine benutzerdefinierte Tabelle, raw_documents, mit einer name Spalte und einer Spalte content , die die Binärdaten enthält.

Schreibe binäre Daten auf ein Volume als FILE EXTERNAL

Um die Dateien als externe Dateien auf ein Unity-Catalog-Volume zu schreiben, übergebe ein destination_path an create_file, wie im folgenden Code:

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()

Schreibe binäre Daten in den verwalteten Speicher als FILE MANAGED

Um die Dateien stattdessen als verwaltete Dateien zu speichern, rufe create_file nur mit dem Binärinhalt. Wenn Sie destination_path, lädt Unity Catalog die Inhalte an den verwalteten Speicherort hoch:

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()

Nächste Schritte