Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 54 additions & 0 deletions sdk/python/feast/openlineage/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,10 @@ def get_full_lineage_graph(
"schema": schema,
"description": ds.get("description"),
"source_type": ds.get("source_type"),
"owner": ds.get("owner_name"),
"owner_type": ds.get("owner_type"),
"lifecycle_state": ds.get("lifecycle_state"),
"current_version": ds.get("current_version"),
"facets": facets,
}
)
Expand Down Expand Up @@ -377,6 +381,56 @@ def get_full_lineage_graph(
"total_nodes": total_nodes,
}

# ── Dataset versioning endpoints ──

@router.get("/lineage/openlineage/datasets/{namespace}/{name}/versions")
def list_dataset_versions(
namespace: str,
name: str,
limit: int = Query(50, ge=1, le=500),
offset: int = Query(0, ge=0),
):
"""List version history for a dataset."""
ns_filter = _get_namespace_filter(get_allowed_namespaces)
if ns_filter is not None and namespace not in ns_filter:
return {"versions": [], "total": 0}
versions = store.get_dataset_versions(namespace, name, limit, offset)
for v in versions:
v["schema"] = _safe_parse_json(v.pop("schema_json", None))
v["facets"] = _safe_parse_json(v.pop("facets_json", None))
return {"versions": versions, "total": len(versions)}

# ── Column-level lineage endpoint ──

@router.get("/lineage/openlineage/datasets/{namespace}/{name}/columns")
def get_column_lineage(
namespace: str,
name: str,
direction: str = Query("both", pattern="^(both|upstream|downstream)$"),
):
"""Get column-level lineage for a dataset.

Shows which input fields map to which output fields, including
transformation types when available.
"""
ns_filter = _get_namespace_filter(get_allowed_namespaces)
if ns_filter is not None and namespace not in ns_filter:
return {"column_lineage": []}
cl = store.get_column_lineage(namespace, name, direction)
return {"column_lineage": cl}

# ── Execution tree endpoint ──

@router.get("/lineage/openlineage/runs/{run_id}/tree")
def get_run_tree(run_id: str):
"""Get the full execution tree (all runs sharing the same root).

Each run includes parent_run_id and root_run_id for hierarchy
reconstruction on the client side.
"""
tree = store.get_run_tree(run_id)
return {"runs": tree, "total": len(tree)}

# ── Run history endpoints ──

