For AI agents: the complete documentation index is available at https://docs.ovhcloud.com/fr/llms.txt, the full documentation bundle is available at https://docs.ovhcloud.com/fr/llms-full.txt, and this page is available as Markdown at https://docs.ovhcloud.com/fr/guides/public-cloud/data-platform/developers-python-sdk-lineage.md.

Suivre le lineage des données dans une action Custom

Voir en Markdown

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.

Info

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 COMPLETE en cas de sortie réussie
  • Émet un événement FAIL en 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
from forepaas.dwh import connect
from forepaas.dwh.lineage import lineage_run

def my_custom(event):
    with lineage_run("custom_titanic_transform",
                     inputs=["default_dataset/titanic"],
                     outputs=["default_dataset/titanic_survivors"]):

        connector = connect("dwh/default_dataset/")
        connector.query("""
            CREATE TABLE IF NOT EXISTS titanic_survivors AS
            SELECT passengerid, name, sex, age, pclass, fare, embarked
            FROM titanic
            WHERE survived = 1
        """)
Info

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

from forepaas.dwh import connect
from forepaas.dwh.lineage import lineage_run

def my_custom(event):
    with lineage_run("custom_titanic_transform",
                     inputs=["default_dataset/titanic"],
                     outputs=["default_dataset/titanic_survivors"]):

        connector = connect("dwh/default_dataset/")
        connector.query("""
            CREATE TABLE IF NOT EXISTS titanic_survivors AS
            SELECT passengerid, name, sex, age, pclass, fare, embarked
            FROM titanic
            WHERE survived = 1
        """)

2. bulk_insert avec de nouvelles colonnes

from forepaas.dwh import connect, bulk_insert
from forepaas.dwh.lineage import lineage_run

def my_custom(event):
    with lineage_run("custom_titanic_newcolumns",
                     inputs=["default_dataset/titanic"],
                     outputs=["default_dataset/titanic"]):

        connector = connect("dwh/default_dataset/")
        df = connector.query("SELECT * FROM titanic")

        df["newsurvived"] = df["survived"].apply(lambda x: "Yes" if x == 1 else "No")
        df["newclass"] = df["pclass"].apply(lambda x: f"Class {x}")

        bulk_insert(connector, "titanic", df)

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.

from forepaas.dwh import connect, bulk_insert
from forepaas.dwh.lineage import lineage_run

def my_custom(event):
    connector = connect("dwh/default_dataset/")

    with lineage_run("custom_titanic_newcolumns",
                     inputs=["default_dataset/titanic"],
                     outputs=["default_dataset/titanic"],
                     connector=connector):

        df = connector.query("SELECT * FROM titanic")

        df["newsurvived"] = df["survived"].apply(lambda x: "Yes" if x == 1 else "No")
        df["newclass"] = df["pclass"].apply(lambda x: f"Class {x}")

        bulk_insert(connector, "titanic", df)

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.

from forepaas.dwh import connect, bulk_insert
from forepaas.dwh.lineage import lineage_run

def my_custom(event):
    cn_default = connect("dwh/default_dataset/")
    cn_analytics = connect("dwh/analytics_dataset/")

    with lineage_run("enrich_orders",
                     inputs=["default_dataset/raw_orders", "default_dataset/customers"],
                     outputs=["analytics_dataset/enriched_orders"],
                     connector={
                         "default_dataset": cn_default,
                         "analytics_dataset": cn_analytics,
                     }):

        orders = cn_default.query("SELECT * FROM raw_orders")
        customers = cn_default.query("SELECT * FROM customers")
        enriched = orders.merge(customers, on="customer_id", how="left")

        bulk_insert(cn_analytics, "enriched_orders", enriched)

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.

from forepaas.dwh import connect
from forepaas.dwh.lineage import lineage_run, schema_facet

def my_custom(event):
    with lineage_run("custom_titanic_transform",
                     inputs=[schema_facet("default_dataset/titanic", [
                         ("passengerid", "Integer"),
                         ("survived", "Integer"),
                         ("name", "String"),
                         ("sex", "String"),
                         ("age", "Number"),
                     ])],
                     outputs=[schema_facet("default_dataset/titanic_survivors", [
                         ("passengerid", "Integer"),
                         ("name", "String"),
                         ("sex", "String"),
                         ("age", "Number"),
                         ("pclass", "Integer"),
                         ("fare", "Number"),
                         ("embarked", "String"),
                     ])]):

        connector = connect("dwh/default_dataset/")
        connector.query("""
            CREATE TABLE IF NOT EXISTS titanic_survivors AS
            SELECT passengerid, name, sex, age, pclass, fare, embarked
            FROM titanic
            WHERE survived = 1
        """)

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.

from forepaas.dwh import connect
from forepaas.dwh.lineage import lineage_run, column_lineage_facet

