カスタムアクションでのデータ系譜の追跡
Data Platformは、LoadおよびAggregateアクション(PythonおよびPySparkを含む)のデータ系譜イベントを自動的に記録します(スキーマおよび列レベルの系譜を含む)
目的
Data Platform は、Load および Aggregate アクションに対して、Python と PySpark の両方(スキーマおよび列レベルのラインナンスを含む)について、自動的にラインナンスイベントを記録します。 カスタムアクション(Python および PySpark)および ノートブック については、ラインナンスはオプションです:どのデータセットを入力および出力として宣言するかを決定します。
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)。プラットフォームが自動的に正しい名前空間を付加します。
例
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
パラメータ
LineageRun オブジェクトを生成し、run_id 属性を持たせます。
schema_facet(name, fields)
from forepaas.dwh.lineage import schema_facet
データセット辞書をSchemaDatasetFacetで構築します。_producerと_schemaURLを自動的に追加します。
パラメータ
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を自動的に追加します。名前空間はプラットフォームの設定から注入されます。
パラメータ
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 を使用してください。
パラメータ
各エントリは、inputs または outputs のいずれかで指定できます。
- 文字列: データセット名(名前空間は自動的に追加されます):
inputs = ["default_dataset/raw_orders"]
- 辞書: ファセットを付け加えたり、名前空間を上書きする必要がある場合:
inputs = [{"name": "default_dataset/raw_orders", "namespace": "custom_ns"}]
- ヘルパー:
schema_facet(..) または column_lineage_facet(..) は、辞書を返すものです。
すべての形式を同じリスト内で混在させることができます。
さらに深く掘り下げる
もしトレーニングや技術的なサポートが必要な場合は、営業担当者にお問い合わせください、またはこのリンクをクリックして見積もりを依頼し、プロフェッショナルサービスの専門家にプロジェクトのカスタム分析を依頼してください。
質問をする、フィードバックを提供する、またはData Platformを構築するチームと直接交流するには、専用のDiscordチャネルをご利用ください。
サポートが必要な場合は、ヘルプセンターでリクエストを作成してください。
当社のユーザーコミュニティに参加してください。