Skip to content
Merged
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
81 changes: 79 additions & 2 deletions demo/SDP_META_INTERACTIVE_DEMO.py
Original file line number Diff line number Diff line change
Expand Up @@ -1601,6 +1601,33 @@ def _load_sample_onboarding_text(onboarding_format, conf_ext, git_branch):

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

# MAGIC %md
# MAGIC ### 2.5 DQE Without Quarantine Metadata
# MAGIC
# MAGIC The customers and transactions Silver DQE files contain drop rules but
# MAGIC no `expect_or_quarantine` rules. Their onboarding entries therefore do
# MAGIC not need quarantine table metadata. SDP-META retains the DQE rules and
# MAGIC leaves `quarantineTargetDetails` empty.

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

display(
spark.sql(
f"""
SELECT
dataFlowId,
targetDetails['table'] AS silver_table,
quarantineTargetDetails,
dataQualityExpectations
FROM {uc_catalog_name}.{uc_schema_name}.silver_dataflowspec
WHERE dataFlowId IN ('100', '101')
ORDER BY dataFlowId
"""
)
)

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

# MAGIC %md
# MAGIC ---
# MAGIC ## Stage 3: Create Lakeflow Spark Declarative Pipeline
Expand Down Expand Up @@ -4122,7 +4149,57 @@ def _expect_nonempty(fqn):
f".{domain}_quarantine"
)

# 5. Customers / transactions / products / stores — count varies
# 5. Silver DQE without quarantine metadata — both flows contain
# drop rules but no quarantine rules or targets. Onboarding must
# retain the DQE without requiring quarantine metadata.
try:
optional_rows = {
row.dataFlowId: row
for row in spark.sql(
f"""
SELECT dataFlowId, dataQualityExpectations,
quarantineTargetDetails
FROM {uc_catalog_name}.{uc_schema_name}.silver_dataflowspec
WHERE dataFlowId IN ('100', '101')
"""
).collect()
}
if set(optional_rows) != {"100", "101"}:
failures.append(
"Silver DQE without quarantine: expected DataflowSpecs "
"100 and 101"
)
else:
customer_dqe = json.loads(
optional_rows["100"].dataQualityExpectations
)
transaction_dqe = json.loads(
optional_rows["101"].dataQualityExpectations
)
if customer_dqe.get("expect_or_quarantine"):
failures.append(
"customers Silver DQE unexpectedly enables quarantine"
)
if transaction_dqe.get("expect_or_quarantine"):
failures.append(
"transactions Silver DQE unexpectedly enables quarantine"
)
if optional_rows["100"].quarantineTargetDetails:
failures.append(
"customers Silver quarantine metadata "
"should remain empty"
)
if optional_rows["101"].quarantineTargetDetails:
failures.append(
"transactions Silver quarantine metadata "
"should remain empty"
)
except Exception as exc:
failures.append(
f"Silver DQE without quarantine validation failed: {exc}"
)

# 6. Customers / transactions / products / stores — count varies
# with ``data_source``: ``github`` uses fixed CSVs from the
# repo, ``dbdatagen`` uses random synthetic data with no fixed
# seed. Existence + non-empty is the strongest universal check;
Expand All @@ -4137,7 +4214,7 @@ def _expect_nonempty(fqn):
f"{uc_catalog_name}.{silver_schema}.{domain}"
)

