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
336 changes: 336 additions & 0 deletions demos/dqx_demo_ai_assisted_unity_catalog_metadata.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,336 @@
# Databricks notebook source
# MAGIC %md
# MAGIC # AI-assisted rule generation — with and without Unity Catalog metadata
# MAGIC
# MAGIC Builds a synthetic bronze → silver → gold invoice pipeline, then calls
# MAGIC `DQGenerator.generate_dq_rules_ai_assisted` **twice** against the same table:
# MAGIC
# MAGIC 1. **Baseline** — `unity_catalog_metadata_config=None`; the LLM sees only column names
# MAGIC and types.
# MAGIC 2. **Enriched** — `UnityCatalogMetadataConfig()` with defaults; the LLM additionally
# MAGIC sees UC table / column comments, column tags, bounded column upstream lineage, and
# MAGIC an external upstream relationship from a SAP-origin `ExternalMetadata` object.
# MAGIC
# MAGIC The two generated rule sets are printed side-by-side at the end so you can see how
# MAGIC the Unity Catalog metadata shifts the model's choices. Comments / tags are deliberately
# MAGIC brief — the goal is to add semantics the LLM would otherwise have to guess, not to
# MAGIC write encyclopaedia entries.
# MAGIC
# MAGIC Disable a specific enrichment by setting its sub-model to `None`, e.g.
# MAGIC `UnityCatalogMetadataConfig(column_upstream_lineage=None, external_lineage=None)`.

# COMMAND ----------

# MAGIC %md
# MAGIC ## Install DQX with LLM extras + dbldatagen

# COMMAND ----------

%pip install 'databricks-labs-dqx[llm]' dbldatagen

%restart_python

# COMMAND ----------

dbutils.widgets.text("demo_catalog", "main", "Catalog Name")
dbutils.widgets.text("demo_schema", "default", "Schema Name")
dbutils.widgets.text("model_name", "databricks/claude-sonnet-5-5", "Model Name")

# COMMAND ----------

demo_catalog_name = dbutils.widgets.get("demo_catalog")
demo_schema_name = dbutils.widgets.get("demo_schema")
model_name = dbutils.widgets.get("model_name")

qualified = f"{demo_catalog_name}.{demo_schema_name}"
bronze_table = f"{qualified}.invoices_bronze"
silver_table = f"{qualified}.invoices_silver"
dim_customer_table = f"{qualified}.dim_customer"
dim_service_table = f"{qualified}.dim_service"
fact_invoice_table = f"{qualified}.fact_invoice"

print(
"Tables:\n"
f" bronze: {bronze_table}\n"
f" silver: {silver_table}\n"
f" dim_customer:{dim_customer_table}\n"
f" dim_service: {dim_service_table}\n"
f" fact_invoice:{fact_invoice_table}"
)

# COMMAND ----------

# MAGIC %md
# MAGIC ## Bronze — synthetic invoices via `dbldatagen`

# COMMAND ----------

import dbldatagen as dg
from datetime import datetime, timedelta

ROW_COUNT = 10_000

bronze_generator = (
dg.DataGenerator(sparkSession=spark, name="invoices_bronze", rowCount=ROW_COUNT, partitions=4)
.withColumn("invoice_id", "string", template=r"INV-\d{8}", random=True)
.withColumn(
"invoice_timestamp",
"timestamp",
begin=datetime.utcnow() - timedelta(days=30),
end=datetime.utcnow() + timedelta(days=2),
random=True,
)
.withColumn("client_id", "string", values=[f"C{n:04d}" for n in range(1, 250)] + [None], random=True)
.withColumn("client_name", "string", values=[" alice corp ", "Beta LLC", "gamma inc", None], random=True)
.withColumn("service_id", "string", values=[f"S{n:03d}" for n in range(1, 40)], random=True)
.withColumn("service_name", "string", values=[" widgets ", "Support", "consulting"], random=True)
.withColumn("amount", "integer", minValue=-50, maxValue=5000, random=True)
.withColumn("currency_code", "string", values=["USD", "EUR", "GBP", "CHF", "JPY", "ZZZ"], random=True)
.withColumn("signed", "boolean", random=True)
)

bronze_generator.build().write.mode("overwrite").format("delta").saveAsTable(bronze_table)

# COMMAND ----------

# MAGIC %md
# MAGIC ## Silver + gold layers