@router.get("/lineage/openlineage/runs")
Expand Down
139 changes: 134 additions & 5 deletions sdk/python/feast/openlineage/emitter.py
Original file line number Diff line number Diff line change
Expand Up @@ -567,10 +567,13 @@ def emit_saved_dataset_lineage(
"""
Emit lineage for a saved dataset definition.

Creates a definition job with inputs derived from:
- FeatureService (when feature_service_name is set)
- FeatureViews (extracted from feature refs in the format "view:feat")
- DataSources (matched by comparing storage location against registered sources)
Emits two events:
1. A **RunEvent** that captures the definition-time topology:
FeatureService/FeatureView → SavedDataset job.
2. A **DatasetEvent** with SymlinksDatasetFacet that links the
logical SavedDataset name to its physical storage URI, plus
OwnershipDatasetFacet and LifecycleStateChangeDatasetFacet
when the information is available.

Args:
saved_dataset: The SavedDataset object
Expand Down Expand Up @@ -675,7 +678,7 @@ def emit_saved_dataset_lineage(
**self._job_kind_facets(FeastJobKind.DEFINITION, project),
}

return self._client.emit_run_event(
run_result = self._client.emit_run_event(
job_name=f"saved_dataset_{saved_dataset.name}",
run_id=str(uuid.uuid4()),
event_type=RunState.COMPLETE,
Expand All @@ -684,12 +687,138 @@ def emit_saved_dataset_lineage(
job_facets=job_facets,
namespace=namespace,
)

# Emit a DatasetEvent with symlink, ownership, and lifecycle facets
self._emit_saved_dataset_metadata(
saved_dataset, project, namespace, registered_data_sources
)

return run_result
except Exception as e:
logger.error(
f"Error emitting saved dataset lineage for {saved_dataset.name}: {e}"
)
return False

def _emit_saved_dataset_metadata(
self,
saved_dataset: Any,
project: str,
namespace: str,
registered_data_sources: Optional[List[Any]] = None,
):
"""Emit a DatasetEvent for a SavedDataset with OL-standard facets.

Includes:
- SymlinksDatasetFacet linking logical name to physical storage URI
- OwnershipDatasetFacet from tags (``owner`` key) or the object owner
- LifecycleStateChangeDatasetFacet (CREATE)
- SchemaDatasetFacet from features
"""
try:
from openlineage.client.facet_v2 import (
lifecycle_state_change_dataset,
ownership_dataset,
schema_dataset,
symlinks_dataset,
)
except ImportError:
logger.debug(
"openlineage facet_v2 modules not available; "
"skipping DatasetEvent for SavedDataset"
)
return

ds_facets: Dict[str, Any] = {}

# SymlinksDatasetFacet: link logical name → physical storage URI
storage_uri = self._get_saved_dataset_storage_uri(saved_dataset)
if storage_uri:
ds_facets["symlinks"] = symlinks_dataset.SymlinksDatasetFacet(
identifiers=[
symlinks_dataset.Identifier(
namespace=namespace,
name=storage_uri,
type="TABLE",
)
]
)

# OwnershipDatasetFacet
owner = None
if hasattr(saved_dataset, "tags") and saved_dataset.tags:
owner = saved_dataset.tags.get("owner")
if not owner and hasattr(saved_dataset, "owner") and saved_dataset.owner:
owner = saved_dataset.owner
if owner:
ds_facets["ownership"] = ownership_dataset.OwnershipDatasetFacet(
owners=[
ownership_dataset.Owner(
name=owner,
type="FEAST_TAG",
)
]
)

# LifecycleStateChangeDatasetFacet
ds_facets["lifecycleStateChange"] = (
lifecycle_state_change_dataset.LifecycleStateChangeDatasetFacet(
lifecycleStateChange=lifecycle_state_change_dataset.LifecycleStateChange.CREATE,
)
)

# SchemaDatasetFacet
if saved_dataset.features:
ds_facets["schema"] = schema_dataset.SchemaDatasetFacet(
fields=[
schema_dataset.SchemaDatasetFacetFields(name=f, type="UNKNOWN")
for f in saved_dataset.features
]
)

try:
self._client.emit_dataset_event(
dataset_name=saved_dataset.name,
namespace=namespace,
facets=ds_facets,
)
except Exception as e:
logger.warning(
f"Failed to emit DatasetEvent for SavedDataset "
f"{saved_dataset.name}: {e}"
)

@staticmethod
def _get_saved_dataset_storage_uri(saved_dataset: Any) -> Optional[str]:
"""Extract a physical storage URI from a SavedDataset's storage."""
if not hasattr(saved_dataset, "storage") or not saved_dataset.storage:
return None
try:
storage_proto = saved_dataset.storage.to_proto()
except Exception:
return None

if hasattr(storage_proto, "file_url") and storage_proto.file_url:
return storage_proto.file_url
if hasattr(storage_proto, "bigquery_path") and storage_proto.bigquery_path:
return storage_proto.bigquery_path
if hasattr(storage_proto, "redshift_path") and storage_proto.redshift_path:
return storage_proto.redshift_path
if hasattr(storage_proto, "s3_path") and storage_proto.s3_path:
return storage_proto.s3_path
if hasattr(storage_proto, "trino_path") and storage_proto.trino_path:
return storage_proto.trino_path

for attr_name in dir(storage_proto):
if attr_name.startswith("_"):
continue
if "path" in attr_name or "url" in attr_name or "uri" in attr_name:
val = getattr(storage_proto, attr_name, None)
if val and isinstance(val, str):
return val

return None

def emit_materialize_start(
self,
feature_views: List["FeatureView"],
Expand Down
121 changes: 121 additions & 0 deletions sdk/python/feast/openlineage/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,25 @@

Uses SQLAlchemy Core (Table + Column) pattern consistent with Feast's
existing SQL registry in feast/infra/registry/sql.py.

Tables
------
Core:
openlineage_events, openlineage_jobs, openlineage_datasets,
openlineage_runs, openlineage_run_io, openlineage_lineage_edges,
openlineage_dataset_symlinks

Extended (OL spec facets):
openlineage_dataset_versions – per-run snapshots of dataset state
openlineage_column_lineage – column-level field transformations
openlineage_dataset_ownership – extracted OwnershipDatasetFacet rows
"""

from sqlalchemy import (
BigInteger,
Column,
Index,
Integer,
MetaData,
String,
Text,
Expand Down Expand Up @@ -71,6 +84,11 @@
Column("feast_object_type", String(100), nullable=True),
Column("feast_object_name", String(255), nullable=True),
Column("feast_project", String(255), nullable=True),
# Extended: ownership, lifecycle, version tracking
Column("owner_name", String(512), nullable=True),
Column("owner_type", String(100), nullable=True),
Column("lifecycle_state", String(50), nullable=True),
Column("current_version", Integer, nullable=True),
Column("updated_at", BigInteger, nullable=False),
)

Expand All @@ -84,6 +102,9 @@
Column("started_at", BigInteger, nullable=True),
Column("ended_at", BigInteger, nullable=True),
Column("facets_json", Text, nullable=True),
# Extended: execution hierarchy from ParentRunFacet
Column("parent_run_id", String(255), nullable=True),
Column("root_run_id", String(255), nullable=True),
Column("updated_at", BigInteger, nullable=False),
)

Expand Down Expand Up @@ -145,6 +166,66 @@
)


openlineage_dataset_versions = (
"openlineage_dataset_versions",
ol_metadata,
Column("version_id", Integer, primary_key=True, autoincrement=True),
Column("dataset_namespace", String(512), nullable=False),
Column("dataset_name", String(512), nullable=False),
Column("version", Integer, nullable=False),
Column("created_by_run_id", String(255), nullable=True),
Column("schema_json", Text, nullable=True),
Column("facets_json", Text, nullable=True),
Column("created_at", BigInteger, nullable=False),
UniqueConstraint(
"dataset_namespace",
"dataset_name",
"version",
name="uq_dataset_version",
),
)

openlineage_column_lineage = (
"openlineage_column_lineage",
ol_metadata,
Column("dataset_namespace", String(512), nullable=False),
Column("dataset_name", String(512), nullable=False),
Column("output_field", String(512), nullable=False),
Column("input_namespace", String(512), nullable=False),
Column("input_name", String(512), nullable=False),
Column("input_field", String(512), nullable=False),
Column("transformation_type", String(100), nullable=True),
Column("transformation_description", Text, nullable=True),
Column("run_id", String(255), nullable=True),
Column("updated_at", BigInteger, nullable=False),
UniqueConstraint(
"dataset_namespace",
"dataset_name",
"output_field",
"input_namespace",
"input_name",
"input_field",
name="uq_column_lineage",
),
)

openlineage_dataset_ownership = (
"openlineage_dataset_ownership",
ol_metadata,
Column("dataset_namespace", String(512), nullable=False),
Column("dataset_name", String(512), nullable=False),
Column("owner_name", String(512), nullable=False),
Column("owner_type", String(100), nullable=True),
Column("updated_at", BigInteger, nullable=False),
UniqueConstraint(
"dataset_namespace",
"dataset_name",
"owner_name",
name="uq_dataset_owner",
),
)


def _build_tables():
"""Build Table objects from the tuple definitions above."""
from sqlalchemy import Table
Expand All @@ -158,6 +239,9 @@ def _build_tables():
("run_io", openlineage_run_io),
("lineage_edges", openlineage_lineage_edges),
("dataset_symlinks", openlineage_dataset_symlinks),
("dataset_versions", openlineage_dataset_versions),
("column_lineage", openlineage_column_lineage),
("dataset_ownership", openlineage_dataset_ownership),
]:
tbl_name = tbl_def[0]
meta = tbl_def[1]
Expand Down Expand Up @@ -216,3 +300,40 @@ def _build_tables():
OL_TABLES["dataset_symlinks"].c.linked_namespace,
OL_TABLES["dataset_symlinks"].c.linked_name,
)

idx_dataset_versions_ds = Index(
"idx_ol_dv_dataset",
OL_TABLES["dataset_versions"].c.dataset_namespace,
OL_TABLES["dataset_versions"].c.dataset_name,
)
idx_dataset_versions_run = Index(
"idx_ol_dv_run",
OL_TABLES["dataset_versions"].c.created_by_run_id,
)
idx_column_lineage_ds = Index(
"idx_ol_cl_dataset",
OL_TABLES["column_lineage"].c.dataset_namespace,
OL_TABLES["column_lineage"].c.dataset_name,
)
idx_column_lineage_input = Index(
"idx_ol_cl_input",
OL_TABLES["column_lineage"].c.input_namespace,
OL_TABLES["column_lineage"].c.input_name,
)
idx_dataset_ownership_ds = Index(
"idx_ol_do_dataset",
OL_TABLES["dataset_ownership"].c.dataset_namespace,
OL_TABLES["dataset_ownership"].c.dataset_name,
)
idx_runs_parent = Index(
"idx_ol_runs_parent",
OL_TABLES["runs"].c.parent_run_id,
)
idx_runs_root = Index(
"idx_ol_runs_root",
OL_TABLES["runs"].c.root_run_id,
)
idx_datasets_lifecycle = Index(
"idx_ol_datasets_lifecycle",
OL_TABLES["datasets"].c.lifecycle_state,
)
Loading
Loading