# 6. Multi-source AUTO CDC (Stage 11) — every region seeds the
# 7. Multi-source AUTO CDC (Stage 11) — every region seeds the
# SAME shape: 3 INSERTs + 1 UPDATE + 1 DELETE = 5 raw bronze
# rows. The silver target is SCD-1 with apply_as_deletes, so the
# final live row count = (3 regions × 3 inserted) − (3 regions ×
Expand Down
7 changes: 5 additions & 2 deletions docs/docs/concepts/dataflowspec.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,8 @@ These fields are required on every flow entry.
| `bronze_cluster_by_auto` | boolean | No | Enable auto liquid clustering. |
| `bronze_data_quality_expectations_json_<env>` | string | No | Path to a DQE file (JSON or YAML). |
| `bronze_catalog_quarantine_<env>` | string | No | Catalog for the quarantine table (defaults to bronze catalog). |
| `bronze_database_quarantine_<env>` | string | No | Schema for the quarantine table. Required when DQE has `drop` expectations. |
| `bronze_quarantine_table` | string | No | Quarantine table name. |
| `bronze_database_quarantine_<env>` | string | No | Schema for the quarantine table. Needed to create an output for non-empty `expect_or_quarantine` rules. |
| `bronze_quarantine_table` | string | No | Quarantine table name. Needed to create an output for non-empty `expect_or_quarantine` rules. |
| `bronze_quarantine_table_comment` | string | No | Quarantine table comment. |
| `bronze_quarantine_table_path_<env>` | string | No | External path for the quarantine table (non-UC). |
| `bronze_quarantine_table_cluster_by` | list | No | Liquid clustering columns for the quarantine table. |
Expand All @@ -68,6 +68,9 @@ These fields are required on every flow entry.
| `silver_table_path_<env>` | string | Non-UC | External path for the Silver table (required without Unity Catalog). |
| `silver_transformation_json_<env>` | string | Yes* | Path to a transformations file defining `select_exp` and `where_clause`. *Not required when using `silver_cdc_apply_changes_flows`. |
| `silver_data_quality_expectations_json_<env>` | string | No | Path to a DQE file. |
| `silver_database_quarantine_<env>` | string | No | Schema for the Silver quarantine table. Needed to create an output for non-empty `expect_or_quarantine` rules. |
| `silver_quarantine_table` | string | No | Silver quarantine table name. Needed to create an output for non-empty `expect_or_quarantine` rules. |
| `silver_quarantine_table_path_<env>` | string | No | External path for the Silver quarantine table (non-UC). |
| `silver_cdc_apply_changes` | map | No | Single-source CDC config. See [CDC](../guides/cdc.md). |
| `silver_cdc_apply_changes_flows` | map | No | Multi-source CDC flow group. See [Multi-source CDC](../guides/multi-source-cdc.md). |
| `silver_reader_options` | map | No | Additional Spark reader options for the Silver source. |
Expand Down
27 changes: 25 additions & 2 deletions docs/docs/reference/dq-rules.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ Data quality rules are defined in a separate JSON or YAML file and referenced fr
| `expect_or_fail` | Halt the entire pipeline update | A violated rule indicates a critical upstream data problem |

:::tip
Prefer `expect_or_quarantine` over `expect_or_drop` when you want to inspect failed rows later. The quarantine table has the same schema as the target table plus an `_error` column.
Prefer `expect_or_quarantine` over `expect_or_drop` when you want to inspect failed rows later.
:::

:::warning
Expand Down Expand Up @@ -84,7 +84,30 @@ For silver:

## Quarantine behavior

When `expect_or_drop` rules are configured and a quarantine table is defined (`bronze_quarantine_table`, `bronze_database_quarantine_{env}`), rows that fail are written to the quarantine table rather than discarded. The quarantine table has the same schema as the main bronze table plus a `_error` column.
Quarantine target fields remain optional during onboarding for backward
compatibility. A DQE file containing only `expect`, `expect_or_drop`, or
`expect_or_fail` does not use any quarantine fields.

To create an output for Bronze quarantine rules, configure
`bronze_quarantine_table` and
`bronze_database_quarantine_{env}`. For Silver, configure
`silver_quarantine_table` and `silver_database_quarantine_{env}`. These fields
identify the quarantine target persisted in the DataflowSpec. Non-Unity
Catalog targets also use the corresponding
`bronze_quarantine_table_path_{env}` or
`silver_quarantine_table_path_{env}`.

Legacy onboarding files with non-empty `expect_or_quarantine` rules but no
target continue to onboard successfully. They do not produce a quarantine
table until target metadata is supplied, and onboarding emits a warning. If
`expect_or_quarantine` is the only non-empty constraint block, the pipeline
also does not declare the main output table. Add the quarantine table and
database fields to avoid a pipeline with no declared output. Missing targets
remain accepted only for backward compatibility.

Existing onboarding files may include optional quarantine metadata before
quarantine rules are added. SDP-META preserves that metadata, but no
quarantine output is created until `expect_or_quarantine` is non-empty.