# COMMAND ----------

spark.sql(
f"""
CREATE OR REPLACE TABLE {silver_table} USING DELTA AS
SELECT DISTINCT invoice_id, invoice_timestamp, client_id,
INITCAP(TRIM(client_name)) AS client_name,
service_id, INITCAP(TRIM(service_name)) AS service_name,
amount, currency_code, signed
FROM {bronze_table}
WHERE invoice_id IS NOT NULL
"""
)
spark.sql(
f"""
CREATE OR REPLACE TABLE {dim_customer_table} USING DELTA AS
SELECT DISTINCT client_id, client_name FROM {silver_table} WHERE client_id IS NOT NULL
"""
)
spark.sql(
f"""
CREATE OR REPLACE TABLE {dim_service_table} USING DELTA AS
SELECT DISTINCT service_id, service_name FROM {silver_table} WHERE service_id IS NOT NULL
"""
)
spark.sql(
f"""
CREATE OR REPLACE TABLE {fact_invoice_table} USING DELTA AS
SELECT invoice_id, invoice_timestamp, client_id, service_id, amount, currency_code, signed
FROM {silver_table}
"""
)

# COMMAND ----------

# MAGIC %md
# MAGIC ## Baseline — generate rules without Unity Catalog metadata
# MAGIC
# MAGIC `unity_catalog_metadata_config=None` is the default. At this point the tables have no
# MAGIC comments, no tags, and no external lineage, so the LLM has only column names + types
# MAGIC to work with.

# COMMAND ----------

import yaml
from databricks.sdk import WorkspaceClient

from databricks.labs.dqx.config import InputConfig, LLMModelConfig, UnityCatalogMetadataConfig
from databricks.labs.dqx.engine import DQEngine
from databricks.labs.dqx.profiler.generator import DQGenerator

ws = WorkspaceClient()
engine = DQEngine(ws, spark)
generator = DQGenerator(ws, spark, llm_model_config=LLMModelConfig(model_name=model_name))

USER_INPUT = "generate quality checks for an invoice fact table"

baseline_checks = generator.generate_dq_rules_ai_assisted(
user_input=USER_INPUT,
input_config=InputConfig(location=fact_invoice_table),
unity_catalog_metadata_config=None,
)
print(yaml.safe_dump(baseline_checks, sort_keys=False))

# COMMAND ----------

# MAGIC %md
# MAGIC ## Attach concise UC comments on every layer

# COMMAND ----------

spark.sql(f"COMMENT ON TABLE {bronze_table} IS 'Raw SAP SD invoice events, pre-dedup'")
spark.sql(f"COMMENT ON TABLE {silver_table} IS 'Deduplicated invoices with normalised names'")
spark.sql(f"COMMENT ON TABLE {fact_invoice_table} IS 'One row per confirmed invoice line'")

for tbl in (bronze_table, silver_table, fact_invoice_table):
spark.sql(f"ALTER TABLE {tbl} ALTER COLUMN invoice_id COMMENT 'Invoice primary key'")
spark.sql(f"ALTER TABLE {tbl} ALTER COLUMN invoice_timestamp COMMENT 'Invoice issuance time, UTC'")
spark.sql(f"ALTER TABLE {tbl} ALTER COLUMN client_id COMMENT 'Client identifier (SAP KUNAG)'")
spark.sql(f"ALTER TABLE {tbl} ALTER COLUMN amount COMMENT 'Positive integer amount in minor units'")
spark.sql(f"ALTER TABLE {tbl} ALTER COLUMN currency_code COMMENT 'ISO-4217 three-letter currency'")

# COMMAND ----------

# MAGIC %md
# MAGIC ## Attach a handful of UC tags
# MAGIC
# MAGIC Tags land in `system.information_schema.table_tags` / `column_tags`, which the Spark
# MAGIC tag collector reads. Requires the appropriate UC privileges — the demo swallows the
# MAGIC error and continues so the rest of the flow still runs.

# COMMAND ----------

try:
spark.sql(f"ALTER TABLE {fact_invoice_table} SET TAGS ('domain' = 'finance')")
spark.sql(f"ALTER TABLE {fact_invoice_table} ALTER COLUMN client_id SET TAGS ('pii' = 'client')")
spark.sql(f"ALTER TABLE {fact_invoice_table} ALTER COLUMN amount SET TAGS ('financial' = 'true')")
spark.sql(f"ALTER TABLE {fact_invoice_table} ALTER COLUMN invoice_id SET TAGS ('source' = 'sap_sd')")
except Exception as exc:
print(f"Skipping tag DDL (requires UC privileges): {exc}")

