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 externes avec Auto Loader, analyse chaque document avec des fonctions d’IA, le classe selon un 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 les PDF de contrats à partir d'un volume en tant que références FILE externes 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 externes 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é.
  • Avoir les autorisations nécessaires pour créer des tables dans un schéma et pour créer un pipeline.
  • Utilisez le canal de distribution « Aperçu ».

Le dataset samples.sec.contracts est disponible default dans tous les Workspace ; aucune configuration supplémentaire n’est donc requise. [[ ## completed ##]] Comme les fichiers se trouvent déjà dans un volume Unity Catalog, ce tutoriel les stocke en tant que références FILE EXTERNAL sans copier leur contenu. 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érez des fichiers PDF bruts en tant que références de FICHIER externes

Utilisez Auto Loader pour lire de manière incrémentielle les 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. Déclarer la colonne comme FILE EXTERNAL référence chaque fichier sur place, sans copier son contenu.

SQL
CREATE OR REFRESH STREAMING TABLE raw_contracts (
path STRING,
size BIGINT,
modification_time TIMESTAMP,
file FILE EXTERNAL
)
AS SELECT *
FROM STREAM read_files(
'/Volumes/samples/sec/contracts/',
format => 'file');
  • Fonctionne pour les fichiers volumineux : un PDF volumineux reste dans le volume, tandis que la ligne de 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.
  • 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 externes, ai_parse_document et ai_classify acheminent chaque document, et la vue matérialisée Gold consulting_agreements fait apparaître les champs extraits.

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