Suivre le lineage des données dans une action Custom
La Data Platform enregistre automatiquement les événements de lineage pour les actions Load et Aggregate, en Python comme en PySpark (y compris le lineage au niveau schéma et colonne)
Objectif
La Data Platform enregistre automatiquement les événements de lineage pour les actions Load et Aggregate, en Python comme en PySpark (y compris le lineage au niveau schéma et colonne). Pour les actions Custom (Python et PySpark) et les notebooks, le lineage est opt-in : vous décidez quels datasets déclarer comme entrées et sorties.
Vous voulez visualiser le lineage ? Explorez-le directement dans la vue Lineage intégrée du Lakehouse Manager, ou transférez tous les événements de lineage vers votre propre solution compatible OpenLineage (par exemple Marquez) : configurez le consommateur OpenLineage dans les Connectors, puis utilisez l'action DPE Send OpenLineage Events pour envoyer les événements selon une planification ou en continu.
Démarrage rapide avec lineage_run
Le gestionnaire de contexte lineage_run est la méthode recommandée pour suivre le lineage. Il gère automatiquement le cycle de vie complet :
- Génère un
run_id(UUID v4) - Émet un événement
STARTà l'entrée - Émet un événement
COMPLETEen cas de sortie réussie - Émet un événement
FAILen cas d'exception (puis relève votre erreur) - Les erreurs de lineage ne font jamais planter votre action. Tous les appels d'émission sont encapsulés dans un
try/except
Passez les noms de tables sous forme de chaînes simples (database/table). La platform attache automatiquement le namespace correct.
Exemples
1. Transformation SQL simple
2. bulk_insert avec de nouvelles colonnes
3. Détection automatique du schéma avec connector=
Lorsque vous passez un connecteur, lineage_run appelle automatiquement connector.get_table_schema() pour chaque entrée et sortie, et enrichit l'événement COMPLETE avec des facets de schéma (noms et types de colonnes). L'événement START est émis avec des entrées/sorties simples (sans schéma). Si get_table_schema échoue pour une table, celle-ci est conservée telle quelle sans faire planter l'exécution.
4. Connecteurs multiples (inter-bases de données)
Lorsque les entrées et sorties couvrent différentes bases de données, passez un dict associant chaque préfixe de base de données à son connecteur. Chaque table est routée vers le connecteur correct pour la détection du schéma. Les tables sans préfixe correspondant sont conservées telles quelles.
5. Schéma manuel avec schema_facet
Si vous souhaitez déclarer les schémas explicitement (sans connecteur), utilisez l'assistant schema_facet. Il construit le dict OpenLineage correct avec _producer et _schemaURL automatiquement.
6. Jointure avec lineage au niveau colonne
Pour les jointures ou les transformations complexes, utilisez column_lineage_facet pour déclarer de quelles tables et colonnes d'entrée proviennent les colonnes de sortie.
Lorsque connector= et un column_lineage_facet sont utilisés ensemble, le connecteur ajoute automatiquement des facets de schéma aux assets qui n'en possèdent pas déjà un. Les assets ayant déjà des facets (comme le lineage de colonnes) sont conservés tels quels.
7. Accéder à l'ID de run
Le gestionnaire de contexte renvoie un objet LineageRun avec un attribut run_id.
Référence API
lineage_run(job_name, ...)
Paramètres
Renvoie un objet LineageRun avec un attribut run_id.
schema_facet(name, fields)
Construit un dict de dataset avec un SchemaDatasetFacet. Ajoute _producer et _schemaURL automatiquement.
Paramètres
Renvoie un dict utilisable dans inputs ou outputs.
Exemple
column_lineage_facet(name, mappings)
Construit un dict de dataset avec un ColumnLineageDatasetFacet. Ajoute _producer et _schemaURL automatiquement. Le namespace est injecté depuis la configuration de la platform.
Paramètres
Chaque valeur dans mappings est un dict avec :
source(requis) : nom du dataset sourcefield(requis) : nom de la colonne sourceoperation(optionnel) : description de la transformation (par exemple"SUM","JOIN")
Renvoie un dict utilisable dans outputs.
Exemple
emit_lineage(event_type, job_name, ...)
Fonction bas niveau qui envoie un unique RunEvent OpenLineage. Utilisez lineage_run à la place pour la plupart des cas d'usage.
Paramètres
Format entrée/sortie
Chaque entrée dans inputs ou outputs peut être :
- Une chaîne : le nom du dataset (le namespace est ajouté automatiquement) :
- Un dict : lorsque vous devez attacher des facets ou redéfinir le namespace :
- Un assistant :
schema_facet(..)oucolumn_lineage_facet(..)qui renvoient des dicts.
Tous ces formats peuvent être mélangés dans la même liste.
Aller plus loin
Si vous avez besoin d'une formation ou d'une assistance technique pour la mise en oeuvre de nos solutions, contactez votre commercial ou cliquez sur ce lien pour obtenir un devis et demander une analyse personnalisée de votre projet à nos experts de l’équipe Professional Services.
Posez vos questions, faites-nous part de vos commentaires et interagissez directement avec l’équipe qui développe la Data Platform sur le canal Discord dédié.
Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre centre d'aide.
Rejoignez notre communauté d'utilisateurs.