Aller au contenu principal

Tutoriel : créer un pipeline de traitement de fichiers avec le type FILE

info

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.

Apprenez à créer un pipeline en médaillon avec Lakeflow pipeline qui traite des documents non structurés de bout en bout. Cet exemple utilise le dataset d’exemple samples.sec.contracts, une collection d’accords juridiques déposés auprès de la SEC et stockés sous forme de PDF dans un volume Unity Catalog.

Le pipeline ingère les PDF en tant que références FILE gérées avec Auto Loader, analyse chaque document avec des fonctions d’IA, le classe par type d’accord et extrait des champs structurés pour chaque type.

Pour la référence de type, consultez typeFILE.

Dans ce tutoriel, vous allez :

  • Ingérez de manière incrémentielle des PDF de contrats à partir d'un volume sous forme de références FILE gérées avec Auto Loader.
  • Analysez chaque document avec la fonctionai_parse_document et classez-le avec la fonctionai_classify.
  • Extrayez des champs structurés pour chaque type d’accord avec la fonctionai_extract.

Le résultat est un pipeline de style médaillon : Bronze (références FILE gérées brutes), Silver (documents analysés et classés) et Gold (champs extraits par type d’accord). Voir Qu’est-ce que l’architecture lakehouse en médaillon ? pour plus d’informations. La couche Bronze est une table de streaming qui ingère les fichiers de manière incrémentielle, et les couches Silver et Gold sont des vues matérialisées qui ne se recalculent que lorsque leurs entrées changent.

Exigences

Pour terminer ce didacticiel, vous devez remplir les conditions suivantes :

  • Soyez en Connexion à un Workspace Databricks avec Unity Catalog activé.
  • Activez le type FILE pour votre workspace. Les administrateurs du workspace peuvent l’activer depuis la page Aperçus . Consultez Gérer les aperçus Databricks.
  • Avoir les autorisations nécessaires pour créer des tables dans un schéma et pour créer un pipeline.
  • Disposez d’un volume Unity Catalog sur lequel vous pouvez écrire. Vous déclarez ce volume comme étant le FileSpace de la table bronze, et Unity Catalog y copie les fichiers ingérés en tant que stockage géré.
  • Utilisez le canal de distribution « Aperçu ».

Le dataset samples.sec.contracts est disponible par default dans tous les workspaces. Ce tutoriel stocke les fichiers PDF ingérés en tant que références FILE MANAGED : Unity Catalog copie chaque fichier dans le volume que vous déclarez comme FileSpace de la table et le gère avec la table. Ainsi, la suppression de lignes rend les fichiers référencés éligibles au nettoyage de la mémoire (garbage collection), et la table ainsi que ses fichiers restent synchronisés. Pour adapter le pipeline à vos propres fichiers PDF, pointez le chemin source vers un volume contenant vos fichiers. Pour d’autres options d’ingestion, consultez Ingérer des fichiers en tant que type FILE.

Créer le pipeline de traitement de fichiers

Le pipeline traite les documents en trois étapes.

Étape 1. Bronze : ingérer des PDF bruts en tant que références de FICHIER gérées

Utilisez Auto Loader pour lire de manière incrémentielle les fichiers PDF de contrat à partir du volume. La lecture de fichiers avec format => 'file' capture une référence et des métadonnées pour chaque fichier sans matérialiser ses octets. La déclaration de la colonne en tant que FILE MANAGED copie chaque fichier dans le FileSpace de la table, le volume que vous définissez avec la propriété de table databricks.filespace-preview, afin que Unity Catalog gère les fichiers avec la table.

