Ingest souborů jako typ SOUBORU

Important

Tato funkce je v beta verzi. Správci pracovního prostoru můžou řídit přístup k této funkci ze stránky Previews . Viz Manage Azure Databricks preview.

Typ FILE ukládá a dotazuje odkazy na nestrukturované soubory (dokumenty, obrázky a zvuk) v tabulkách. Tato stránka ukazuje, jak objevovat soubory, načítat je jako FILE reference a postupně přijímat nové soubory, jakmile přicházejí.

Pro referenci o FILE typu viz FILE typ. Přehled přístupů k získávání nestrukturovaných dat naleznete v článku TYP souboru a nestrukturovaná data.

Note

FILE Sloupce nemají definované pořadí. Sloupec FILE nelze použít jako dělící sloupec, shlukovací sloupec ani Z-order klíč. Další informace najdete v tématu Omezení.

Režimy úložiště

Reference FILE může být uložena ve dvou režimech:

  • FILE EXTERNAL odkazuje na soubory, které již existují v jednom svazku Unity Catalog. Databricks nepodporuje ukládání FILE EXTERNAL referencí pro soubory uložené mimo svazky.
  • FILE MANAGED ukládá kopie souborů do úložiště spravovaného Unity Catalog. Soubory ze zdrojů mimo svazky, jako je SharePoint, Google Drive nebo SFTP, musí být přijaty a uloženy jako FILE MANAGED.

Použití list_files k objevování souborů

Použijte tabulkovoulist_files funkci tabulkové hodnoty k nalezení souborů dostupných na cestě. Vrací jeden řádek na soubor s , pathsize, modification_time, a odkazemFILE:

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

Pro objevení souborů ve zdroji, který vyžaduje připojení k Unity Catalog, například SharePoint, Google Drive nebo SFTP, přidejte parametr:connection

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

list_files ve výchozím nastavení objevuje soubory rekurzivně. Pro více informací viz list_files funkce s tabulkovými hodnotami.

Ingest souborů jako reference na SOUBORY

Vyberte přístup ke vstupu podle toho, kde soubory ukládáte. Pro odkazování na soubory již v Unity katalogu použijte FILE EXTERNAL. Pro ingest souborů z externího zdroje je zkopírovat do spravovaného úložiště jako FILE MANAGED.

Ingest svazkové soubory jako SOUBOR EXTERNÍ

Pro ingest souborů, které již existují v Unity Catalog svazku, použijte CREATE TABLE AS SELECT příkaz (CTAS) s .list_files Tím se vytvoří tabulka se sloupcem FILE EXTERNAL , který odkazuje na každý soubor přímo na místě, aniž by se kopíroval jeho obsah. Následující příklad vytváří tabulku documents s názvem souboru, metadaty a odkazem FILE pro každý soubor:

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

Ingest externích zdrojových souborů jako FILE MANAGED

Pro generování FILE referencí pro soubory ve zdroji, jako je SharePoint, Google Drive nebo SFTP, nejprve soubory uložte a uložte jako FILE MANAGED. FILE EXTERNAL není podporováno pro soubory uložené mimo svazky.

Následující příklad obsahuje soubory ze SharePoint do tabulkyFILE 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()

Používejte pipeline pro postupné přijímání nových souborů

Pro vstup nových souborů při jejich příchodu použijte streamovací tabulku v Lakeflow pipeline, která čte zdroj pomocí STREAM read_files(..., format => 'file'). Každá aktualizace pipeline zpracovává pouze soubory přidané po poslední aktualizaci. See read_files a Spark deklarativní pipeline.

Pro postupné streamování souborů ze zdroje, jako je Google Drive:

  1. Nastavte kanál pipeline na PREVIEW. Ingestování FILE referencí v pipeline vyžaduje kanál.PREVIEW

  2. Definujte streamovací tabulku, která čte zdroj s STREAM read_files(..., format => 'file'), jako v následujícím kódu:

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

Aplikujte aktualizace a vymazání s AUTO CDC

Streamovací ingest přidává nové soubory, ale nezaznamenává aktualizace ani mazání ze zdroje. Pro aplikaci těchto změn si přečtěte zdrojový změnový kanál pomocí AUTO CDC.

Výstraha

Databricks doporučuje nejprve umístit data o změně do spravované tabulky, jako je uvedeno v následujícím příkladu, a poté aplikovat AUTO CDC na tuto tabulku. Přímé přijetí AUTO CDC na STREAM read_files(..., readChangeFeed => true) opětovné čtení zdroje pro každý následný tok, což může zvýšit náklady na zpracování.

Přijímejte feed change ve dvou krocích. Následující příklad přijímá feed změn ze SharePoint a poté jej aplikuje na cílovou streamovací tabulku jako SCD typ 1:

  1. Zapište data o změně do tabulky streamování se spravovanými soubory, jak je uvedeno v následujícím kódu. Nastavte readChangeFeed => true na read_files vrácení feed změn, který zahrnuje sloupce _file_id, _sequence, a _is_deleted metadata.

    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. Použijte AUTO CDC změny z této tabulky na cílovou streamovací tabulku, jak je uvedeno v následujícím kódu. Použijte _file_id jako klíč _sequence , jako sloupec sekvence a _is_deleted k identifikaci smazání.

    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
    )
    

Převod inline binárních dat na FILE reference

Pokud tabulka již ukládá obsah souboru jako inline binární data, použijte create_file funkci k zápisu těchto dat do úložiště a vytvoření FILE reference.

Následující příklady používají uživatelsky generovanou tabulku raw_documents, s jedním sloupcem name a sloupcem content , který obsahuje binární data.

Zápis binárních dat do svazku jako SOUBOR EXTERNÍ

Pro zápis souborů do svazku Unity Catalog jako externích souborů předejte a destination_path do create_file, jak je uvedeno v následujícím kódu:

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

Zapisovat binární data do spravovaného úložiště jako FILE MANAGED

Pro ukládání souborů jako spravovaných souborů volejte create_file pouze s binárním obsahem. Když vynecháte destination_path, Unity Catalog nahraje obsah do spravovaného úložiště:

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

Další kroky