Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
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.
Apprenez à construire un pipeline Medallion avec Lakeflow pipeline qui traite des documents non structurés de bout en bout. Cet exemple utilise l’échantillon samples.sec.contracts de jeu de données, une collection d’accords juridiques déposés auprès de la SEC stockés sous forme de PDF dans un volume du catalogue Unity.
Le pipeline ingère les PDF comme références gérées FILE avec Auto Loader, analyse chaque document avec des fonctions IA, le classe en un type d’accord, et extrait des champs structurés pour chaque type.
Pour la référence typographique, voir FILE type.
Dans ce tutoriel, vous allez :
- Ingérer progressivement les PDF contractuels d’un volume comme références gérées
FILEavec Auto Loader. - Analysez chaque document avec
ai_parse_documentfonction et classez-le avecai_classifyfonction. - Extraire des champs structurés pour chaque type d’accord avec
ai_extractune fonction.
Le résultat est un pipeline de type médaillon : bronze (références brutes gérées FILE ), argent (documents analysés et classifiés) et or (champs extraits par type d’accord). Pour plus d’informations, consultez Qu'est-ce que l'architecture Lakehouse Medallion. La couche bronze est une table de flux qui ingère progressivement les fichiers, et les couches argent et or sont des vues matérialisées qui ne se recalculent que lorsque leurs entrées changent.
Exigences
Pour suivre ce tutoriel, vous devez répondre aux exigences suivantes :
- Soyez connecté à un espace de travail Azure Databricks avec Unity Catalog activé.
- Activez le
FILEtype pour votre espace de travail. Les administrateurs de l’espace de travail peuvent l’activer depuis la page des Aperçus . Consultez Gérer les préversions d’Azure Databricks. - Avoir des autorisations pour créer des tables dans un schéma et pour créer un pipeline.
- Ayez un volume du catalogue Unity où vous pouvez écrire. Vous déclarez ce volume comme celui de
FileSpacela table de bronze , et le catalogue Unity copie les fichiers ingérés en tant que stockage géré. - Utilisez le canal Aperçu.
Le samples.sec.contracts jeu de données est disponible par défaut dans tous les espaces de travail. Ce tutoriel stocke les PDF ingérés comme FILE MANAGED références : Unity Catalog copie chaque fichier dans le volume que vous déclarez comme appartenant à la FileSpace table et le gère avec la table, de sorte que supprimer les lignes rend les fichiers référencés éligibles à la collecte des déchets et que la table et ses fichiers restent synchronisés. Pour adapter le pipeline à vos propres PDF, pointez le chemin source vers un volume contenant vos fichiers. Pour d’autres options d’ingestion, voir les fichiers Ingest comme type de fichier.
Créer le pipeline de traitement des fichiers
Le pipeline traite les documents en trois étapes.
Étape 1. Bronze : ingérer les PDF bruts comme références FICHIERS gérées
Utilisez Auto Loader pour lire progressivement les PDF contractuels 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 MANAGED copie chaque fichier dans la table FileSpace, le volume que vous définissez avec la databricks.filespace-preview propriété de table, afin que Unity Catalogue 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');
Python
from pyspark import pipelines as dp
@dp.table(
name="raw_contracts",
schema="path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED",
table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
)
def raw_contracts():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "file")
.load("/Volumes/samples/sec/contracts/")
)
-
Fonctionne pour les grands fichiers : un grand PDF se trouve dans la table
FileSpace, tandis que la ligne de table ne stocke qu’une référence légèreFILE(uri,size,content_type,checksum). Comparez cela avec leBINARYtype, qui inligne les octets de la ligne. -
Cycle de vie des fichiers gérés : Unity Catalog copie chaque fichier ingéré dans les
FileSpacetables et le gère avec la table : supprimer des lignes rend les fichiers référencés éligibles à la collecte des déchets, afin que la table et ses fichiers restent synchronisés. Pour plus de détails, voir FICHIER MANAGED et FICHIER EXTERNE. -
Traitement incrémental : la table de streaming ingère progressivement de nouveaux fichiers à leur arrivée dans la source, sans retraiter les fichiers existants. Le
samples.sec.contractsjeu de données dans cet exemple est statique, mais avec une source active, de nouveaux fichiers sont détectés à chaque mise à jour du pipeline. Pour également propager les modifications et suppressions de source, ingérez le flux de modifications avecAUTO CDC. Voir Appliquer les mises à jour et suppressions avec AUTO CDC.
Étape 2. Argent : analyser et classifier les documents
Passez chaque FILE fonction à ai_parse_document chaque pour convertir le PDF brut en une structure VARIANT contenant des éléments de document, des métadonnées de mise en page et du texte. Comme ai_parse_document il accepte une FILE colonne, il lit le document directement depuis le stockage et ne charge jamais les octets dans la mémoire du cluster.
SQL
CREATE OR REFRESH MATERIALIZED VIEW parsed_contracts AS
SELECT
path,
ai_parse_document(file) AS parsed
FROM raw_contracts;
Python
@dp.materialized_view(name="parsed_contracts")
def parsed_contracts():
return (
spark.read.table("raw_contracts")
.selectExpr("path", "ai_parse_document(file) AS parsed")
)
Remarque
Définir l’étape d’analyse comme une vue matérialisée sur la raw_contracts table de flux incrémentalise le calcul. Chaque mise à jour du pipeline s’exécute ai_parse_document uniquement sur les fichiers ajoutés depuis la dernière mise à jour, pas sur l’ensemble de la table. Comme ai_parse_document c’est l’étape la plus coûteuse, cela évite de réparer les documents que vous avez déjà traités. Un rafraîchissement progressif des vues matérialisées nécessite un calcul sans serveur ; Faites tourner le pipeline en serverless. Consultez Spark Declarative Pipelines.
Ensuite, passez la sortie analysée à ai_classify la fonction pour assigner à chaque document l’un des cinq types d’accord. Les documents avec erreurs d’analyse sont filtrés avant la classification. Cet exemple est épinglé ai_classify à la version 2.1, qui renvoie la classification comme un objet par étiquette, donc lis l’étiquette à partir de la value clé.
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);
Python
@dp.materialized_view(name="classified_contracts")
def classified_contracts():
return (
spark.read.table("parsed_contracts")
.filter("is_variant_null(parsed:error_status)")
.selectExpr(
"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""")
)
Tip
Pour améliorer la précision de la classification, ajoutez des descriptions d’étiquettes et une instructions option à ai_classify. Consultez Fonction ai_classify.
Étape 3. Or : extraire les champs par type d’accord
Chaque type d’accord possède son propre ensemble de champs pertinents. Filtrez les documents classifiés à un type, passez le contenu analysé pour ai_extract qu’il fonctionne avec un schéma des champs souhaités, puis aplatissez la réponse en colonnes tapées. Cet exemple relie ai_extract à la version 2.1, dans laquelle chaque champ extrait est un objet, donc lis sa value clé.
L’exemple suivant construit la table d’or 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;
Python
@dp.materialized_view(name="consulting_agreements")
def consulting_agreements():
return (
spark.read.table("classified_contracts")
.filter("contract_type = 'consulting_agreement'")
.selectExpr(
"path",
"""ai_extract(
parsed,
'["company_name", "consultant_name", "compensation_amount", "effective_date"]',
map('version', '2.1')
) AS fields""")
.selectExpr(
"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")
)
Avec ces instructions, vous obtenez un pipeline entièrement incrémental : à mesure que de nouveaux PDF de contrats arrivent dans le volume, Auto Loader les ingère comme références gérées FILE , ai_parse_document puis ai_classify route chaque document, et la consulting_agreements vue matérialisée en or fait apparaître les champs extraits.
Exemples de notebooks
Les carnets suivants contiennent la pipeline complète de ce tutoriel. Ces notebooks sont du code source pipeline, pas des notebooks exécutables. Importez le carnet pour votre langage, puis spécifiez son chemin dans le champ Code source lorsque vous configurez le pipeline. Consultez Configurer des pipelines.
SQL
Carnet de lecture SQL pipeline de traitement de fichiers
Obtenir un ordinateur portable
Python
Carnet Python de pipeline de traitement de fichiers
Obtenir un ordinateur portable
Explorez par vous-même
Le pipeline classe les documents en cinq types d’accord mais extrait des champs uniquement consulting_agreementpour . Pour l’étendre, répétez l’étape d’or pour chaque type restant, en modifiant le contract_type filtre et le ai_extract schéma pour correspondre aux champs pertinents à 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 additionnelles
-
Type
FILE - Ingérer les fichiers comme type de fichier
- Démarrage rapide des fonctions FILE
- En savoir plus sur le chargeur automatique. Consultez Qu’est-ce que Auto Loader ?.