Ingérer des fichiers en tant que type FILE
Bêta
Cette fonctionnalité est en bêta. Les administrateurs de Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Aperçus . Voir Gérer les prévisualisations Databricks.
Le type FILE stocke et query les références vers des fichiers non structurés (documents, images et audio) dans des tables. Cette page montre comment découvrir des fichiers, les ingérer en tant que références FILE et ingérer de manière incrémentielle les nouveaux fichiers à mesure qu’ils arrivent.
Pour obtenir des informations de référence sur le type FILE, consultez le typeFILE. Pour une vue d'ensemble des approches d'ingestion de données non structurées, consultez le type FILE et les données non structurées.
FILE les colonnes n’ont pas d’ordre défini. Vous ne pouvez pas utiliser une colonne FILE comme colonne de partitionnement, colonne de clustering ou clé Z-order. Pour plus d’informations, voir Limites.
Modes de stockage
Une référence FILE peut être stockée dans l'un des deux modes suivants :
FILE EXTERNALréférence des fichiers qui existent déjà dans un volume Unity Catalog. Databricks ne prend pas en charge le stockage de référencesFILE EXTERNALpour les fichiers stockés en dehors des volumes.FILE MANAGEDstocke des copies de fichiers dans le stockage géré par Unity Catalog. Les fichiers provenant de sources extérieures aux volumes, tels que SharePoint, Google Drive ou SFTP, doivent être ingérés et stockés en tant queFILE MANAGED.
Utilisez list_files pour découvrir des fichiers
Utilisez la list_files fonction de valeur tabulaire pour découvrir les fichiers disponibles à un chemin d'accès. Il renvoie une ligne par fichier avec ses références path, size, modification_time et FILE :
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, telle que SharePoint, Google Drive ou SFTP, ajoutez le parameter connection :
SELECT * FROM list_files('https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection');
list_files découvre les fichiers de manière récursive par default. Pour en savoir plus, consultez la fonction à valeur tabulairelist_files.
Ingérer les fichiers en tant que références FILE
Sélectionnez une approche d'ingestion en fonction de l'endroit où vous stockez vos fichiers. Pour référencer des fichiers déjà présents dans un volume Unity Catalog, utilisez FILE EXTERNAL. Pour ingérer des fichiers à partir d’une source externe, copiez-les dans le stockage géré en tant que FILE MANAGED.
Ingérer les fichiers de volume en tant que FILE EXTERNAL
Pour ingérer des fichiers qui existent déjà dans un volume Unity Catalog, utilisez une instruction CREATE TABLE AS SELECT (CTAS) avec list_files. Ceci crée une table avec une colonne FILE EXTERNAL qui référence chaque fichier sur place, sans copier son contenu. L’exemple suivant crée une table documents avec le nom de fichier, les métadonnées et une référence FILE pour chaque fichier :
CREATE TABLE documents AS
SELECT _metadata.file_name, *
FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');
Ingérer des fichiers sources externes en tant que FILE MANAGED
Pour générer des références FILE pour des fichiers dans une source telle que SharePoint, Google Drive ou SFTP, ingérez d’abord les fichiers et stockez-les en tant que FILE MANAGED. FILE EXTERNAL n’est pas pris en charge pour les fichiers stockés en dehors des volumes.
L’exemple suivant ingère des fichiers depuis SharePoint dans une table FILE MANAGED :
- SQL
- Python
- Scala
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');
(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())
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()
Utiliser des pipelines pour ingérer de manière incrémentielle de nouveaux fichiers
Pour ingérer de nouveaux fichiers dès leur arrivée, utilisez une table de streaming dans un pipeline Lakeflow qui lit la source avec STREAM read_files(..., format => 'file'). Chaque mise à jour de pipeline traite uniquement les fichiers ajoutés après la dernière mise à jour. Consultez read_files et Spark Declarative Pipelines.
Pour Stream des fichiers de manière incrémentielle à partir d'une source telle que Google Drive : [[ ## completed ##]]
- Définissez le canal de distribution du pipeline sur
PREVIEW. L’ingestion de référencesFILEdans un pipeline nécessite le canal de distributionPREVIEW. - Définissez un tableau de streaming qui lit la source avec
STREAM read_files(..., format => 'file'), comme dans le code suivant :
- SQL
- Python
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');
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")
)
Appliquer les mises à jour et les suppressions avec AUTO CDC
Une ingestion en streaming ajoute de nouveaux fichiers mais ne capture pas les mises à jour ou les suppressions provenant de la source. Pour appliquer ces modifications, lisez le flux de modification source avec AUTO CDC.
Databricks recommande de déposer d’abord les données de changement dans une table gérée, comme dans l’exemple suivant, puis d’appliquer AUTO CDC à cette table. L’application de AUTO CDC directement à STREAM read_files(..., readChangeFeed => true) relit le flux de changement source pour chaque flux en aval, ce qui pourrait augmenter les coûts de traitement.
Ingérez le flux de données de modification en deux étapes. L'exemple suivant ingère le flux de données de modification depuis SharePoint, puis l'applique à une table de streaming cible en tant que SCD de type 1 :
- Écrivez les données de modification dans une table de streaming avec des fichiers gérés, comme dans le code suivant. Définissez
readChangeFeed => truesurread_filespour renvoyer le flux de modification, qui inclut les colonnes de métadonnées_file_id,_sequenceet_is_deleted.
- SQL
- Python
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);
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/")
)
- Utilisez
AUTO CDCpour appliquer les changements de cette table à une table de streaming cible, comme dans le code suivant. Utilisez_file_idcomme clé,_sequencecomme colonne de séquence et_is_deletedpour identifier les suppressions.
- SQL
- Python
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;
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 FILE
Si une table stocke déjà le contenu des fichiers sous forme de données binaires en ligne, utilisez la fonctioncreate_file pour écrire ces données dans le stockage et produire une référence FILE.
Les exemples suivants utilisent une table générée par l'utilisateur, raw_documents, avec une colonne name et une colonne content qui contient les données binaires.
Écrire des données binaires dans un volume en tant que FILE EXTERNAL
Pour écrire les fichiers dans un volume Unity Catalog en tant que fichiers externes, transmettez un destination_path à create_file, comme dans le code suivant :
- SQL
- Python
- Scala
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;
(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())
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()
Écrire des données binaires dans le stockage géré en tant que FILE MANAGED
Pour stocker les fichiers en tant que fichiers gérés, appelez create_file avec uniquement le contenu binaire. Lorsque vous omettez destination_path, Unity Catalog effectue un upload du contenu vers l’emplacement de stockage géré :
[[ ## completed ##]]
- SQL
- Python
- Scala
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;
(spark.read.table("raw_documents")
.selectExpr("name", "create_file(content => content) AS file")
.writeTo("managed_documents").append())
spark.read.table("raw_documents")
.selectExpr("name", "create_file(content => content) AS file")
.writeTo("managed_documents").append()
Étapes suivantes
FILEType- Type de FILE et données non structurées
- Didacticiel : créer un pipeline de traitement de fichiers avec le type de fichier
- En savoir plus sur Auto Loader. Voir Qu’est-ce qu’Auto Loader ?.