4 weeks ago
This post walks through the architecture of The Third Eye, a continuous AI health and governed ModelOps control plane built entirely on native Databricks capabilities. The core constraint: it never recomputes a metric Databricks already computes. It reads Lakehouse Monitoring's own output tables, the lineage system tables, Unity Gateway's usage and guardrail tables, and MLflow/UC registry objects, correlates across them, scores composite risk with business context, and drives a governed action loop on top. If you're running enough models that per-model dashboards have stopped being useful, this is the layer I think is missing from most Databricks AI/ML deployments.
The single most useful artifact in the design phase was writing down, explicitly, what Databricks already computes versus what actually needs building. Skipping this step is how governance projects balloon into duplicate monitoring stacks.
| Drift / distribution stats | Lakehouse Monitoring (or Data Profiling): {output_schema}.{table}_profile_metrics, {output_schema}.{table}_drift_metrics | Per-column stats, consecutive and baseline drift, on a configured schedule | Read the tables directly. No PSI/KS reimplementation. |
| Model quality / accuracy | Lakehouse Monitoring, InferenceLog analysis type | Accuracy per model_id/version once ground truth is joined | Read from the profile table. |
| Lineage (table + column) | system.access.table_lineage, system.access.column_lineage, Lineage REST API | Automatic lineage across jobs, notebooks, pipelines, dashboards, DBSQL | Read directly for correlation. 1-year rolling retention on system tables (indefinite via Catalog Explorer/API since Sept 1, 2024); REST API returns one hop per call, so walk recursively for multi-hop. |
| Inference request/response logs | Classic inference tables, or the newer Unified Trace Table (Unity Gateway, OpenTelemetry, Beta) | Full request/response payloads, latency, status, model version served | Read. Prefer the Unified Trace Table for anything Gateway-routed. |
| Cost / usage / tokens | system.serving.served_entities, system.serving.endpoint_usage | Per-endpoint and per-served-entity usage, plus a usage_context map for custom attribution | Read directly for cost panels. No custom cost math. |
| PII/PHI, unsafe content | Unity Gateway AI Guardrails | Detection/blocking/filtering at the gateway | Read violation events as a risk signal. |
| Model registry | MLflow Model Registry + UC model objects | Registered models, versions, aliases, experiments | Read via MLflow API / UC objects for the asset inventory. |
| AI asset discovery | Unity Gateway AI Asset Registry | Catalog of governed models, agents, MCP servers, tools | Read as the primary zero-touch discovery source. |
| Fairness / bias | Lakehouse Monitoring fairness/bias support for classification models | Bias metrics on schedule, if configured | Read if configured; provision the monitor via API if not. |
The actual new engineering is seven things: cross-signal correlation, criticality-weighted composite scoring, confidence scoring on top of that, LLM-generated root-cause narrative, zero-touch discovery/registration glue, a governed action layer (retrain, champion/challenger, gated promotion), and one unified multi-channel alert digest.
Everything lives in governance.model_health as Delta tables. The design principle: Third Eye's tables are pointers and derived state, not copies of Databricks' own data. The native tables stay the system of record.
CREATE CATALOG IF NOT EXISTS governance; CREATE SCHEMA IF NOT EXISTS governance.model_health; CREATE TABLE governance.model_health.model_registry_map ( model_id STRING NOT NULL, model_name STRING, model_version STRING, serving_endpoint STRING, gateway_registered BOOLEAN, owning_team STRING, business_domain STRING, criticality_tier STRING, -- config, defaulted from tag/domain lakehouse_monitor_configured BOOLEAN, created_at TIMESTAMP, is_active BOOLEAN ) USING DELTA; CREATE TABLE governance.model_health.signal_index ( model_id STRING, signal_type STRING, -- drift | accuracy | cost | usage | guardrail | lineage_change source_table STRING, -- fully qualified native table this reads from computed_at TIMESTAMP, latest_value DOUBLE, -- normalized numeric snapshot for scoring raw_reference STRING -- pointer back to the full native record ) USING DELTA PARTITIONED BY (signal_type); CREATE TABLE governance.model_health.lineage_events ( model_id STRING, upstream_table STRING, event_type STRING, -- schema_change | new_write | etc event_time TIMESTAMP, entity_type STRING -- JOB | NOTEBOOK | PIPELINE | DASHBOARD_V3 | DBSQL_QUERY ) USING DELTA; CREATE TABLE governance.model_health.risk_scores ( model_id STRING, computed_at TIMESTAMP, drift_component DOUBLE, quality_component DOUBLE, cost_component DOUBLE, guardrail_component DOUBLE, criticality_weight DOUBLE, health_score DOUBLE, -- 0-100, weighted composite confidence DOUBLE, -- 0-1 health_tier STRING -- healthy | watch | at_risk | critical ) USING DELTA; CREATE TABLE governance.model_health.incidents ( incident_id STRING NOT NULL, model_id STRING, opened_at TIMESTAMP, trigger_signals STRING, -- JSON array of co-occurring signal_index rows lineage_context STRING, -- JSON, linked lineage_events if in-window root_cause_narrative STRING, -- LLM-generated root_cause_confidence DOUBLE, recommended_action STRING, -- investigate | no_action | remediate status STRING -- open | acknowledged | resolved ) USING DELTA; CREATE TABLE governance.model_health.remediation_suggestions ( incident_id STRING, generated_at TIMESTAMP, explanation STRING, -- plain-language, from foundation model suggested_actions STRING, -- JSON array, ranked by effort urgency STRING -- urgent | can_wait ) USING DELTA; CREATE TABLE governance.model_health.risk_weights_config ( criticality_tier STRING NOT NULL, w_drift DOUBLE, w_quality DOUBLE, w_cost DOUBLE, w_guardrail DOUBLE, criticality_weight DOUBLE ) USING DELTA;
Registering a model normally, via mlflow.register_model() or a UC model registration, should be the only onboarding step. An hourly job does the rest:
from databricks.sdk import WorkspaceClient def sync_model_registry(mlflow_client, gateway_client, w: WorkspaceClient): mlflow_models = mlflow_client.search_registered_models() gateway_assets = gateway_client.list_ai_asset_registry() known = spark.table("governance.model_health.model_registry_map") \ .select("model_id").collect() known_ids = {r.model_id for r in known} new_rows = [] for m in mlflow_models: if m.model_id not in known_ids: new_rows.append(build_registry_row(m, gateway_assets)) if new_rows: spark.createDataFrame(new_rows).write.mode("append") \ .saveAsTable("governance.model_health.model_registry_map") for row in new_rows: lineage = walk_lineage(row["serving_endpoint"], hops=1) # system table or REST API write_lineage_events(row["model_id"], lineage) if not row["lakehouse_monitor_configured"]: w.lakehouse_monitors.create( table_name=row["serving_endpoint"], assets_dir=f"/monitors/{row['model_id']}", output_schema_name="governance.model_health", inference_log=InferenceLogProfileType( problem_type="regression", # or classification prediction_col="prediction", timestamp_col="ts", granularities=["1 day"], model_id_col="model_version", ), )
The new model_registry_map row is the only thing that activates dashboard inclusion, Genie scope, and monitoring. The dashboard and Genie's metric views are parameterized off this table's contents, not hardcoded per model.
Read-only normalizers that write into signal_index. Example for the drift adapter:
def adapt_drift_signals(output_schema: str):
drift_df = spark.table(f"{output_schema}.drift_metrics") \
.filter(F.col("drift_type") == "consecutive") \
.select(
F.col("model_id"),
F.lit("drift").alias("signal_type"),
F.lit(f"{output_schema}.drift_metrics").alias("source_table"),
F.current_timestamp().alias("computed_at"),
F.col("js_distance").alias("latest_value"), # or your chosen distance metric
F.col("column_name").alias("raw_reference"),
)
drift_df.write.mode("append").saveAsTable("governance.model_health.signal_index")The cost/usage and guardrail adapters follow the same shape, reading from system.serving.endpoint_usage and the Gateway guardrail event tables respectively. Normalize into the same six-column signal_index shape so the correlation engine below doesn't need to know which native table a signal came from.
Runs as a Databricks Workflow task after each Lakehouse Monitoring refresh cycle. It's a time-window join, not a new statistical method:
WITH recent_signals AS (
SELECT *
FROM governance.model_health.signal_index
WHERE computed_at >= current_timestamp() - INTERVAL 2 HOURS
AND latest_value > (
SELECT threshold FROM governance.model_health.signal_thresholds t
WHERE t.signal_type = signal_index.signal_type
)
),
lineage_in_window AS (
SELECT le.model_id, le.upstream_table, le.event_type, le.event_time
FROM governance.model_health.lineage_events le
JOIN recent_signals rs
ON le.model_id = rs.model_id
AND le.event_time BETWEEN rs.computed_at - INTERVAL 2 HOURS AND rs.computed_at
)
SELECT
rs.model_id,
to_json(collect_list(struct(rs.signal_type, rs.latest_value, rs.computed_at))) AS trigger_signals,
to_json(collect_list(struct(lw.upstream_table, lw.event_type, lw.event_time))) AS lineage_context
FROM recent_signals rs
LEFT JOIN lineage_in_window lw ON rs.model_id = lw.model_id
GROUP BY rs.model_id
HAVING size(collect_list(lw.upstream_table)) > 0 -- signal + corroborating lineage event
OR count(DISTINCT rs.signal_type) >= 2 -- or 2+ independent signal types agreeIf a model has one weak signal and nothing corroborating it, this query produces no row for it, and no incident opens. That HAVING clause is the alert-fatigue control. It's as load-bearing as the join itself.
from pyspark.sql import functions as F
def compute_health_score(signals_df, weights_df):
joined = signals_df.join(weights_df, on="criticality_tier")
return joined \
.withColumn(
"health_score",
100 - (
F.col("w_drift") * F.col("drift_component")
+ F.col("w_quality") * (1 - F.col("quality_component"))
+ F.col("w_cost") * F.col("cost_component")
+ F.col("w_guardrail") * F.col("guardrail_component")
) * F.col("criticality_weight")
) \
.withColumn(
"confidence",
F.least(
F.col("evidence_volume") / F.lit(30.0),
F.col("baseline_window_completeness"),
F.when(F.col("ground_truth_available"), F.lit(1.0)).otherwise(F.lit(0.6)),
)
) \
.withColumn(
"health_tier",
F.when(F.col("health_score") >= 85, "healthy")
.when(F.col("health_score") >= 65, "watch")
.when(F.col("health_score") >= 40, "at_risk")
.otherwise("critical")
)w_drift, w_quality, w_cost, w_guardrail, and criticality_weight all come from risk_weights_config, keyed on tier. Never hardcoded, so a large drift on a Tier 3 experimental model can rank below a small drift on a Tier 1 regulated one.
Once an incident opens, a Workflow task calls a Foundation Model endpoint with the evidence bundle, entirely inside Databricks:
SELECT
incident_id,
ai_query(
'databricks-meta-llama-3-70b-instruct',
concat(
'Evidence bundle: ', trigger_signals, ' Lineage context: ', lineage_context, '. ',
'In 3 sentences: explain the likely root cause, ',
'list remediation actions ranked by effort, ',
'and classify urgency as urgent or can_wait.'
)
) AS explanation
FROM governance.model_health.incidents
WHERE status = 'open'
AND incident_id NOT IN (SELECT incident_id FROM governance.model_health.remediation_suggestions)Write the result to remediation_suggestions. The dashboard, Genie, and Copilot Studio are pure readers of this table. None of them re-derive or reformat the explanation independently, which avoids three slightly different stories about the same incident.
CREATE VIEW governance.model_health.vw_fleet_risk AS
SELECT
m.model_name, m.business_domain, m.criticality_tier,
r.health_score, r.health_tier, r.confidence, r.computed_at
FROM governance.model_health.risk_scores r
JOIN governance.model_health.model_registry_map m USING (model_id)
QUALIFY ROW_NUMBER() OVER (PARTITION BY m.model_id ORDER BY r.computed_at DESC) = 1;Auto-retraining, blast-radius analysis, and multi-workspace federation share the same data model and scoring engine but don't start until the above runs clean end to end on synthetic data.
Unity Catalog and Delta Lake; Lakehouse Monitoring / Data Profiling; system.access.table_lineage and column_lineage; Unity Gateway (AI Asset Registry, Unified Trace Table, system.serving.*, AI Guardrails); MLflow Model Registry; Databricks Workflows; Foundation Model APIs via ai_query(); Databricks AI/BI (Lakeview) dashboards; Genie Space with metric views; Microsoft Copilot Studio fanned out to Slack/Google Chat; PySpark and Databricks SQL for the adapters and scoring jobs.
I'd welcome feedback from anyone running Lakehouse Monitoring or Unity Gateway at fleet scale, particularly on the correlation window sizing in the query above, and on whether the signal_index pointer-table pattern holds up against your own retention requirements.