def my_custom(event):
    connector = connect("dwh/default_dataset/")

    with lineage_run("test_with_join",
                     inputs=["default_dataset/titanic",
                             "default_dataset/chicago_calendar_full"],
                     outputs=[column_lineage_facet("default_dataset/titanic_enriched", {
                         "passengerid": {
                             "source": "default_dataset/titanic",
                             "field": "passengerid",
                         },
                         "name": {
                             "source": "default_dataset/titanic",
                             "field": "name",
                         },
                         "humidity": {
                             "source": "default_dataset/chicago_calendar_full",
                             "field": "humidity",
                             "operation": "JOIN",
                         },
                         "temperature": {
                             "source": "default_dataset/chicago_calendar_full",
                             "field": "temperature",
                             "operation": "JOIN",
                         },
                     })],
                     connector=connector):

        connector.query("""
            CREATE TABLE IF NOT EXISTS titanic_enriched AS
            SELECT t.passengerid, t.name, c.humidity, c.temperature
            FROM titanic t
            INNER JOIN chicago_calendar_full c
              ON t.passengerid = c.passengerid
        """)
Info

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.

from forepaas.dwh.lineage import lineage_run

def my_custom(event):
    with lineage_run("my_job",
                     inputs=["default_dataset/source"],
                     outputs=["default_dataset/target"]) as run:

        print(f"Run ID: {run.run_id}")
        # ... traitement ...

Référence API

lineage_run(job_name, ...)

from forepaas.dwh.lineage import lineage_run

Paramètres

NomTypeRequisDescription
job_namestrOuiUn identifiant stable pour le job
inputslistNonDatasets lus par le job (chaînes ou dicts)
outputslistNonDatasets écrits par le job (chaînes ou dicts)
run_idstrNonUUID pour ce run. Généré automatiquement si omis
job_facetsdictNonFacets de job OpenLineage additionnels
connectorconnector ou dictNonConnecteur pour l'enrichissement automatique du schéma. Passez un connecteur unique ou un dict associant un préfixe de base de données à un connecteur (voir exemple 4)

Renvoie un objet LineageRun avec un attribut run_id.

schema_facet(name, fields)

from forepaas.dwh.lineage import schema_facet

Construit un dict de dataset avec un SchemaDatasetFacet. Ajoute _producer et _schemaURL automatiquement.

Paramètres

NomTypeRequisDescription
namestrOuiNom du dataset (par exemple "default_dataset/orders")
fieldslist[tuple]OuiListe de tuples (column_name, column_type)

Renvoie un dict utilisable dans inputs ou outputs.

Exemple

schema_facet("default_dataset/orders", [
    ("order_id", "Integer"),
    ("amount", "Number"),
])
# Renvoie :
# {
#     "name": "default_dataset/orders",
#     "facets": {
#         "schema": {
#             "_producer": "https://gitlab.forepaas.com/...",
#             "_schemaURL": "https://openlineage.io/spec/facets/1-2-0/SchemaDatasetFacet.json",
#             "fields": [
#                 {"name": "order_id", "type": "Integer"},
#                 {"name": "amount", "type": "Number"}
#             ]
#         }
#     }
# }

column_lineage_facet(name, mappings)

from forepaas.dwh.lineage import column_lineage_facet

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

NomTypeRequisDescription
namestrOuiNom du dataset (par exemple "analytics_dataset/enriched")
mappingsdictOuiDict associant les noms de colonnes de sortie aux informations de champ d'entrée

Chaque valeur dans mappings est un dict avec :

  • source (requis) : nom du dataset source
  • field (requis) : nom de la colonne source
  • operation (optionnel) : description de la transformation (par exemple "SUM", "JOIN")

Renvoie un dict utilisable dans outputs.

Exemple

column_lineage_facet("analytics_dataset/enriched", {
    "order_id": {"source": "raw_orders", "field": "order_id"},
    "total":    {"source": "raw_orders", "field": "amount", "operation": "SUM"},
})

emit_lineage(event_type, job_name, ...)

from forepaas.dwh.lineage import emit_lineage, EventType

Fonction bas niveau qui envoie un unique RunEvent OpenLineage. Utilisez lineage_run à la place pour la plupart des cas d'usage.

Paramètres

NomTypeRequisDescription
event_typeEventType ou strOuiL'un de START, COMPLETE, ABORT, FAIL, OTHER
job_namestrOuiUn identifiant stable pour le job
run_idstrNonUUID pour ce run. Généré automatiquement si omis
inputslistNonDatasets lus par le job
outputslistNonDatasets écrits par le job
job_facetsdictNonFacets de job OpenLineage additionnels
event_timestrNonTimestamp ISO-8601. Par défaut datetime.utcnow()

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) :
    inputs = ["default_dataset/raw_orders"]
  • Un dict : lorsque vous devez attacher des facets ou redéfinir le namespace :
    inputs = [{"name": "default_dataset/raw_orders", "namespace": "custom_ns"}]
  • Un assistant : schema_facet(..) ou column_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.

Cette page vous a-t-elle aidé ?