Bestanden invoeren als het FILE-type

Important

Deze functie bevindt zich in de bètaversie. Werkruimtebeheerders kunnen de toegang tot deze functie beheren vanaf de pagina Previews . Zie Azure Databricks previews beheren.

Het FILE type slaat en zoekt verwijzingen naar ongestructureerde bestanden (documenten, afbeeldingen en audio) in tabellen. Deze pagina laat zien hoe je bestanden kunt ontdekken, ze als FILE referenties kunt invoeren en nieuwe bestanden stapsgewijs kunt invoeren zodra ze binnenkomen.

Voor de referentie over het FILE type, zie FILE type. Voor een overzicht van benaderingen voor het consumeren van ongestructureerde data, zie BESTANDSTYPE en ongestructureerde data.

Note

FILE Kolommen hebben geen vaste ordening. Je kunt een FILE kolom niet gebruiken als partitioneringskolom, clusterkolom of Z-orde sleutel. Zie Limietenvoor meer informatie.

Opslagmodi

Een FILE referentie kan worden opgeslagen in een van twee modi:

  • FILE EXTERNAL verwijst naar bestanden die al bestaan in een Unity Catalog-volume. Databricks ondersteunt geen opslag FILE EXTERNAL van referenties voor bestanden die buiten volumes zijn opgeslagen.
  • FILE MANAGED slaat kopieën van bestanden op in door Unity Catalog beheerde opslag. Bestanden van bronnen buiten volumes, zoals SharePoint, Google Drive of SFTP, moeten worden ingevoerd en opgeslagen als FILE MANAGED.

Gebruik list_files om bestanden te ontdekken

Gebruik de list_files tabelwaardige functie tabelwaarde functie om de beschikbare bestanden op een pad te ontdekken. Het geeft één rij per bestand terug met zijn path, size, modification_time, en een FILE referentie:

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

Om bestanden te ontdekken in een bron die een Unity Catalog-verbinding vereist, zoals SharePoint, Google Drive of SFTP, voeg je de parameter connection toe:

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

list_files Ontdekt bestanden standaard recursief. Voor meer informatie, zie list_files tabel-gewaardeerde functie.

Bestanden invoeren als FILE-referenties

Kies een invoermethode op basis van waar je je bestanden opslaat. Om bestanden te refereren die al in een Unity Catalog-volume staan, gebruik FILE EXTERNAL. Om bestanden van een externe bron te importeren, kopieer je ze naar beheerde opslag als FILE MANAGED.

Importeer volumebestanden als EXTERN BESTAND

Om bestanden die al bestaan in een Unity Catalog-volume in te voeren, gebruik je een CREATE TABLE AS SELECT (CTAS)-instructie met list_files. Dit creëert een tabel met een FILE EXTERNAL kolom die elk bestand ter plaatse verwijst, zonder de inhoud te kopiëren. Het volgende voorbeeld maakt een documents tabel met de bestandsnaam, metadata en een FILE referentie voor elk bestand:

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

Voer externe bronbestanden in als FILE MANAGED

Om referenties te genereren FILE voor bestanden in een bron zoals SharePoint, Google Drive of SFTP, neem je eerst de bestanden in en sla je ze op als FILE MANAGED. FILE EXTERNAL wordt niet ondersteund voor bestanden die buiten volumes zijn opgeslagen.

Het volgende voorbeeld voert bestanden van SharePoint in in een FILE MANAGED tabel:

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

Gebruik pijplijnen om nieuwe bestanden incrementeel te importeren

Om nieuwe bestanden te importeren zodra ze binnenkomen, gebruik je een streamingtabel in een Lakeflow-pipeline die de bron leest met STREAM read_files(..., format => 'file'). Elke pipeline-update verwerkt alleen de bestanden die na de laatste update zijn toegevoegd. Zie read_files en Ontsteek Declaratieve Pijplijnen.

Om bestanden stapsgewijs te streamen vanaf een bron zoals Google Drive:

  1. Stel het kanaal van de pijplijn in op PREVIEW. Het opnemen FILE van referenties in een pijplijn vereist het PREVIEW kanaal.

  2. Definieer een streamingtabel die de bron leest met STREAM read_files(..., format => 'file'), zoals in de volgende 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")
      )
    

Pas updates en verwijderingen toe bij AUTO CDC

Een streaming-ingest voegt nieuwe bestanden toe, maar registreert geen updates of verwijderingen van de bron. Om die wijzigingen toe te passen, lees je de bronwijzigingsfeed met AUTO CDC.

Warning

Databricks raadt aan om de wijzigingsdata eerst in een beheerde tabel te plaatsen, zoals in het volgende voorbeeld, en vervolgens toe te passen AUTO CDC op die tabel. Directe toepassing AUTO CDC op STREAM read_files(..., readChangeFeed => true) herlezen van de bronwijzigingsfeed voor elke downstream stroom, wat de verwerkingskosten kan verhogen.

Voer de wijzigingsfeed in twee stappen in. Het volgende voorbeeld voert de wijzigingsfeed van SharePoint in en past deze vervolgens toe op een doelstroomtabel als SCD type 1:

  1. Schrijf de wijzigingsgegevens in een streamingtabel met beheerde bestanden, zoals in de volgende code. Zet readChangeFeed => true op read_files om de wijzigingsfeed terug te geven, die de _file_idkolommen , _sequence, en _is_deleted metadata bevat.

    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. Gebruik AUTO CDC om de wijzigingen van die tabel toe te passen op een doelstroomtabel, zoals in de volgende code. Gebruik _file_id als sleutel, _sequence als sequentiekolom en _is_deleted om verwijderingen te identificeren.

    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
    )
    

Converteer inline binaire data naar FILE-referenties

Als een tabel al bestandsinhoud opslaat als inline binaire data, gebruik create_file dan de functie om die data naar opslag te schrijven en een FILE referentie te produceren.

De volgende voorbeelden gebruiken een door de gebruiker gegenereerde tabel, raw_documents, met een name kolom en een content kolom die de binaire gegevens bevat.

Schrijf binaire gegevens naar een volume als BESTAND EXTERN

Om de bestanden als externe bestanden naar een Unity Catalog-volume te schrijven, geef een destination_path door aan create_file, zoals in de volgende 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()

Schrijf binaire data naar beheerde opslag als FILE MANAGED

Om de bestanden als beheerde bestanden op te slaan, roep create_file je alleen de binaire inhoud aan. Wanneer je weglaat destination_path, uploadt Unity Catalog de inhoud naar de beheerde opslaglocatie:

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

Volgende stappen