# COMMAND ----------

# MAGIC %md
# MAGIC ## External lineage: register the SAP source
# MAGIC
# MAGIC Create an `ExternalMetadata` object representing the upstream SAP SD invoice table
# MAGIC (VBRK) and an `ExternalLineageRelationship` from it into the bronze UC table. This
# MAGIC enrichment is invisible to the recursive-CTE walker (external lineage is not in
# MAGIC `system.access.*_lineage`) and is only reachable via the SDK — see
# MAGIC https://docs.databricks.com/aws/en/data-governance/unity-catalog/external-lineage.

# COMMAND ----------

from databricks.sdk.service.catalog import (
ColumnRelationship,
CreateRequestExternalLineage,
ExternalLineageExternalMetadata,
ExternalLineageObject,
ExternalLineageTableInfo,
ExternalMetadata,
SystemType,
)

EXTERNAL_METADATA_NAME = "dqx_demo_sap_sd_invoices"
external_metadata_created = False
external_relationship_created = False

try:
ws.external_metadata.create_external_metadata(
ExternalMetadata(
name=EXTERNAL_METADATA_NAME,
system_type=SystemType.SAP,
entity_type="TABLE",
description="SAP SD invoice headers (table VBRK)",
)
)
external_metadata_created = True

bronze_cat, bronze_sch, bronze_name = bronze_table.split(".")
ws.external_lineage.create_external_lineage_relationship(
CreateRequestExternalLineage(
source=ExternalLineageObject(
external_metadata=ExternalLineageExternalMetadata(name=EXTERNAL_METADATA_NAME),
),
target=ExternalLineageObject(
table=ExternalLineageTableInfo(
catalog_name=bronze_cat, schema_name=bronze_sch, name=bronze_name
),
),
columns=[
ColumnRelationship(source="VBELN", target="invoice_id"),
ColumnRelationship(source="KUNAG", target="client_id"),
ColumnRelationship(source="NETWR", target="amount"),
ColumnRelationship(source="WAERK", target="currency_code"),
],
)
)
external_relationship_created = True
print("External lineage SAP → bronze created.")
except Exception as exc:
print(f"Skipping external lineage setup (API unavailable or missing permission): {exc}")

# COMMAND ----------

# MAGIC %md
# MAGIC ## Enriched — generate rules with Unity Catalog metadata
# MAGIC
# MAGIC `UnityCatalogMetadataConfig()` enables every enrichment with sensible defaults
# MAGIC (`column_upstream_lineage=ColumnUpstreamLineageConfig(depth=2, lookback_days=30,
# MAGIC max_nodes=50)` and `external_lineage=ExternalLineageConfig(max_relationships=50)`).

# COMMAND ----------

enriched_checks = generator.generate_dq_rules_ai_assisted(
user_input=USER_INPUT,
input_config=InputConfig(location=fact_invoice_table),
unity_catalog_metadata_config=UnityCatalogMetadataConfig(),
)
print(yaml.safe_dump(enriched_checks, sort_keys=False))

# COMMAND ----------

# MAGIC %md
# MAGIC ## Side-by-side comparison

# COMMAND ----------

print("--- baseline (no UC metadata) ---")
print(yaml.safe_dump(baseline_checks, sort_keys=False))
print("--- enriched (UnityCatalogMetadataConfig defaults) ---")
print(yaml.safe_dump(enriched_checks, sort_keys=False))

# COMMAND ----------

# MAGIC %md
# MAGIC ## Apply the enriched rules

# COMMAND ----------

valid_df, invalid_df = engine.apply_checks_by_metadata_and_split(
spark.table(fact_invoice_table), enriched_checks
)
print(f"valid rows: {valid_df.count()}")
print(f"invalid rows: {invalid_df.count()}")

# COMMAND ----------

# MAGIC %md
# MAGIC ## Cleanup — drop the external metadata registration

# COMMAND ----------

if external_relationship_created:
try:
from databricks.sdk.service.catalog import DeleteRequestExternalLineage