SQL
CREATE OR REFRESH STREAMING TABLE raw_contracts (
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(
'/Volumes/samples/sec/contracts/',
format => 'file');
  • Fonctionne pour les fichiers volumineux : un fichier PDF volumineux réside dans le FileSpace de la table, tandis que la ligne de la table stocke uniquement une référence FILE légère (uri, size, content_type, checksum). Comparez ceci avec le type BINARY, qui intègre les octets en ligne dans la ligne.
  • Cycle de vie des fichiers gérés : Unity Catalog copie chaque fichier ingéré dans le FileSpace de la table et le gère avec la table : la suppression de lignes rend les fichiers référencés éligibles au nettoyage de la mémoire, afin que la table et ses fichiers restent synchronisés. Pour plus de détails, consultez FILE MANAGED et FILE EXTERNAL.
  • Traitement incrémentiel : la table de streaming ingère les nouveaux fichiers de manière incrémentielle à mesure qu'ils arrivent dans la source, sans retraiter ceux qui existent déjà. Le dataset samples.sec.contracts dans cet exemple est statique, mais avec une source en direct, les nouveaux fichiers sont récupérés à chaque mise à jour du pipeline. Pour également propager les modifications et suppressions de la source, ingérez le flux de modification avec AUTO CDC. Voir Appliquer les mises à jour et les suppressions avec AUTO CDC.

Étape 2. Silver : analyser et classer les documents

Transmettez chaque FILE à la fonctionai_parse_document pour convertir le PDF brut en un VARIANT structuré contenant des éléments de document, des métadonnées de Layout et du texte. Comme ai_parse_document accepte une colonne FILE, il lit le document directement depuis le stockage et ne charge jamais les octets dans la mémoire des clusters.

SQL
CREATE OR REFRESH MATERIALIZED VIEW parsed_contracts AS
SELECT
path,
ai_parse_document(file) AS parsed
FROM raw_contracts;
remarque

Définir l'étape d'analyse en tant que vue matérialisée sur la table de streaming raw_contracts rend le calcul incrémentiel. Chaque mise à jour de pipeline exécute ai_parse_document uniquement sur les fichiers ajoutés depuis la dernière mise à jour, et non sur la table entière. Comme ai_parse_document est l'étape la plus coûteuse, cela évite de réanalyser les documents que vous avez déjà traités. Le refresh incrémentiel des vues matérialisées nécessite un compute Serverless ; exécutez le pipeline sur Serverless. See Spark Declarative Pipelines.

Ensuite, transmettez le résultat analysé à la fonctionai_classify pour attribuer à chaque document l’un des cinq types d’accord. Les documents présentant des erreurs d’analyse sont filtrés avant la classification. Cet exemple pin ai_classify à la version 2.1, qui renvoie la classification sous forme d’objet par étiquette ; lisez donc l’étiquette à partir de la clé value.

SQL
CREATE OR REFRESH MATERIALIZED VIEW classified_contracts AS
SELECT
path,
parsed,
ai_classify(
parsed,
'["affiliate_agreement", "marketing_agreement", "consulting_agreement", "hosting_agreement", "escrow_agreement"]',
map('version', '2.1')
):response[0].value::STRING AS contract_type
FROM parsed_contracts
WHERE is_variant_null(parsed:error_status);
astuce

Pour améliorer la précision de la classification, ajoutez des descriptions d'étiquettes et une option instructions à ai_classify. Voir la fonctionai_classify.

Étape 3. Gold : extraire les champs par type d'accord

Chaque type d'accord possède son propre ensemble de champs pertinents. Filtrez les documents classés selon un type, transmettez le contenu analysé à la fonctionai_extract avec un schéma des champs souhaités, puis aplatissez la réponse en colonnes typées. Cet exemple pin ai_extract à la version 2.1, dans laquelle chaque champ extrait est un objet ; lisez donc sa clé value.

L’exemple suivant crée la table Gold pour les accords de conseil :

SQL
CREATE OR REFRESH MATERIALIZED VIEW consulting_agreements AS
WITH extracted AS (
SELECT
path,
ai_extract(
parsed,
'["company_name", "consultant_name", "compensation_amount", "effective_date"]',
map('version', '2.1')
) AS fields
FROM classified_contracts
WHERE contract_type = 'consulting_agreement'
)
SELECT
path,
fields:response.company_name.value::STRING AS company_name,
fields:response.consultant_name.value::STRING AS consultant_name,
fields:response.compensation_amount.value::STRING AS compensation_amount,
fields:response.effective_date.value::STRING AS effective_date
FROM extracted;

Avec ces instructions, vous disposez d'un pipeline entièrement incrémentiel : à mesure que de nouveaux PDF de contrats arrivent dans le volume, Auto Loader les ingère en tant que références FILE gérées, ai_parse_document et ai_classify acheminent chaque document, et la vue matérialisée Gold consulting_agreements fait apparaître les champs extraits.

Exemples de notebooks

Les Notebooks suivants contiennent le pipeline complet de ce tutoriel. Ces Notebooks sont des codes source de pipeline, et non des Notebooks exécutables. Importez le Notebook correspondant à votre langue, puis spécifiez son chemin dans le champ Code source lors de la configuration du pipeline. Voir Configurer les pipelines.

Notebook SQL de pipeline de traitement de fichiers

Explorer par vous-même

Le pipeline classe les documents en cinq types d’accords, mais n’extrait les champs que pour consulting_agreement. Pour l’étendre, répétez l’étape Gold pour chaque type restant, en modifiant le filtre contract_type et le schéma ai_extract pour qu’ils correspondent aux champs pertinents pour ce type. Par exemple :

  • affiliate_agreement: party_1_name, party_2_name, commission_rate, payment_frequency
  • marketing_agreement: party_1_name, party_2_name, effective_date, territory
  • hosting_agreement: provider_name, customer_name, effective_date, term_length
  • escrow_agreement: owner_name, licensee_name, escrow_agent_name, software_name

Ressources supplémentaires