For AI agents: the complete documentation index is available at https://docs.dataplatform.ovh.net/ja/llms.txt, the full documentation bundle is available at https://docs.dataplatform.ovh.net/ja/llms-full.txt, and this page is available as Markdown at https://docs.dataplatform.ovh.net/ja/developers-python-sdk-lineage.md.
  • 🇯🇵 日本語
  • カスタムアクションでのデータ系譜の追跡

    Data Platformは、LoadおよびAggregateアクション(PythonおよびPySparkを含む)のデータ系譜イベントを自動的に記録します(スキーマおよび列レベルの系譜を含む)

    目的

    Data Platform は、Load および Aggregate アクションに対して、Python と PySpark の両方(スキーマおよび列レベルのラインナンスを含む)について、自動的にラインナンスイベントを記録します。 カスタムアクション(Python および PySpark)および ノートブック については、ラインナンスはオプションです:どのデータセットを入力および出力として宣言するかを決定します。

    Info

    系譜を可視化したいですか? Lakehouse Manager に組み込まれた系譜ビューで直接確認できます。あるいは、すべての系譜イベントを独自の OpenLineage 互換ソリューション(例:Marquez)に転送することもできます。その場合は、Connectors で OpenLineage コンシューマーを設定し、次に Send OpenLineage Events DPE アクションを使用して、スケジュールに従って、または継続的にイベントを送信します。

    lineage_runを使用したクイックスタート

    lineage_run コンテキストマネージャーは、ラインナージを追跡するための推奨方法です。これは、ライフサイクル全体を自動的に処理します。

    • UUID v4 の run_id を生成します
    • 入力時に START イベントを発生させます
    • 正常終了時に COMPLETE イベントを発生させます
    • 例外が発生した場合、FAIL イベントを発生させます(その後、エラーを再スローします)
    • ライナージエラーはアクションをクラッシュさせません。すべての emit 呼び出しは 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

    テーブル名をプレーンな文字列として渡します (database/table)。プラットフォームが自動的に正しい名前空間を付加します。

    1. シンプルなSQL変換

    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

    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. connector=を使用した自動スキーマ検出

    コネクタを渡すと、lineage_run は自動的に各入力と出力に対して connector.get_table_schema() を呼び出し、COMPLETE イベントにスキーマのファセット(列名と型)を追加します。START イベントは、スキーマのないシンプルな入力/出力で発生します。get_table_schema がテーブルで失敗した場合、そのテーブルはそのまま保持され、クラッシュしません。

    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. 複数コネクタ(クロスデータベース)

    入力と出力が異なるデータベースにまたがる場合、各データベースのプレフィックスをそのコネクタにマッピングする辞書を渡します。各テーブルはスキーマ検出のために正しいコネクタにルーティングされます。一致するプレフィックスのないテーブルはそのまま保持されます。

    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. schema_facetを使用した手動スキーマ

    スキーマを明示的に宣言する場合(コネクタなし)、schema_facet ヘルパーを使用します。これにより、_producer_schemaURL を自動的に持つ正しい OpenLineage 辞書が構築されます。

    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. カラム系譜を使用した結合

    結合や複雑な変換を行う場合は、column_lineage_facet を使用して、出力列がどの入力テーブルや列から来るかを宣言します。

    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

    connector=column_lineage_facet を同時に使用する場合、コネクタは既にスキーマファセットを持っていないアセットに自動的にスキーマファセットを追加します。既存のファセット(例:列の系譜)を持つアセットは、そのまま保持されます。

    7. 実行IDにアクセス

    コンテキストマネージャーは、LineageRun オブジェクトを生成し、そのオブジェクトには 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}")
            # ... processing ...

    APIリファレンス

    lineage_run(job_name, ...)

    from forepaas.dwh.lineage import lineage_run

    パラメータ

    名前必須説明
    job_namestrYesジョブの安定した識別子です
    inputslistNoジョブで読み取られるデータセット(文字列または辞書)です
    outputslistNoジョブで書き込まれるデータセット(文字列または辞書)です
    run_idstrNoこの実行のUUIDです。省略された場合は自動生成されます
    job_facetsdictNo追加のOpenLineageジョブファセットです
    connectorconnector or dictNo自動スキーマエンリッチメント用のコネクタです。単一のコネクタまたはデータベースプレフィックスをコネクタにマッピングする辞書を渡します(例4を参照してください)

    LineageRun オブジェクトを生成し、run_id 属性を持たせます。

    schema_facet(name, fields)

    from forepaas.dwh.lineage import schema_facet

    データセット辞書をSchemaDatasetFacetで構築します。_producer_schemaURLを自動的に追加します。

    パラメータ

    名前必須説明
    namestrYesデータセット名(例:"default_dataset/orders"
    fieldslist[tuple]Yes(column_name, column_type) タプルのリスト

    inputs または outputs で使用するための辞書を返します。

    schema_facet("default_dataset/orders", [
        ("order_id", "Integer"),
        ("amount", "Number"),
    ])
    # Returns:
    # {
    #     "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

    データセットの辞書をColumnLineageDatasetFacetで構築します。_producer_schemaURLを自動的に追加します。名前空間はプラットフォームの設定から注入されます。

    パラメータ

    名前必須説明
    namestrYesデータセット名(例:"analytics_dataset/enriched"
    mappingsdictYes出力列名を入力フィールド情報にマッピングする辞書

    mappings 内の各値は、以下のような辞書です:

    • source (必須): ソースデータセット名
    • field (必須): ソース列名
    • operation (任意): 変換の説明 (例: "SUM", "JOIN")

    outputs で使用するのに適した辞書を返します。

    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

    単一の OpenLineage RunEvent を送信する低レベル関数です。ほとんどのユースケースでは、代わりに lineage_run を使用してください。

    パラメータ

    名前必須説明
    event_typeEventType または strYesSTARTCOMPLETEABORTFAILOTHER のいずれかです
    job_namestrYesジョブの安定した識別子です
    run_idstrNoこの実行のUUIDです。省略された場合は自動生成されます
    inputslistNoジョブで読み取られるデータセットです
    outputslistNoジョブで書き込まれるデータセットです
    job_facetsdictNo追加のOpenLineageジョブファセットです
    event_timestrNoISO-8601形式のタイムスタンプです。デフォルトはdatetime.utcnow()です

    入力/出力形式

    各エントリは、inputs または outputs のいずれかで指定できます。

    • 文字列: データセット名(名前空間は自動的に追加されます):
      inputs = ["default_dataset/raw_orders"]
    • 辞書: ファセットを付け加えたり、名前空間を上書きする必要がある場合:
      inputs = [{"name": "default_dataset/raw_orders", "namespace": "custom_ns"}]
    • ヘルパー: schema_facet(..) または column_lineage_facet(..) は、辞書を返すものです。

    すべての形式を同じリスト内で混在させることができます。

    さらに深く掘り下げる

    もしトレーニングや技術的なサポートが必要な場合は、営業担当者にお問い合わせください、またはこのリンクをクリックして見積もりを依頼し、プロフェッショナルサービスの専門家にプロジェクトのカスタム分析を依頼してください。

    質問をする、フィードバックを提供する、またはData Platformを構築するチームと直接交流するには、専用のDiscordチャネルをご利用ください。

    サポートが必要な場合は、ヘルプセンターでリクエストを作成してください。

    当社のユーザーコミュニティに参加してください。