:::tip
Use the quarantine table to inspect and reprocess failed rows.
Expand Down
54 changes: 54 additions & 0 deletions integration_tests/conf/json/cloudfiles-onboarding.template
Original file line number Diff line number Diff line change
Expand Up @@ -192,5 +192,59 @@
"silver_quarantine_table":"transactions_quarantine",
"silver_quarantine_table_cluster_by":["id","customer_id"],
"silver_quarantine_table_cluster_by_auto": true
},
{
"data_flow_id": "190",
"data_flow_group": "QUARANTINE_COMPAT",
"source_system": "MYSQL",
"source_format": "cloudFiles",
"source_details": {
"source_database": "APP",
"source_table": "CUSTOMERS",
"source_path_it": "{uc_volume_path}/integration_tests/resources/data/customers",
"source_schema_path": "{uc_volume_path}/integration_tests/resources/customers.ddl"
},
"bronze_catalog_it": "{uc_catalog_name}",
"bronze_database_it": "{bronze_schema}",
"bronze_table": "quarantine_compat_customers",
"bronze_reader_options": {
"cloudFiles.format": "json",
"cloudFiles.inferColumnTypes": "true",
"cloudFiles.rescuedDataColumn": "_rescued_data"
},
"bronze_table_path_it": "{uc_volume_path}/data/bronze/quarantine_compat_customers",
"bronze_data_quality_expectations_json_it": "{uc_volume_path}/integration_tests/conf/json/dqe/customers/bronze_data_quality_expectations.json",
"silver_catalog_it": "{uc_catalog_name}",
"silver_database_it": "{silver_schema}",
"silver_table": "customers",
"silver_transformation_json_it": "{uc_volume_path}/integration_tests/conf/json/silver_transformations.json",
"silver_data_quality_expectations_json_it": "{uc_volume_path}/integration_tests/conf/json/dqe/customers/silver_data_quality_expectations.json"
},
{
"data_flow_id": "191",
"data_flow_group": "DQE_NO_QUARANTINE",
"source_system": "MYSQL",
"source_format": "cloudFiles",
"source_details": {
"source_database": "APP",
"source_table": "CUSTOMERS",
"source_path_it": "{uc_volume_path}/integration_tests/resources/data/customers",
"source_schema_path": "{uc_volume_path}/integration_tests/resources/customers.ddl"
},
"bronze_catalog_it": "{uc_catalog_name}",
"bronze_database_it": "{bronze_schema}",
"bronze_table": "dqe_no_quarantine_customers",
"bronze_reader_options": {
"cloudFiles.format": "json",
"cloudFiles.inferColumnTypes": "true",
"cloudFiles.rescuedDataColumn": "_rescued_data"
},
"bronze_table_path_it": "{uc_volume_path}/data/bronze/dqe_no_quarantine_customers",
"bronze_data_quality_expectations_json_it": "{uc_volume_path}/integration_tests/conf/json/dqe/customers/no_quarantine_rules.json",
"silver_catalog_it": "{uc_catalog_name}",
"silver_database_it": "{silver_schema}",
"silver_table": "customers",
"silver_transformation_json_it": "{uc_volume_path}/integration_tests/conf/json/silver_transformations.json",
"silver_data_quality_expectations_json_it": "{uc_volume_path}/integration_tests/conf/json/dqe/customers/no_quarantine_rules.json"
}
]
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
{
"expect_or_drop": {
"valid_id": "id IS NOT NULL",
"valid_operation": "operation IN ('APPEND', 'DELETE', 'UPDATE')"
}
}
46 changes: 46 additions & 0 deletions integration_tests/conf/yml/cloudfiles-onboarding.template.yml
Original file line number Diff line number Diff line change
Expand Up @@ -157,3 +157,49 @@
- id
- customer_id
silver_quarantine_table_cluster_by_auto: true
- data_flow_id: '190'
data_flow_group: QUARANTINE_COMPAT
source_system: MYSQL
source_format: cloudFiles
source_details:
source_database: APP
source_table: CUSTOMERS
source_path_it: '{uc_volume_path}/integration_tests/resources/data/customers'
source_schema_path: '{uc_volume_path}/integration_tests/resources/customers.ddl'
bronze_catalog_it: '{uc_catalog_name}'
bronze_database_it: '{bronze_schema}'
bronze_table: quarantine_compat_customers
bronze_reader_options:
cloudFiles.format: json
cloudFiles.inferColumnTypes: 'true'
cloudFiles.rescuedDataColumn: _rescued_data
bronze_table_path_it: '{uc_volume_path}/data/bronze/quarantine_compat_customers'
bronze_data_quality_expectations_json_it: '{uc_volume_path}/integration_tests/conf/yml/dqe/customers/bronze_data_quality_expectations.yml'
silver_catalog_it: '{uc_catalog_name}'
silver_database_it: '{silver_schema}'
silver_table: customers
silver_transformation_json_it: '{uc_volume_path}/integration_tests/conf/yml/silver_transformations.yml'
silver_data_quality_expectations_json_it: '{uc_volume_path}/integration_tests/conf/yml/dqe/customers/silver_data_quality_expectations.yml'
- data_flow_id: '191'
data_flow_group: DQE_NO_QUARANTINE
source_system: MYSQL
source_format: cloudFiles
source_details:
source_database: APP
source_table: CUSTOMERS
source_path_it: '{uc_volume_path}/integration_tests/resources/data/customers'
source_schema_path: '{uc_volume_path}/integration_tests/resources/customers.ddl'
bronze_catalog_it: '{uc_catalog_name}'
bronze_database_it: '{bronze_schema}'
bronze_table: dqe_no_quarantine_customers
bronze_reader_options:
cloudFiles.format: json
cloudFiles.inferColumnTypes: 'true'
cloudFiles.rescuedDataColumn: _rescued_data
bronze_table_path_it: '{uc_volume_path}/data/bronze/dqe_no_quarantine_customers'
bronze_data_quality_expectations_json_it: '{uc_volume_path}/integration_tests/conf/yml/dqe/customers/no_quarantine_rules.yml'
silver_catalog_it: '{uc_catalog_name}'
silver_database_it: '{silver_schema}'
silver_table: customers
silver_transformation_json_it: '{uc_volume_path}/integration_tests/conf/yml/silver_transformations.yml'
silver_data_quality_expectations_json_it: '{uc_volume_path}/integration_tests/conf/yml/dqe/customers/no_quarantine_rules.yml'
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
expect_or_drop:
valid_id: id IS NOT NULL
valid_operation: operation IN ('APPEND', 'DELETE', 'UPDATE')
72 changes: 72 additions & 0 deletions integration_tests/notebooks/cloudfile_runners/validate.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,11 @@
# Databricks notebook source
import pandas as pd
import json

