Skip to main content

dlt-expectations-patterns

Spark Declarative Pipeline (SDP, formerly DLT) expectations patterns for data quality with Unity Catalog Delta table storage. Use when implementing Silver layer SDP/DLT pipelines, creating portable data quality rules, or needing runtime-updateable expectations without code deployment. Supports severity-based filtering (critical vs warning) and quarantine patterns. Standardizes on `import dlt` for the DQ-rules framework; the modern `dp` API (`from pyspark import pipelines as dp`) also supports expectations (`dp.expect_all_or_drop`) and is Databricks' recommended forward path.

跳到安装

来源信息

仓库
databricks-solutions/vibe-coding-workshop-template
最近来源活动
2026年8月31日 04:03
检测到的 SKILL.md 语言
英语
星标
6
分支
7

安装方式

默认使用会先检查来源的 Prompt;你也可以切换为直接命令,或下载本地副本。

检查来源文件

决定是否安装前,请先阅读 SKILL.md,以及 SkillsMP 当前展示的配套文件。

文件资源管理器
4 个文件

正在显示 SKILL.md

SKILL.md
来源说明 · 只读预览
name
dlt-expectations-patterns
description
Spark Declarative Pipeline (SDP, formerly DLT) expectations patterns for data quality with Unity Catalog Delta table storage. Use when implementing Silver layer SDP/DLT pipelines, creating portable data quality rules, or needing runtime-updateable expectations without code deployment. Supports severity-based filtering (critical vs warning) and quarantine patterns. Standardizes on `import dlt` for the DQ-rules framework; the modern `dp` API (`from pyspark import pipelines as dp`) also supports expectations (`dp.expect_all_or_drop`) and is Databricks' recommended forward path.
clients
["ide_cli","genie_code"]
bundle_resource
pipelines
deploy_verb
bundle_deploy
deploy_note
DLT expectations run inside the Silver pipeline; deploy via `bundle deploy --target dev` (runDatabricksCli on Genie Code).
coverage
full
metadata
{"author":"prashanth subrahmanyam","version":"1.0","domain":"silver","role":"worker","pipeline_stage":3,"pipeline_stage_name":"silver","called_by":["silver-layer-setup"],"standalone":true,"last_verified":"2026-08-30","volatility":"medium","upstream_sources":[{"name":"databricks-agent-skills","repo":"databricks/databricks-agent-skills","paths":"[Truncated]","relationship":"extended","last_synced":"2026-08-30","sync_commit":"ca92a6c"}]}
# SDP/DLT Expectations Patterns > **Naming:** Databricks rebranded DLT to **Spark Declarative Pipelines (SDP)**. The modern Python API is `from pyspark import pipelines as dp` with `@dp.table()` decorators, and **Databricks recommends `dp`**. Expectations decorators **are available in `dp`** — `@dp.expect_all_or_drop()` and `@dp.expect_all()` mirror the `@dlt.*` versions (see the [Lakeflow pipelines Python reference](https://docs.databricks.com/aws/en/ldp/developer/python-ref)). This skill standardizes on `import dlt` so the rules-loader and decorators stay uniform across the DQ-rules framework; migrating to `dp` is a mechanical swap (`import dlt` → `from pyspark import pipelines as dp`, `@dlt.` → `@dp.`). ## Overview All Silver layer tables use data quality expectations loaded from a **Unity Catalog Delta table**. This skill standardizes the Delta table-based approach for portable, maintainable, and runtime-updateable data quality management. **Key Patterns:** 1. **Delta Table for Rules Storage** - Single source of truth in Unity Catalog 2. **Rules Loader Module** - Pure Python functions to load rules at runtime 3. **`@dlt.expect_all_or_drop()` Decorator** - Strict enforcement pattern 4. **Direct Publishing Mode** - Fully qualified table names with `get_source_table()` helper 5. **Severity-Based Filtering** - Critical vs warning rules **Official Reference:** [Portable and Reusable Expectations](https://docs.databricks.com/aws/en/ldp/expectation-patterns#portable-and-reusable-expectations) ## When to Use This Skill Use this skill when: - Implementing Silver layer DLT pipelines with data quality expectations - Creating portable data quality rules that can be shared across pipelines - Needing runtime-updateable expectations without code deployment - Requiring severity-based filtering (critical vs warning rules) - Implementing quarantine patterns for failed validations ## Benefits of Delta Table-Based Rules ### Why Delta Table Instead of Hardcoded Rules? | Aspect | Hardcoded Rules | Delta Table Rules | |---|---|---| | **Updateability** | Requires code changes + redeployment | UPDATE table, rules apply immediately | | **Auditability** | Git history only | Delta time travel + Git history | | **Portability** | Copied across environments | Shared table across pipelines | | **Documentation** | In code comments | Queryable with SQL | | **Maintenance** | Edit multiple notebooks | Single table UPDATE | | **Governance** | No access control | Unity Catalog permissions | **Recommended by Databricks:** "Store expectation definitions separately from pipeline logic to easily apply expectations to multiple datasets or pipelines. Update, audit, and maintain expectations without modifying pipeline source code." ## Quick Reference ### DLT Direct Publishing Mode (Modern Pattern) **DEPRECATED Patterns (Do NOT use):** - ❌ `LIVE.` prefix for table references (e.g., `LIVE.bronze_transactions`) - ❌ `target:` field in DLT pipeline configuration **MODERN Pattern (Always use):** - ✅ **Fully qualified table names**: `{catalog}.{schema}.{table_name}` - ✅ **`schema:` field** in DLT pipeline configuration (not `target`) - ✅ **Helper function** to build table names from configuration ### Helper Function Pattern ```python def get_source_table(table_name, source_schema_key="bronze_schema"): """Get fully qualified table name from DLT configuration.""" spark = SparkSession.getActiveSession() catalog = spark.conf.get("catalog") schema = spark.conf.get(source_schema_key) return f"{catalog}.{schema}.{table_name}" # Use in DLT table @dlt.table(...) def silver_transactions(): return dlt.read_stream(get_source_table("bronze_transactions")) ``` ### DLT Pipeline Configuration ```yaml resources: pipelines: silver_dlt_pipeline: name: "[${bundle.target} ${var.user_prefix}] Silver Layer Pipeline" # ✅ CORRECT: Use 'schema' (Direct Publishing Mode) catalog: ${var.catalog} schema: ${var.silver_schema} # ❌ WRONG: Don't use 'target' (deprecated) # target: ${var.catalog}.${var.silver_schema} configuration: catalog: ${var.catalog} bronze_schema: ${var.bronze_schema} silver_schema: ${var.silver_schema} serverless: true edition: ADVANCED ``` ## Critical Rules ### 1. Rules Loader Module (Pure Python, NO Notebook Header) **⚠️ CRITICAL: Pure Python file (NO `# Databricks notebook source` header)** ```python # File: dq_rules_loader.py (NO notebook header!) from pyspark.sql import SparkSession # Module-level cache for rules (loaded once at import time) _rules_cache = {} _cache_initialized = False def _load_all_rules() -> None: """Load all rules from Delta table into module-level cache.""" global _rules_cache, _cache_initialized if _cache_initialized: return spark = SparkSession.getActiveSession() if spark is None: return try: rules_table = f"{catalog}.{schema}.dq_rules" # ✅ Use toPandas() instead of .collect() to avoid DLT warning pdf = spark.sql(f"SELECT * FROM {rules_table}").toPandas() # Populate cache for _, row in pdf.iterrows(): cache_key = (row['table_name'], row['severity']) if cache_key not in _rules_cache: _rules_cache[cache_key] = {} _rules_cache[cache_key][row['rule_name']] = row['constraint_sql'] _cache_initialized = True except Exception as e: print(f"Note: Could not load DQ rules: {e}") def get_critical_rules_for_table(table_name: str) -> dict: """Get critical DQ rules from cache (no Spark operations).""" if not _cache_initialized: _load_all_rules() return _rules_cache.get((table_name, "critical"), {}) def get_warning_rules_for_table(table_name: str) -> dict: """Get warning DQ rules from cache (no Spark operations).""" if not _cache_initialized: _load_all_rules() return _rules_cache.get((table_name, "warning"), {}) ``` **See:** `references/expectation-patterns.md` for complete loader implementation ### 2. DLT Table Decorators **CRITICAL Rules (use `@dlt.expect_all_or_drop()`):** - Primary key fields (must be present and non-empty) - Foreign key fields (must be present for referential integrity) - Required date fields (must be present and >= minimum valid date) - Non-nullable business fields (quantity != 0, price > 0, etc.) **WARNING Rules (use `@dlt.expect_all()`):** - Reasonableness checks (quantity between 1 and 10000) - Recency checks (date within last 90 days) - Format preferences (UPC length between 12 and 14) - Coordinate ranges (latitude/longitude within valid bounds) ```python from dq_rules_loader import ( get_critical_rules_for_table, get_warning_rules_for_table ) @dlt.table(...) @dlt.expect_all_or_drop(get_critical_rules_for_table("silver_transactions")) @dlt.expect_all(get_warning_rules_for_table("silver_transactions")) def silver_transactions(): return dlt.read_stream(get_source_table("bronze_transactions")) ``` > **`dp` equivalent (recommended forward path):** the modern API supports the same expectations. Swap the import and prefix — `@dp.table(...)`, `@dp.expect_all_or_drop(...)`, `@dp.expect_all(...)`, and `spark.readStream.table(get_source_table(...))` instead of `dlt.read_stream(...)`. The rules loader is unchanged. See the [Lakeflow pipelines Python reference](https://docs.databricks.com/aws/en/ldp/developer/python-ref). ### 3. Avoiding `DataFrame.collect()` Warning **Problem:** DLT shows warning when `.collect()` is used in rules loader **Solution:** Use module-level cache with `toPandas()` ```python # ❌ WRONG: Direct .collect() shows warning def get_rules(table_name: str, severity: str) -> dict: df = spark.read.table(rules_table).filter(...).collect() # Warning! return {row['rule_name']: row['constraint_sql'] for row in df} # ✅ CORRECT: Use toPandas() with module-level cache _rules_cache = {} _cache_initialized = False def _load_all_rules(): pdf = spark.sql(f"SELECT * FROM {rules_table}").toPandas() # No warning! # Populate cache... def get_rules(table_name: str, severity: str) -> dict: if not _cache_initialized: _load_all_rules() return _rules_cache.get((table_name, severity), {}) # From cache! ``` > **Why this stays valid on `dp`/SDP:** the Lakeflow pipelines Python reference lists `toPandas()` and `.collect()` among operations to avoid **inside dataset definitions** (the functions decorated with `@dlt.table`/`@dp.table`). This loader calls `toPandas()` at **module import time** (decorator-evaluation), *outside* any dataset function, and caches the result — so it runs once and never during streaming execution. That's the intended pattern for both APIs. **See:** `references/expectation-patterns.md` for complete pattern ## Core Patterns ### Sourcing Column Names AND Enum Values from Data (Extract, Don't Generate) **First pin the column names, then the values.** Before authoring ANY rule, run `DESCRIBE TABLE` on each Bronze table and keep the column list in memory — every column referenced in a `constraint_sql` MUST exist in that `DESCRIBE` output. A rule that names a column absent from the live schema is a **hard error**, not a near-miss to "fix later" (the live failure: rules written for `price`/`latitude` when the schema had `base_price`/`property_latitude`). PRD/CSV names describe *intent*; only `DESCRIBE` gives the real column names. ```sql DESCRIBE TABLE {catalog}.{bronze_schema}.{bronze_table}; -- pin column names FIRST ``` Then, before authoring `col IN (...)` rules, extract the actual values from Bronze rather than reading them from schema CSV comments: ```sql SELECT DISTINCT col_name FROM {catalog}.{bronze_schema}.{bronze_table} WHERE col_name IS NOT NULL; ``` CSV column comments describe *intent*; production data may include extra values, typos, or legacy states. This follows the "Extract, Don't Generate" principle from `skills/databricks-expert-agent` — applied to both constraint **column names** and value literals. ### Pattern 1: Create DQ Rules Delta Table ```sql CREATE OR REPLACE TABLE {catalog}.{schema}.dq_rules ( table_name STRING NOT NULL, rule_name STRING NOT NULL, constraint_sql STRING NOT NULL, severity STRING NOT NULL, description STRING, created_timestamp TIMESTAMP NOT NULL, updated_timestamp TIMESTAMP NOT NULL, CONSTRAINT pk_dq_rules PRIMARY KEY (table_name, rule_name) NOT ENFORCED ) USING DELTA CLUSTER BY AUTO ``` 🔴 **Author this `dq_rules` table EXACTLY as shown.** Do NOT add columns the template does not list (no `is_active`, no status flags) and do NOT add a `DEFAULT` clause to any column. A `DEFAULT <expr>` (e.g. `is_active BOOLEAN NOT NULL DEFAULT true`) requires the `delta.feature.allowColumnDefaults` table feature — it is OFF by default and the DDL fails. If you genuinely need an active flag, declare the column without `DEFAULT` and set its value in the INSERT, never in the DDL (see `common/unity-catalog-constraints` → "Never Use `DEFAULT` Column Clauses in DDL"). This is a real regression: an invented `is_active ... DEFAULT true` column failed the DQ-setup job. **See:** `references/expectation-patterns.md` for complete schema and population examples ### Pattern 2: Apply Rules in DLT Tables ```python import dlt from dq_rules_loader import ( get_critical_rules_for_table, get_warning_rules_for_table, get_quarantine_condition ) @dlt.table( name="silver_transactions", table_properties={ "quality": "silver", "delta.enableChangeDataFeed": "true", "layer": "silver" }, cluster_by_auto=True ) @dlt.expect_all_or_drop(get_critical_rules_for_table("silver_transactions")) @dlt.expect_all(get_warning_rules_for_table("silver_transactions")) def silver_transactions(): """ Data Quality Rules (loaded from dq_rules Delta table): CRITICAL (Record DROPPED/QUARANTINED if fails): - Transaction ID, store number, UPC must be present - Quantity cannot be zero, price must be positive WARNING (Logged but record passes): - Quantity within reasonable range (-20 to 50) - Price within reasonable range ($0.01 to $500) """ return ( dlt.read_stream(get_source_table("bronze_transactions")) .withColumn("processed_timestamp", current_timestamp()) ) ``` **See:** `references/expectation-patterns.md` for complete DLT table examples ### Pattern 3: Quarantine Table ```python @dlt.table( name="silver_transactions_quarantine", comment="Quarantine table for records that failed CRITICAL data quality checks", table_properties={ "quality": "quarantine", "layer": "silver" }, cluster_by_auto=True ) def silver_transactions_quarantine(): """Quarantine failed records with rich diagnostic information.""" from dq_rules_loader import get_quarantine_condition return ( dlt.read_stream(get_source_table("bronze_transactions")) .filter(get_quarantine_condition("silver_transactions"))
在 GitHub 查看
这个 SKILL.md 很大,SkillsMP 这里只预览前一段内容。 在 GitHub 查看