cancel
Showing results for 
Search instead for 
Did you mean: 
Community Articles
Dive into a collaborative space where members like YOU can exchange knowledge, tips, and best practices. Join the conversation today and unlock a wealth of collective wisdom to enhance your experience and drive success.
cancel
Showing results for 
Search instead for 
Did you mean: 

The Missing Correlation Layer in Databricks ModelOps | Building an AI Health Control Plane for DBx

ravikr1
Databricks Partner

Summary

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.

Scope boundary: decide this before writing any code

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.

Signal Databricks-native source What it gives you What Third Eye does
Drift / distribution statsLakehouse Monitoring (or Data Profiling): {output_schema}.{table}_profile_metrics, {output_schema}.{table}_drift_metricsPer-column stats, consecutive and baseline drift, on a configured scheduleRead the tables directly. No PSI/KS reimplementation.
Model quality / accuracyLakehouse Monitoring, InferenceLog analysis typeAccuracy per model_id/version once ground truth is joinedRead from the profile table.
Lineage (table + column)system.access.table_lineage, system.access.column_lineage, Lineage REST APIAutomatic lineage across jobs, notebooks, pipelines, dashboards, DBSQLRead 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 logsClassic inference tables, or the newer Unified Trace Table (Unity Gateway, OpenTelemetry, Beta)Full request/response payloads, latency, status, model version servedRead. Prefer the Unified Trace Table for anything Gateway-routed.
Cost / usage / tokenssystem.serving.served_entities, system.serving.endpoint_usagePer-endpoint and per-served-entity usage, plus a usage_context map for custom attributionRead directly for cost panels. No custom cost math.
PII/PHI, unsafe contentUnity Gateway AI GuardrailsDetection/blocking/filtering at the gatewayRead violation events as a risk signal.
Model registryMLflow Model Registry + UC model objectsRegistered models, versions, aliases, experimentsRead via MLflow API / UC objects for the asset inventory.
AI asset discoveryUnity Gateway AI Asset RegistryCatalog of governed models, agents, MCP servers, toolsRead as the primary zero-touch discovery source.
Fairness / biasLakehouse Monitoring fairness/bias support for classification modelsBias metrics on schedule, if configuredRead 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.

Unity Catalog governance schema

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.

 
sql
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;

Discovery/sync job (PySpark)

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:

 
python
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.

Signal adapters

Read-only normalizers that write into signal_index. Example for the drift adapter:

 
python
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.

Correlation engine

Runs as a Databricks Workflow task after each Lakehouse Monitoring refresh cycle. It's a time-window join, not a new statistical method:

 
sql
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 agree

If 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.

Scoring engine (PySpark)

 
python
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.

Root-cause narrative via ai_query()

Once an incident opens, a Workflow task calls a Foundation Model endpoint with the evidence bundle, entirely inside Databricks:

 
sql
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.

Serving layer

  • AI/BI Dashboard: Fleet Overview (KPI tiles for % healthy/watch/at_risk/critical, health_score trend, open incident count), Model Leaderboard (sortable by health_score/criticality, filterable by domain/team, drift sparkline, last-retrain date, cost trend), and a Model Detail drill-through (health_score time series with component breakdown, feature-level drift pulled straight from the native Lakehouse Monitoring table, a lineage snippet, open/past incidents with root-cause narrative, cost/usage trend).
  • Genie Space: metric views over risk_scores, incidents, and signal_index, plus an instruction doc mapping business vocabulary ("risky," "drifting," "stale," "costly") to specific columns/thresholds so Genie resolves questions without ad hoc joins per query. A metric view definition looks like:
 
sql
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;
  • Copilot Studio: two message types, both pure readers of the governed schema. A scheduled KPI digest posted to Teams and fanned out to Slack/Google Chat via webhook actions in the same topic, and a breach notification triggered by an incidents creation event, reading the linked remediation_suggestions row and formatting it as an adaptive card.

Platform limits to design around

  • Lineage isn't preserved across renames of catalogs/schemas/tables/columns, and doesn't exist before September 1, 2024.
  • Lineage system tables carry a rolling 1-year retention (Catalog Explorer/API retain indefinitely since Sept 2024), so snapshot anything needed for long-term trend charts into your own tables.
  • The Lineage REST API returns one hop per call. Walk recursively for multi-hop graphs, or query the system tables directly for a full join.
  • Trace/inference table delivery for Gateway-routed traffic is best-effort, with delays up to roughly an hour, and inference tables aren't guaranteed for 401/403/429/500 responses.

Build order

  1. Scaffolding: the DDL above, plus risk_weights_config and criticality defaults.
  2. Discovery/sync job: build against a synthetic/mock registry first if you don't have live Gateway access yet.
  3. Signal adapters: the largest glue surface, so budget the most time here.
  4. Correlation and scoring engine: the one truly novel computation, and it deserves the most test coverage.
  5. Root-cause narrative job.
  6. AI/BI dashboard (three pages), Genie Space (metric views plus instruction doc), alerting plus Copilot Studio topics.
  7. Model passport generator: a formatted read of the governed schema per model. Cheap, high demo value.
  8. README, architecture diagram, and a demo script that injects a synthetic drift/lineage-change scenario end to end: discovery, then correlation, then dashboard update, then Genie investigation, then a Teams alert with remediation.

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.

Tech stack

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.ChatGPT Image Sep 13, 2026, 09_00_11 PM.pngChatGPT Image Sep 13, 2026, 08_54_12 PM.pngChatGPT Image Sep 13, 2026, 08_52_09 PM.pngChatGPT Image Sep 13, 2026, 08_49_11 PM.pngChatGPT Image Sep 13, 2026, 08_47_53 PM.png

0 REPLIES 0