Ingérer les fichiers comme type de fichier

Important

Cette fonctionnalité est en version bêta. Les administrateurs d’espace de travail peuvent contrôler l’accès à cette fonctionnalité à partir de la page Aperçus . Consultez Gérer les préversions d’Azure Databricks.

Le FILE type stocke et interroge des références à des fichiers non structurés (documents, images et audio) dans des tableaux. Cette page montre comment découvrir des fichiers, les intégrer comme FILE références, et ingérer progressivement de nouveaux fichiers au fur et à mesure qu’ils arrivent.

Pour la référence sur le FILE type, voir FILE type. Pour une vue d’ensemble des approches d’ingestion de données non structurées, voir TYPE de fichier et données non structurées.

Remarque

FILE Les colonnes n’ont pas d’ordre défini. Vous ne pouvez pas utiliser une FILE colonne comme colonne de partition, colonne de clustering ou clé d’ordre Z. Pour en savoir plus, consultez Limites.

Modes de stockage

Une FILE référence peut être stockée selon l’un des deux modes suivants :

  • FILE MANAGEDstocke des copies de fichiers dans un stockage géré par le catalogue Unity : les permissions sont gérées via la table, et supprimer des lignes rend les fichiers référencés éligibles à la collecte des ordures, afin que la table et ses fichiers restent synchronisés. Les fichiers provenant de sources extérieures aux volumes, tels que SharePoint, Google Drive ou SFTP, doivent être ingérés et stockés sous forme FILE MANAGEDde fichiers .
  • FILE EXTERNAL référence des fichiers déjà existants dans un volume du catalogue Unity. Databricks ne supporte pas le stockage FILE EXTERNAL de références pour les fichiers stockés en dehors des volumes.

Azure Databricks recommande FILE MANAGED des charges de travail bénéficiant des autorisations au niveau des fichiers et d’une conformité intégrée. Pour une comparaison entre la gouvernance et le comportement du cycle de vie, voir TYPE de fichier et données non structurées.

Utilisation list_files pour découvrir des fichiers

Utilisez la list_files fonction à valeurs de table fonction à valeurs de table pour découvrir les fichiers disponibles sur un chemin. Il retourne une ligne par fichier avec ses path, size, modification_time, et une FILE référence :

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

Pour découvrir des fichiers dans une source nécessitant une connexion Unity Catalog, comme SharePoint, Google Drive ou SFTP, ajoutez le connection paramètre :

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

list_files découvre les fichiers de façon récursive par défaut. Pour en savoir plus, voir list_files fonction à valeurs de table.

Ingestion des fichiers en tant que références FICHIERS

Choisissez une approche d’ingestion en fonction de l’endroit où vous stockez vos fichiers. Pour ingérer des fichiers provenant d’une source externe, copiez-les dans un stockage géré sous forme FILE MANAGEDde fichiers . Pour référencer des fichiers déjà présents dans un volume du catalogue Unity sans les copier, utilisez FILE EXTERNAL.

Ingérer les fichiers sources externes sous la forme FILE MANAGED

Pour générer FILE des références pour des fichiers dans une source telle que SharePoint, Google Drive ou SFTP, ingérez d’abord les fichiers et stockez-les sous forme FILE MANAGEDde fichiers . FILE EXTERNAL n’est pas prise en charge pour les fichiers stockés en dehors des volumes.

L’exemple suivant intègre des fichiers de SharePoint dans une FILE MANAGED table :

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

Ingérer les fichiers de volume en FICHIER EXTERNE

Pour ingérer des fichiers déjà existant dans un volume du catalogue Unity, utilisez une CREATE TABLE AS SELECT instruction (CTAS) avec list_files. Cela crée un tableau avec une FILE EXTERNAL colonne qui fait référence à chaque fichier en ligne, sans en copier le contenu. L’exemple suivant crée une documents table avec le nom du fichier, les métadonnées et une FILE référence pour chaque fichier :

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

Utilisez les pipelines pour ingérer progressivement de nouveaux fichiers

Pour ingérer de nouveaux fichiers au fur et à mesure qu’ils arrivent, utilisez une table de streaming dans un pipeline Lakeflow qui lit la source avec STREAM read_files(..., format => 'file'). Chaque mise à jour du pipeline ne traite que les fichiers ajoutés après la dernière mise à jour. Voir read_files et déclencher les pipelines déclaratifs.

Pour diffuser progressivement des fichiers depuis une source telle que Google Drive :

  1. Réglez le canal du pipeline à PREVIEW. Ingérer FILE des références dans un pipeline nécessite le PREVIEW canal.

  2. Définissons une table de flux qui lit la source avec STREAM read_files(..., format => 'file'), comme dans le code suivant :

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

Appliquez les mises à jour et suppressions avec AUTO CDC

Une ingestion de streaming ajoute de nouveaux fichiers mais ne capture pas les mises à jour ou suppressions de la source. Pour appliquer ces changements, lisez le flux source de modifications avec AUTO CDC.

Avertissement

Databricks recommande de placer d’abord les données de changement dans une table gérée, comme dans l’exemple suivant, puis de l’appliquer AUTO CDC à cette table. Appliquer AUTO CDC directement à STREAM read_files(..., readChangeFeed => true) relire l’alimentation de changement de source pour chaque flux en aval, ce qui peut augmenter les coûts de traitement.

Ingérer le flux de changement en deux étapes. L’exemple suivant ingère le flux de changement de SharePoint, puis l’applique à une table de streaming cible sous le nom de type SCD 1 :

  1. Écrivez les données de changement dans une table de flux avec des fichiers gérés, comme dans le code suivant. Activez readChangeFeed => trueread_files pour retourner le flux de changement, qui inclut les _file_idcolonnes , _sequence, et _is_deleted les métadonnées.

    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. Utilisez AUTO CDC pour appliquer les modifications de cette table à une table de streaming cible, comme dans le code suivant. À utiliser _file_id comme clé, _sequence comme colonne de séquence, et _is_deleted pour identifier les suppressions.

    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
    )
    

Convertir les données binaires en ligne en références FICHIER

Si une table stocke déjà le contenu du fichier sous forme de données binaires en ligne, utilisez create_file la fonction pour écrire ces données et produire une FILE référence.

Les exemples suivants utilisent une table générée par l’utilisateur, raw_documents, avec une name colonne et une content colonne contenant les données binaires.

Écrire des données binaires sur un stockage géré sous le nom de FICHIER GÉRÉ

Pour stocker les fichiers en tant que fichiers gérés, appelez create_file uniquement avec le contenu binaire. Lorsque vous omettez destination_path, Unity Catalog télécharge le contenu dans l’emplacement de stockage géré :

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

Écrire des données binaires dans un volume en tant que FICHIER EXTERNE

Pour écrire les fichiers dans un volume du catalogue Unity en tant que fichiers externes, passez a destination_path à create_file, comme dans le code suivant :

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

Étapes suivantes