Tutoriel : Construis un pipeline de traitement de fichiers avec le type 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.

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 :

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 FILE type 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ère FILE (uri, size, content_type, checksum). Comparez cela avec le BINARY type, qui inligne les octets de la ligne.
  • Cycle de vie des fichiers gérés : Unity Catalog copie chaque fichier ingéré dans les FileSpace tables 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.contracts jeu 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 avec AUTO 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