diff --git a/demo/SDP_META_INTERACTIVE_DEMO.py b/demo/SDP_META_INTERACTIVE_DEMO.py index 4cbcc048..e65c1445 100644 --- a/demo/SDP_META_INTERACTIVE_DEMO.py +++ b/demo/SDP_META_INTERACTIVE_DEMO.py @@ -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 @@ -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; @@ -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 × diff --git a/docs/docs/concepts/dataflowspec.md b/docs/docs/concepts/dataflowspec.md index b3dd9eb4..0873b13e 100644 --- a/docs/docs/concepts/dataflowspec.md +++ b/docs/docs/concepts/dataflowspec.md @@ -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_` | string | No | Path to a DQE file (JSON or YAML). | | `bronze_catalog_quarantine_` | string | No | Catalog for the quarantine table (defaults to bronze catalog). | -| `bronze_database_quarantine_` | string | No | Schema for the quarantine table. Required when DQE has `drop` expectations. | -| `bronze_quarantine_table` | string | No | Quarantine table name. | +| `bronze_database_quarantine_` | 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_` | string | No | External path for the quarantine table (non-UC). | | `bronze_quarantine_table_cluster_by` | list | No | Liquid clustering columns for the quarantine table. | @@ -68,6 +68,9 @@ These fields are required on every flow entry. | `silver_table_path_` | string | Non-UC | External path for the Silver table (required without Unity Catalog). | | `silver_transformation_json_` | 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_` | string | No | Path to a DQE file. | +| `silver_database_quarantine_` | 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_` | 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. | diff --git a/docs/docs/reference/dq-rules.md b/docs/docs/reference/dq-rules.md index d47d8c01..9f52f817 100644 --- a/docs/docs/reference/dq-rules.md +++ b/docs/docs/reference/dq-rules.md @@ -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 @@ -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. diff --git a/integration_tests/conf/json/cloudfiles-onboarding.template b/integration_tests/conf/json/cloudfiles-onboarding.template index a443ad7b..ecb8ffca 100644 --- a/integration_tests/conf/json/cloudfiles-onboarding.template +++ b/integration_tests/conf/json/cloudfiles-onboarding.template @@ -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" } ] \ No newline at end of file diff --git a/integration_tests/conf/json/dqe/customers/no_quarantine_rules.json b/integration_tests/conf/json/dqe/customers/no_quarantine_rules.json new file mode 100644 index 00000000..23d82aaa --- /dev/null +++ b/integration_tests/conf/json/dqe/customers/no_quarantine_rules.json @@ -0,0 +1,6 @@ +{ + "expect_or_drop": { + "valid_id": "id IS NOT NULL", + "valid_operation": "operation IN ('APPEND', 'DELETE', 'UPDATE')" + } +} diff --git a/integration_tests/conf/yml/cloudfiles-onboarding.template.yml b/integration_tests/conf/yml/cloudfiles-onboarding.template.yml index 46634f70..0f6b0c2b 100644 --- a/integration_tests/conf/yml/cloudfiles-onboarding.template.yml +++ b/integration_tests/conf/yml/cloudfiles-onboarding.template.yml @@ -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' diff --git a/integration_tests/conf/yml/dqe/customers/no_quarantine_rules.yml b/integration_tests/conf/yml/dqe/customers/no_quarantine_rules.yml new file mode 100644 index 00000000..65de086e --- /dev/null +++ b/integration_tests/conf/yml/dqe/customers/no_quarantine_rules.yml @@ -0,0 +1,3 @@ +expect_or_drop: + valid_id: id IS NOT NULL + valid_operation: operation IN ('APPEND', 'DELETE', 'UPDATE') diff --git a/integration_tests/notebooks/cloudfile_runners/validate.py b/integration_tests/notebooks/cloudfile_runners/validate.py index aa11189f..8487951a 100644 --- a/integration_tests/notebooks/cloudfile_runners/validate.py +++ b/integration_tests/notebooks/cloudfile_runners/validate.py @@ -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") @@ -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 diff --git a/integration_tests/run_integration_tests.py b/integration_tests/run_integration_tests.py index f3bd5378..2bf1e601 100644 --- a/integration_tests/run_integration_tests.py +++ b/integration_tests/run_integration_tests.py @@ -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}" diff --git a/src/databricks/labs/sdp_meta/onboard_dataflowspec.py b/src/databricks/labs/sdp_meta/onboard_dataflowspec.py index c7ad0965..61388358 100644 --- a/src/databricks/labs/sdp_meta/onboard_dataflowspec.py +++ b/src/databricks/labs/sdp_meta/onboard_dataflowspec.py @@ -1529,6 +1529,7 @@ def __get_bronze_dataflow_spec_dataframe(self, onboarding_df, env): data_quality_expectations = None quarantine_target_details = {} quarantine_table_properties = {} + has_quarantine_expectations = False if f"bronze_data_quality_expectations_json_{env}" in onboarding_row: bronze_data_quality_expectations_json = onboarding_row[ f"bronze_data_quality_expectations_json_{env}" @@ -1537,10 +1538,36 @@ def __get_bronze_dataflow_spec_dataframe(self, onboarding_df, env): data_quality_expectations = self.__get_data_quality_expecations( bronze_data_quality_expectations_json ) - if onboarding_row["bronze_quarantine_table"]: - quarantine_target_details, quarantine_table_properties = self.__get_quarantine_details( - env, "bronze", onboarding_row + has_quarantine_expectations = ( + self.__has_quarantine_expectations( + data_quality_expectations ) + ) + quarantine_target_configured = ( + "bronze_quarantine_table" in onboarding_row + and bool(onboarding_row["bronze_quarantine_table"]) + ) + quarantine_database_configured = ( + f"bronze_database_quarantine_{env}" in onboarding_row + and bool(onboarding_row[f"bronze_database_quarantine_{env}"]) + ) + if has_quarantine_expectations and not ( + quarantine_target_configured and quarantine_database_configured + ): + logger.warning( + "Bronze DQE contains non-empty expect_or_quarantine rules " + "but its quarantine target is incomplete. No quarantine " + "table will be created; if these are the only DQE rules, " + "the pipeline will not declare a main output table. Add " + "bronze_quarantine_table and " + "bronze_database_quarantine_%s. Missing targets remain " + "allowed for backward compatibility.", + env, + ) + if has_quarantine_expectations or quarantine_target_configured: + quarantine_target_details, quarantine_table_properties = self.__get_quarantine_details( + env, "bronze", onboarding_row + ) append_flows, append_flows_schemas = self.get_append_flows_json( onboarding_row, "bronze", env @@ -2269,6 +2296,16 @@ def __get_data_quality_expecations(self, file_path): return None return json.dumps(parsed) + @staticmethod + def __has_quarantine_expectations(data_quality_expectations): + """Return whether serialized DQE contains quarantine rules.""" + if not data_quality_expectations: + return False + parsed = json.loads(data_quality_expectations) + if not isinstance(parsed, dict): + return False + return bool(parsed.get("expect_or_quarantine")) + def __get_silver_dataflow_spec_dataframe(self, onboarding_df, env): """Get silver_dataflow_spec method transform onboarding dataframe to silver dataflowSpec dataframe. @@ -2539,6 +2576,7 @@ def __get_silver_dataflow_spec_dataframe(self, onboarding_df, env): silver_quarantine_target_details = None silver_quarantine_table_properties = None silver_quarantine_cluster_by = None + has_quarantine_expectations = False if f"silver_data_quality_expectations_json_{env}" in onboarding_row: silver_data_quality_expectations_json = onboarding_row[ f"silver_data_quality_expectations_json_{env}" @@ -2547,13 +2585,45 @@ def __get_silver_dataflow_spec_dataframe(self, onboarding_df, env): data_quality_expectations = self.__get_data_quality_expecations( silver_data_quality_expectations_json ) - silver_quarantine_target_details, silver_quarantine_table_properties = self.__get_quarantine_details( - env, "silver", onboarding_row + has_quarantine_expectations = ( + self.__has_quarantine_expectations( + data_quality_expectations + ) + ) + quarantine_target_configured = ( + "silver_quarantine_table" in onboarding_row + and bool(onboarding_row["silver_quarantine_table"]) + ) + quarantine_database_configured = ( + f"silver_database_quarantine_{env}" in onboarding_row + and bool(onboarding_row[f"silver_database_quarantine_{env}"]) + ) + if has_quarantine_expectations and not ( + quarantine_target_configured and quarantine_database_configured + ): + logger.warning( + "Silver DQE contains non-empty expect_or_quarantine rules " + "but its quarantine target is incomplete. No quarantine " + "table will be created; if these are the only DQE rules, " + "the pipeline will not declare a main output table. Add " + "silver_quarantine_table and " + "silver_database_quarantine_%s. Missing targets remain " + "allowed for backward compatibility.", + env, ) - silver_quarantine_cluster_by = self.__get_cluster_by_properties( - onboarding_row, + if has_quarantine_expectations or quarantine_target_configured: + ( + silver_quarantine_target_details, silver_quarantine_table_properties, - "silver_quarantine_cluster_by" + ) = self.__get_quarantine_details( + env, "silver", onboarding_row + ) + silver_quarantine_cluster_by = ( + self.__get_cluster_by_properties( + onboarding_row, + silver_quarantine_table_properties, + "silver_quarantine_cluster_by", + ) ) append_flows, append_flow_schemas = self.get_append_flows_json( onboarding_row, layer="silver", env=env diff --git a/tests/test_integration_runner_yaml.py b/tests/test_integration_runner_yaml.py index a3d183aa..7c7e484b 100644 --- a/tests/test_integration_runner_yaml.py +++ b/tests/test_integration_runner_yaml.py @@ -98,6 +98,118 @@ def test_write_onboarding_file_yaml_round_trip(self): finally: os.unlink(tmp.name) + def test_cloudfiles_quarantine_compat_flow_matches_json_and_yaml(self): + json_path = os.path.join( + _PROJECT_ROOT, + "integration_tests/conf/json/cloudfiles-onboarding.template", + ) + yaml_path = os.path.join( + _PROJECT_ROOT, + "integration_tests/conf/yml/cloudfiles-onboarding.template.yml", + ) + with open(json_path) as fh: + json_payload = json.load(fh) + with open(yaml_path) as fh: + yaml_payload = yaml.safe_load(fh) + + json_flow = next( + row for row in json_payload if row["data_flow_id"] == "190" + ) + yaml_flow = next( + row for row in yaml_payload if row["data_flow_id"] == "190" + ) + for flow in (json_flow, yaml_flow): + self.assertEqual( + flow["data_flow_group"], "QUARANTINE_COMPAT" + ) + self.assertFalse( + any("quarantine" in key for key in flow) + ) + self.assertIn( + "bronze_data_quality_expectations_json_it", flow + ) + self.assertIn( + "silver_data_quality_expectations_json_it", flow + ) + + self.assertEqual( + json_flow["bronze_table"], yaml_flow["bronze_table"] + ) + self.assertEqual( + json_flow["silver_table"], yaml_flow["silver_table"] + ) + + json_dqe_only = next( + row for row in json_payload if row["data_flow_id"] == "191" + ) + yaml_dqe_only = next( + row for row in yaml_payload if row["data_flow_id"] == "191" + ) + for flow in (json_dqe_only, yaml_dqe_only): + self.assertEqual( + flow["data_flow_group"], "DQE_NO_QUARANTINE" + ) + self.assertFalse( + any("quarantine" in key for key in flow) + ) + self.assertIn( + "no_quarantine_rules", + flow["bronze_data_quality_expectations_json_it"], + ) + self.assertIn( + "no_quarantine_rules", + flow["silver_data_quality_expectations_json_it"], + ) + + dqe_json_path = os.path.join( + _PROJECT_ROOT, + "integration_tests/conf/json/dqe/customers/" + "no_quarantine_rules.json", + ) + dqe_yaml_path = os.path.join( + _PROJECT_ROOT, + "integration_tests/conf/yml/dqe/customers/" + "no_quarantine_rules.yml", + ) + with open(dqe_json_path) as fh: + json_dqe = json.load(fh) + with open(dqe_yaml_path) as fh: + yaml_dqe = yaml.safe_load(fh) + self.assertEqual(json_dqe, yaml_dqe) + self.assertIn("expect_or_drop", json_dqe) + self.assertNotIn("expect_or_quarantine", json_dqe) + + def test_interactive_demo_dqe_without_quarantine_metadata_parity(self): + json_path = os.path.join( + _PROJECT_ROOT, "demo/conf/json/sample_onboarding.json" + ) + yaml_path = os.path.join( + _PROJECT_ROOT, "demo/conf/yml/sample_onboarding.yml" + ) + with open(json_path) as fh: + json_payload = json.load(fh) + with open(yaml_path) as fh: + yaml_payload = yaml.safe_load(fh) + + for payload in (json_payload, yaml_payload): + customers = next( + row for row in payload if row["data_flow_id"] == "100" + ) + transactions = next( + row for row in payload if row["data_flow_id"] == "101" + ) + for flow in (customers, transactions): + self.assertIn( + "silver_data_quality_expectations_json_prod", flow + ) + self.assertFalse( + any( + key.startswith("silver_") + and "quarantine" in key + for key in flow + ) + ) + class SilverDqeYamlPathRewriteTests(unittest.TestCase): """Verify silver/DQ path rewriting in YAML mode against dedicated .yml siblings. diff --git a/tests/test_onboard_dataflowspec.py b/tests/test_onboard_dataflowspec.py index 1c77d365..cdb11ac5 100644 --- a/tests/test_onboard_dataflowspec.py +++ b/tests/test_onboard_dataflowspec.py @@ -665,6 +665,214 @@ def test_get_data_quality_expectations_unsupported_no_longer_silently_drops(self "tests/resources/schema.ddl" ) + def test_has_quarantine_expectations_rejects_empty_and_non_mapping_values(self): + has_quarantine_expectations = ( + OnboardDataflowspec._OnboardDataflowspec__has_quarantine_expectations + ) + + self.assertFalse(has_quarantine_expectations(None)) + self.assertFalse(has_quarantine_expectations(json.dumps([]))) + + def test_empty_dqe_paths_do_not_require_quarantine_fields(self): + with tempfile.TemporaryDirectory() as tmp_dir: + with open(self.onboarding_json_file, "r") as source: + onboarding_row = copy.deepcopy(json.load(source)[0]) + for field_name in list(onboarding_row): + if "quarantine" in field_name: + del onboarding_row[field_name] + onboarding_row["bronze_data_quality_expectations_json_dev"] = "" + onboarding_row["silver_data_quality_expectations_json_dev"] = "" + + onboarding_path = os.path.join(tmp_dir, "onboarding.json") + with open(onboarding_path, "w") as target: + json.dump([onboarding_row], target) + + params = copy.deepcopy(self.onboarding_bronze_silver_params_map) + params["onboarding_file_path"] = onboarding_path + onboarder = OnboardDataflowspec(self.spark, params) + onboarding_df = onboarder._OnboardDataflowspec__get_onboarding_file_dataframe( + onboarding_path + ) + + bronze_row = onboarder._OnboardDataflowspec__get_bronze_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + silver_row = onboarder._OnboardDataflowspec__get_silver_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + + self.assertIsNone(bronze_row.dataQualityExpectations) + self.assertEqual(bronze_row.quarantineTargetDetails, {}) + self.assertIsNone(silver_row.dataQualityExpectations) + self.assertIsNone(silver_row.quarantineTargetDetails) + + def _stage_onboarding_with_dqe_without_quarantine( + self, tmp_dir, extension, dqe_payload + ): + with open(self.onboarding_json_file, "r") as source: + onboarding_row = copy.deepcopy(json.load(source)[0]) + for field_name in list(onboarding_row): + if "quarantine" in field_name: + del onboarding_row[field_name] + + dqe_path = os.path.join(tmp_dir, f"expectations.{extension}") + onboarding_path = os.path.join(tmp_dir, f"onboarding.{extension}") + with open(dqe_path, "w") as target: + if extension == "json": + json.dump(dqe_payload, target) + else: + yaml.safe_dump(dqe_payload, target, sort_keys=False) + + onboarding_row["bronze_data_quality_expectations_json_dev"] = dqe_path + onboarding_row["silver_data_quality_expectations_json_dev"] = dqe_path + with open(onboarding_path, "w") as target: + if extension == "json": + json.dump([onboarding_row], target) + else: + yaml.safe_dump([onboarding_row], target, sort_keys=False) + return onboarding_path + + def test_dqe_without_quarantine_rules_does_not_require_quarantine_fields(self): + dqe_payloads = { + "expect": {"expect": {"observed_id": "id IS NOT NULL"}}, + "expect_or_drop": { + "expect_or_drop": {"valid_id": "id IS NOT NULL"} + }, + "expect_or_fail": { + "expect_or_fail": {"required_id": "id IS NOT NULL"} + }, + "empty_expect_or_quarantine": {"expect_or_quarantine": {}}, + } + for extension in ("json", "yml"): + for case_name, dqe_payload in dqe_payloads.items(): + with self.subTest(extension=extension, case=case_name): + with tempfile.TemporaryDirectory() as tmp_dir: + onboarding_path = ( + self._stage_onboarding_with_dqe_without_quarantine( + tmp_dir, + extension, + dqe_payload, + ) + ) + params = copy.deepcopy( + self.onboarding_bronze_silver_params_map + ) + params["onboarding_file_path"] = onboarding_path + onboarder = OnboardDataflowspec(self.spark, params) + onboarding_df = onboarder._OnboardDataflowspec__get_onboarding_file_dataframe( + onboarding_path + ) + + bronze_row = onboarder._OnboardDataflowspec__get_bronze_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + silver_row = onboarder._OnboardDataflowspec__get_silver_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + + self.assertEqual( + json.loads(bronze_row.dataQualityExpectations), + dqe_payload, + ) + self.assertEqual( + json.loads(silver_row.dataQualityExpectations), + dqe_payload, + ) + self.assertEqual( + bronze_row.quarantineTargetDetails, {} + ) + self.assertIsNone( + silver_row.quarantineTargetDetails + ) + + def test_legacy_quarantine_rules_without_targets_remain_compatible(self): + dqe_payload = { + "expect_or_quarantine": {"valid_id": "id IS NOT NULL"} + } + with tempfile.TemporaryDirectory() as tmp_dir: + onboarding_path = ( + self._stage_onboarding_with_dqe_without_quarantine( + tmp_dir, "json", dqe_payload + ) + ) + params = copy.deepcopy(self.onboarding_bronze_silver_params_map) + params["onboarding_file_path"] = onboarding_path + onboarder = OnboardDataflowspec(self.spark, params) + onboarding_df = onboarder._OnboardDataflowspec__get_onboarding_file_dataframe( + onboarding_path + ) + + with self.assertLogs( + "databricks.labs.sdp_meta", level="WARNING" + ) as captured: + bronze_row = onboarder._OnboardDataflowspec__get_bronze_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + silver_row = onboarder._OnboardDataflowspec__get_silver_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + + self.assertEqual( + json.loads(bronze_row.dataQualityExpectations), + dqe_payload, + ) + self.assertEqual( + json.loads(silver_row.dataQualityExpectations), + dqe_payload, + ) + self.assertEqual(bronze_row.quarantineTargetDetails, {}) + self.assertEqual(silver_row.quarantineTargetDetails, {}) + warning_output = "\n".join(captured.output) + self.assertIn( + "Bronze DQE contains non-empty expect_or_quarantine rules", + warning_output, + ) + self.assertIn( + "Silver DQE contains non-empty expect_or_quarantine rules", + warning_output, + ) + self.assertIn( + "Missing targets remain allowed for backward compatibility", + warning_output, + ) + + def test_optional_quarantine_metadata_is_preserved_without_dqe(self): + with tempfile.TemporaryDirectory() as tmp_dir: + with open(self.onboarding_json_file, "r") as source: + onboarding_row = copy.deepcopy(json.load(source)[0]) + onboarding_row.pop( + "bronze_data_quality_expectations_json_dev", None + ) + onboarding_row.pop( + "silver_data_quality_expectations_json_dev", None + ) + onboarding_path = os.path.join(tmp_dir, "onboarding.json") + with open(onboarding_path, "w") as target: + json.dump([onboarding_row], target) + + params = copy.deepcopy(self.onboarding_bronze_silver_params_map) + params["onboarding_file_path"] = onboarding_path + onboarder = OnboardDataflowspec(self.spark, params) + onboarding_df = onboarder._OnboardDataflowspec__get_onboarding_file_dataframe( + onboarding_path + ) + + bronze_row = onboarder._OnboardDataflowspec__get_bronze_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + silver_row = onboarder._OnboardDataflowspec__get_silver_dataflow_spec_dataframe( + onboarding_df, "dev" + ).collect()[0] + + self.assertEqual( + bronze_row.quarantineTargetDetails["table"], + onboarding_row["bronze_quarantine_table"], + ) + self.assertEqual( + silver_row.quarantineTargetDetails["table"], + onboarding_row["silver_quarantine_table"], + ) + def test_validate_params_for_onboardBronzeDataflowSpec(self): """Test for onboardDataflowspec parameters.""" onboarding_params_map = copy.deepcopy(self.onboarding_bronze_silver_params_map)