run_id = dbutils.widgets.get("run_id")
uc_enabled = dbutils.widgets.get("uc_enabled").strip().lower() == "true"
uc_catalog_name = dbutils.widgets.get("uc_catalog_name")
sdp_meta_schema = dbutils.widgets.get("sdp_meta_schema")
bronze_schema = dbutils.widgets.get("bronze_schema")
silver_schema = dbutils.widgets.get("silver_schema")
output_file_path = dbutils.widgets.get("output_file_path")
Expand Down Expand Up @@ -44,6 +46,76 @@
except AssertionError:
log_list.append(f"Expected: {counts} Actual: {cnt}. Failed!")

# Backward-compatibility coverage for optional quarantine targets. Flow 190
# carries non-empty quarantine rules but intentionally omits every quarantine
# target field. Onboarding must retain the DQE and persist an empty target for
# both layers rather than failing an existing customer configuration.
log_list.append(
"Validating legacy quarantine rules without target metadata."
)
for layer in ("bronze", "silver"):
spec_table = (
f"{uc_catalog_name}.{sdp_meta_schema}."
f"{layer}_dataflowspec_cdc"
)
rows = spark.sql(
f"""
SELECT dataQualityExpectations, quarantineTargetDetails
FROM {spec_table}
WHERE dataFlowId = '190'
AND dataFlowGroup = 'QUARANTINE_COMPAT'
"""
).collect()
try:
assert len(rows) == 1
dqe = json.loads(rows[0].dataQualityExpectations)
assert dqe.get("expect_or_quarantine")
assert not rows[0].quarantineTargetDetails
log_list.append(
f"{layer.title()} compatibility DataflowSpec retained DQE "
"with an empty quarantine target. Passed!"
)
except (AssertionError, TypeError, ValueError) as exc:
log_list.append(
f"{layer.title()} compatibility DataflowSpec validation "
f"failed: {exc}. Failed!"
)

# Exact optional-metadata scenario: flow 191 has DQE rules, but no
# expect_or_quarantine block and no quarantine fields. Both layers must retain
# the drop rules without synthesizing or requiring a quarantine target.
log_list.append(
"Validating DQE without quarantine rules or target metadata."
)
for layer in ("bronze", "silver"):
spec_table = (
f"{uc_catalog_name}.{sdp_meta_schema}."
f"{layer}_dataflowspec_cdc"
)
rows = spark.sql(
f"""
SELECT dataQualityExpectations, quarantineTargetDetails
FROM {spec_table}
WHERE dataFlowId = '191'
AND dataFlowGroup = 'DQE_NO_QUARANTINE'
"""
).collect()
try:
assert len(rows) == 1
dqe = json.loads(rows[0].dataQualityExpectations)
assert dqe.get("expect_or_drop")
assert not dqe.get("expect_or_quarantine")
assert not rows[0].quarantineTargetDetails
log_list.append(
f"{layer.title()} DQE-only DataflowSpec retained drop "
"rules with an empty quarantine target. Passed!"
)
except (AssertionError, TypeError, ValueError) as exc:
log_list.append(
f"{layer.title()} DQE-only DataflowSpec validation "
f"failed: {exc}. Failed!"
)

# Regression coverage for issue #444: source_metadata nested under a bronze
# append flow must be serialized like top-level source metadata. Every
# transaction row comes from either `transactions/` or `transactions_af/`, so
Expand Down
1 change: 1 addition & 0 deletions integration_tests/run_integration_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -558,6 +558,7 @@ def create_workflow_spec(self, runner_conf: SDPMetaRunnerConf):
base_parameters={
"uc_enabled": "True",
"uc_catalog_name": f"{runner_conf.uc_catalog_name}",
"sdp_meta_schema": f"{runner_conf.sdp_meta_schema}",
"bronze_schema": f"{runner_conf.bronze_schema}",
"silver_schema": (
f"{runner_conf.silver_schema}"
Expand Down
Loading
Loading