bronze_cat, bronze_sch, bronze_name = bronze_table.split(".")
ws.external_lineage.delete_external_lineage_relationship(
DeleteRequestExternalLineage(
source=ExternalLineageObject(
external_metadata=ExternalLineageExternalMetadata(name=EXTERNAL_METADATA_NAME),
),
target=ExternalLineageObject(
table=ExternalLineageTableInfo(
catalog_name=bronze_cat, schema_name=bronze_sch, name=bronze_name
),
),
)
)
except Exception as exc:
print(f"Could not delete external lineage relationship: {exc}")

if external_metadata_created:
try:
ws.external_metadata.delete_external_metadata(name=EXTERNAL_METADATA_NAME)
except Exception as exc:
print(f"Could not delete external metadata: {exc}")
1 change: 1 addition & 0 deletions docs/dqx/docs/demos.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ Import the following notebooks in the Databricks workspace to try DQX Core out:
* [DQX Demo Notebook for Profiling and Applying Checks at Scale on Multiple Tables](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_multi_table_demo.py) - demonstrates how to use DQX as a library at scale to apply checks on multiple tables.
* [DQX Demo Notebook for Row Anomaly Detection](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_row_anomaly_detection_demo.py) - comprehensive demo showing how to use DQX Row Anomaly Detection to detect unusual patterns in your data.
* [DQX Demo Notebook for AI-assisted checks generation](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_demo_ai_assisted_checks_generation.py) - demonstrates how to generate DQX rules/checks with LLM using natural language.
* [DQX Demo Notebook for AI-assisted checks generation with Unity Catalog metadata](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_demo_ai_assisted_unity_catalog_metadata.py) - builds a synthetic invoice pipeline, runs AI-assisted generation twice against the same table (once without `UnityCatalogMetadataConfig`, once with concise comments / tags / column upstream lineage / a SAP-origin external lineage relationship), and prints the two generated rule sets side-by-side.
* [DQX Demo Notebook for Data Contract Integration (ODCS)](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_demo_datacontract_odcs.py) - demonstrates how to generate DQX quality rules from ODCS (Open Data Contract Standard) data contracts, including predefined rules from schema constraints, explicit custom rules, and contract metadata tracking.
* [DQX Demo Notebook for Spark Structured Streaming (Native End-to-End Approach)](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_streaming_demo_native.py) - demonstrates how to use DQX as a library with Spark Structured Streaming, using the built-in end-to-end method to handle both reading and writing.
* [DQX Demo Notebook for Spark Structured Streaming (DIY Approach)](https://github.com/databrickslabs/dqx/blob/v0.16.0/demos/dqx_streaming_demo_diy.py) - demonstrates how to use DQX as a library with Spark Structured Streaming, while handling reading and writing on your own outside DQX using Spark API.
Expand Down
14 changes: 14 additions & 0 deletions docs/dqx/docs/guide/ai_assisted_quality_checks_generation.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,20 @@ The AI-assisted quality checks generation is available via the `DQGenerator` cla
This approach is particularly useful when you have clear business requirements, but need help translating them into technical data quality rules.
It's also useful for generating technical rules without requiring knowledge of DQX-specific syntax.

<Admonition type="tip" title="Enrich the prompt with Unity Catalog metadata (experimental)">
Pass a `UnityCatalogMetadataConfig` to `generate_dq_rules_ai_assisted` to attach UC table and
column comments, column tags from `system.information_schema.*_tags`, bounded column upstream
lineage from `system.access.column_lineage`, and external upstream lineage (non-Databricks
sources such as SAP, Salesforce, Tableau) read via the [Unity Catalog External Lineage API](https://docs.databricks.com/aws/en/data-governance/unity-catalog/external-lineage).

Instantiating `UnityCatalogMetadataConfig()` enables every enrichment path with sensible
defaults. Disable just one walk by setting its sub-model to `None` — for example
`UnityCatalogMetadataConfig(column_upstream_lineage=None)` keeps comments / tags / external
lineage but skips the recursive column-lineage CTE. See
`demos/dqx_demo_ai_assisted_unity_catalog_metadata.py` for a side-by-side comparison that
generates rules for the same table with and without the enrichment.
</Admonition>

<Admonition type="tip" title="When to use AI-Assisted generation">
AI-assisted generation is ideal for:
- Translating business requirements into technical quality rules.
Expand Down
Loading