Skip to content

Commit 09cee82

Browse files
authored
Refactored methods for loading and saving checks (#487)
## Changes <!-- Summary of your changes that are easy to understand. Add screenshots when necessary --> * Moved all load and save checks methods from `DQEngine` to a separate module. The previous approach of continuously expanding DQEngine was not scalable or maintainable. This change will make it easier to add more ways of saving and loading checks. * Refactored code in various places to improve modularity and adhere to the separation of concerns. The original `DQEngine` had grown beyond its intended scope, so responsibilities were redistributed. The change should provide a good basis for extensions without overloading the engine. * All conversion methods for checks (serialization and deserialization) from/to dict and `DQRule` has been centralized in `checks_serializer` module. It is now possible to convert checks defined in dict to `DQRule` and vice versa. Detailed documentation is available [here](https://databrickslabs.github.io/dqx/docs/reference/quality_rules/). * Added support for saving checks to json file/workspace file * Added validation of location when saving and loading checks * Added validation of columns passed to `compare_datasets` check to make sure only simple expressions are allowed * Added parameter to `validate-checks` cli command to enable/disable validating custom check functions (`validate-custom-check-functions`) BREAKING CHANGES! If you are loading or saving checks from a storage (file, workspace file, table, installation), you are affected. We are deprecating the the below methods. We are keeping the methods in the `DQEngine` but you should update your code as these methods will be removed in future versions. * Loading checks to storage has been unified under `load_checks` method. The following methods have been removed from the `DQEngine`: `load_checks_from_local_file`, `load_checks_from_workspace_file`, `load_checks_from_installation`, `load_checks_from_table`. * Saving checks in storage has been unified under `load_checks` method. The following methods have been removed from the `DQEngine`: `save_checks_in_local_file`, `save_checks_in_workspace_file`, `save_checks_in_installation`, `save_checks_in_table`. The `save_checks` and `load_checks` take `config` as a parameter, which determines the storage types used. The following storage configs are currently supported: * `FileChecksStorageConfig`: file in the local filesystem (YAML or JSON) * `WorkspaceFileChecksStorageConfig`: file in the workspace (YAML or JSON) * `TableChecksStorageConfig`: a table * `InstallationChecksStorageConfig`: storage defined in the installation context, using either the `checks_table` or `checks_file` field from the run configuration. * The `load_run_config` method has been moved to `config_loader.RunConfigLoader`, as it is not intended for direct use and falls outside the `DQEngine` core responsibilities. ### Tests <!-- How is this tested? Please see the checklist below and also describe any other relevant tests --> - [x] manually tested - [x] added unit tests - [x] added integration tests
1 parent d9d1000 commit 09cee82

44 files changed

Lines changed: 2095 additions & 1200 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

demos/dqx_demo_asset_bundle/dqx_demo_notebook.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -581,7 +581,7 @@
581581

582582
# DBTITLE 1,Common Imports
583583
from databricks.sdk import WorkspaceClient
584-
from databricks.labs.dqx.config import InputConfig, OutputConfig
584+
from databricks.labs.dqx.config import InputConfig, OutputConfig, WorkspaceFileChecksStorageConfig
585585
from databricks.labs.dqx.engine import DQEngine
586586

587587
# COMMAND ----------
@@ -594,7 +594,7 @@
594594
displayHTML(f'<a href="/#workspace{sensor_rules_file_path}" target="_blank">Quality rules file for sensor dataset</a>')
595595

596596
# Load the checks
597-
maintenance_quality_checks = dq_engine.load_checks_from_workspace_file(workspace_path=sensor_rules_file_path)
597+
maintenance_quality_checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=sensor_rules_file_path))
598598

599599
# Apply the checks and write the output data
600600
dq_engine.apply_checks_by_metadata_and_save_in_table(
@@ -614,7 +614,7 @@
614614
displayHTML(f'<a href="/#workspace{maintenance_rules_file_path}" target="_blank">Quality rules file for maintenance dataset</a>')
615615

616616
# Load the checks
617-
maintenance_quality_checks = dq_engine.load_checks_from_workspace_file(workspace_path=maintenance_rules_file_path)
617+
maintenance_quality_checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=maintenance_rules_file_path))
618618

619619
# Apply the checks and write the output data
620620
dq_engine.apply_checks_by_metadata_and_save_in_table(

demos/dqx_demo_dbt/models/dummy_model_dq.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
11
import yaml
2+
3+
from databricks.labs.dqx.config import WorkspaceFileChecksStorageConfig
24
from databricks.labs.dqx.engine import DQEngine
35
from databricks.sdk import WorkspaceClient
46

@@ -26,7 +28,7 @@ def model(dbt, session):
2628

2729
# Checks can also be loaded from a file in the workspace or delta table
2830
# checks_path = dbt.config.get("checks_file_path") # get from dbt var
29-
# checks = dq_engine.load_checks_from_workspace_file(checks_path)
31+
# checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=checks_path))
3032

3133
# apply quality checks with issues reported in _warnings and _errors columns
3234
df = dq_engine.apply_checks_by_metadata(input_df, checks)

demos/dqx_demo_library.py

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
from databricks.labs.dqx.profiler.profiler import DQProfiler
3131
from databricks.labs.dqx.profiler.generator import DQGenerator
3232
from databricks.labs.dqx.profiler.dlt_generator import DQDltGenerator
33+
from databricks.labs.dqx.config import WorkspaceFileChecksStorageConfig, TableChecksStorageConfig
3334
from databricks.labs.dqx.engine import DQEngine
3435
from databricks.sdk import WorkspaceClient
3536
import os
@@ -81,10 +82,13 @@
8182
user_name = spark.sql("select current_user() as user").collect()[0]["user"]
8283
checks_file = f"{demo_file_directory}/dqx_demo_checks.yml"
8384
dq_engine = DQEngine(ws)
84-
dq_engine.save_checks_in_workspace_file(checks=checks, workspace_path=checks_file)
85+
dq_engine.save_checks(checks=checks, config=WorkspaceFileChecksStorageConfig(location=checks_file))
8586

8687
# save generated checks in a Delta table
87-
dq_engine.save_checks_in_table(checks=checks, table_name=f"{demo_catalog_name}.{demo_schema_name}.dqx_checks_table", mode="overwrite")
88+
dq_engine.save_checks(
89+
checks=checks,
90+
config=TableChecksStorageConfig(location=f"{demo_catalog_name}.{demo_schema_name}.dqx_checks_table", mode="overwrite")
91+
)
8892

8993
# COMMAND ----------
9094

@@ -95,12 +99,13 @@
9599

96100
from databricks.labs.dqx.engine import DQEngine
97101
from databricks.sdk import WorkspaceClient
102+
from databricks.labs.dqx.config import WorkspaceFileChecksStorageConfig
98103

99104
input_df = spark.createDataFrame([[1, 3, 3, 2], [3, 3, None, 1]], schema)
100105

101106
# load checks from a file
102107
dq_engine = DQEngine(WorkspaceClient())
103-
checks = dq_engine.load_checks_from_workspace_file(workspace_path=checks_file)
108+
checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=checks_file))
104109

105110
# Option 1: apply quality rules and quarantine invalid records
106111
valid_df, quarantine_df = dq_engine.apply_checks_by_metadata_and_split(input_df, checks)
@@ -120,12 +125,13 @@
120125

121126
from databricks.labs.dqx.engine import DQEngine
122127
from databricks.sdk import WorkspaceClient
128+
from databricks.labs.dqx.config import TableChecksStorageConfig
123129

124130
input_df = spark.createDataFrame([[1, 3, 3, 2], [3, 3, None, 1]], schema)
125131

126132
# load checks from a Delta table
127133
dq_engine = DQEngine(WorkspaceClient())
128-
checks = dq_engine.load_checks_from_table(table_name=f"{demo_catalog_name}.{demo_schema_name}.dqx_checks_table")
134+
checks = dq_engine.load_checks(config=TableChecksStorageConfig(location=f"{demo_catalog_name}.{demo_schema_name}.dqx_checks_table"))
129135

130136
# Option 1: apply quality rules and quarantine invalid records
131137
valid_df, quarantine_df = dq_engine.apply_checks_by_metadata_and_split(input_df, checks)

demos/dqx_demo_tool.py

Lines changed: 28 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -94,15 +94,19 @@
9494
from databricks.labs.dqx.profiler.profiler import DQProfiler
9595
from databricks.labs.dqx.profiler.generator import DQGenerator
9696
from databricks.labs.dqx.engine import DQEngine
97+
from databricks.labs.dqx.config import InstallationChecksStorageConfig, WorkspaceFileChecksStorageConfig
98+
from databricks.labs.dqx.config_loader import RunConfigLoader
9799
from databricks.labs.dqx.utils import read_input_data
98100
from databricks.sdk import WorkspaceClient
99101

102+
100103
dqx_product_name = dbutils.widgets.get("dqx_product_name")
101104

102-
# setup the DQEngine
103105
ws = WorkspaceClient()
104106
dq_engine = DQEngine(ws)
105-
run_config = dq_engine.load_run_config(run_config_name="default", assume_user=True, product_name=dqx_product_name)
107+
108+
# load the run configuration
109+
run_config = RunConfigLoader(ws).load_run_config(run_config_name="default", product_name=dqx_product_name)
106110

107111
# read the input data, limit to 1000 rows for demo purpose
108112
input_df = read_input_data(spark, run_config.input_config).limit(1000)
@@ -120,9 +124,10 @@
120124
print(yaml.safe_dump(checks))
121125

122126
# save generated checks to location specified in the default run configuration inside workspace installation folder
123-
dq_engine.save_checks_in_installation(checks, run_config_name="default", product_name=dqx_product_name)
124-
# or save it to an arbitrary workspace location
125-
#dq_engine.save_checks_in_workspace_file(checks, workspace_path="/Shared/App1/checks.yml")
127+
dq_engine.save_checks(checks, config=InstallationChecksStorageConfig(run_config_name="default", product_name=dqx_product_name))
128+
129+
# or save checks in arbitrary workspace location
130+
#dq_engine.save_checks(checks, config=WorkspaceFileChecksStorageConfig(location="/Shared/App1/checks.yml"))
126131

127132
# COMMAND ----------
128133

@@ -136,6 +141,8 @@
136141
import yaml
137142
from databricks.labs.dqx.engine import DQEngine
138143
from databricks.sdk import WorkspaceClient
144+
from databricks.labs.dqx.config import InstallationChecksStorageConfig, WorkspaceFileChecksStorageConfig
145+
139146

140147
checks = yaml.safe_load("""
141148
- check:
@@ -180,10 +187,12 @@
180187
assert not status.has_errors
181188

182189
dq_engine = DQEngine(WorkspaceClient())
190+
183191
# save checks to location specified in the default run configuration inside workspace installation folder
184-
dq_engine.save_checks_in_installation(checks, run_config_name="default", product_name=dqx_product_name)
185-
# or save it to an arbitrary workspace location
186-
#dq_engine.save_checks_in_workspace_file(checks, workspace_path="/Shared/App1/checks.yml")
192+
dq_engine.save_checks(checks, config=InstallationChecksStorageConfig(run_config_name="default", product_name=dqx_product_name))
193+
194+
# or save checks in arbitrary workspace location
195+
#dq_engine.save_checks(checks, config=WorkspaceFileChecksStorageConfig(location="/Shared/App1/checks.yml"))
187196

188197
# COMMAND ----------
189198

@@ -195,21 +204,27 @@
195204
from databricks.labs.dqx.engine import DQEngine
196205
from databricks.labs.dqx.utils import read_input_data
197206
from databricks.sdk import WorkspaceClient
207+
from databricks.labs.dqx.config import InstallationChecksStorageConfig, WorkspaceFileChecksStorageConfig
208+
from databricks.labs.dqx.config_loader import RunConfigLoader
198209

199-
run_config = dq_engine.load_run_config(run_config_name="default", assume_user=True, product_name=dqx_product_name)
210+
211+
dq_engine = DQEngine(WorkspaceClient())
212+
213+
# load the run configuration
214+
run_config = RunConfigLoader(ws).load_run_config(run_config_name="default", assume_user=True, product_name=dqx_product_name)
200215

201216
# read the data, limit to 1000 rows for demo purpose
202217
bronze_df = read_input_data(spark, run_config.input_config).limit(1000)
203218

204219
# apply your business logic here
205220
bronze_transformed_df = bronze_df.filter("vendor_id in (1, 2)")
206221

207-
dq_engine = DQEngine(WorkspaceClient())
208-
209222
# load checks from location defined in the run configuration
210-
checks = dq_engine.load_checks_from_installation(assume_user=True, run_config_name="default", product_name=dqx_product_name)
223+
224+
checks = dq_engine.load_checks(config=InstallationChecksStorageConfig(assume_user=True, run_config_name="default", product_name=dqx_product_name))
225+
211226
# or load checks from arbitrary workspace file
212-
# checks = dq_engine.load_checks_from_workspace_file(workspace_path="/Shared/App1/checks.yml")
227+
#checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location="/Shared/App1/checks.yml"))
213228
print(checks)
214229

215230
# Option 1: apply quality rules and quarantine invalid records

demos/dqx_manufacturing_demo.py

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -108,7 +108,7 @@
108108
from pyspark.sql.types import *
109109
from pyspark.sql.functions import *
110110
from datetime import datetime
111-
import delta
111+
112112

113113
if spark.catalog.tableExists(sensor_table) and spark.table(sensor_table).count() > 0:
114114
print(
@@ -634,14 +634,14 @@
634634
# COMMAND ----------
635635

636636
# DBTITLE 1,Common Imports
637-
import os
638637
import yaml
639638
from pprint import pprint
640639

641640
from databricks.sdk import WorkspaceClient
642641
from databricks.labs.dqx.profiler.profiler import DQProfiler
643642
from databricks.labs.dqx.profiler.generator import DQGenerator
644643
from databricks.labs.dqx.engine import DQEngine
644+
from databricks.labs.dqx.config import WorkspaceFileChecksStorageConfig, TableChecksStorageConfig
645645

646646
# COMMAND ----------
647647

@@ -697,7 +697,7 @@
697697
maintenance_dq_rules_yaml = f"{quality_rules_path}/maintenance_dq_rules.yml"
698698

699699
# Save file in a workspace path
700-
dq_engine.save_checks_in_workspace_file(maintenance_checks, workspace_path=maintenance_dq_rules_yaml)
700+
dq_engine.save_checks(maintenance_checks, config=WorkspaceFileChecksStorageConfig(location=maintenance_dq_rules_yaml))
701701

702702
# display the link to the saved checks
703703
displayHTML(f'<a href="/#workspace{maintenance_dq_rules_yaml}" target="_blank">Maintenance Data Quality Rules YAML</a>')
@@ -707,7 +707,7 @@
707707
# DBTITLE 1,Save the quality rules in delta table
708708
# or save in delta table
709709
maintenance_quality_rules_table = f"{database}.{schema}.maintenance_inferred_quality_rules"
710-
dq_engine.save_checks_in_table(table_name=maintenance_quality_rules_table, checks=maintenance_checks, run_config_name="maintenance")
710+
dq_engine.save_checks(maintenance_checks, config=TableChecksStorageConfig(location=maintenance_quality_rules_table, run_config_name="maintenance"))
711711

712712
# COMMAND ----------
713713

@@ -718,10 +718,9 @@
718718
# COMMAND ----------
719719

720720
# Load checks from workspace file
721-
quality_checks = dq_engine.load_checks_from_workspace_file(workspace_path=maintenance_dq_rules_yaml)
722-
721+
quality_checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=maintenance_dq_rules_yaml))
723722
# or Load checks from a table
724-
# quality_checks = dq_engine.load_checks_from_table(table_name=fq_tbl, run_config_name="maintenance")
723+
#quality_checks = dq_engine.load_checks(config=TableChecksStorageConfig(location=maintenance_quality_rules_table, run_config_name="maintenance"))
725724

726725
# Apply checks on input data
727726
valid_df, quarantined_df = dq_engine.apply_checks_by_metadata_and_split(mntnc_bronze_df, quality_checks)
@@ -810,7 +809,7 @@
810809

811810
# save checks in a workspace location
812811
sensor_dq_rules_yaml = f"{quality_rules_path}/sensor_dq_rules.yml"
813-
dq_engine.save_checks_in_workspace_file(sensor_dq_checks, workspace_path=sensor_dq_rules_yaml)
812+
dq_engine.save_checks(sensor_dq_checks, config=WorkspaceFileChecksStorageConfig(location=sensor_dq_rules_yaml))
814813

815814
# display the link to the saved checks
816815
displayHTML(f'<a href="/#workspace{sensor_dq_rules_yaml}" target="_blank">Sensor Data Quality Rules YAML</a>')
@@ -827,7 +826,7 @@
827826
sensor_bronze_df = spark.read.table(sensor_table)
828827

829828
# Load quality rules from YAML file
830-
sensor_dq_checks = dq_engine.load_checks_from_workspace_file(workspace_path=sensor_dq_rules_yaml)
829+
sensor_dq_checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=sensor_dq_rules_yaml))
831830

832831
# Apply checks on input data
833832
valid_df, quarantined_df = dq_engine.apply_checks_by_metadata_and_split(sensor_bronze_df, sensor_dq_checks)
@@ -859,6 +858,7 @@
859858
from pyspark.sql import Column as col
860859
from databricks.labs.dqx.check_funcs import make_condition
861860

861+
862862
def firmware_version_start_with_v(column: str) -> col:
863863
column_expr = F.expr(column)
864864

@@ -886,7 +886,7 @@ def firmware_version_start_with_v(column: str) -> col:
886886

887887
# Save the YAML file with the new custom DQ rule
888888
sensor_custom_dq_rules_yaml = f"{quality_rules_path}/sensor_custom_dq_rules.yml"
889-
dq_engine.save_checks_in_workspace_file(sensor_dq_checks, workspace_path=sensor_custom_dq_rules_yaml)
889+
dq_engine.save_checks(sensor_dq_checks, config=WorkspaceFileChecksStorageConfig(location=sensor_custom_dq_rules_yaml))
890890

891891
# display the link to the saved checks
892892
displayHTML(f'<a href="/#workspace{sensor_custom_dq_rules_yaml}" target="_blank">Sensor Custom Data Quality Rules YAML</a>')
@@ -897,8 +897,7 @@ def firmware_version_start_with_v(column: str) -> col:
897897
# DBTITLE 1,Apply the DQ Rules on Input Data
898898
dq_engine = DQEngine(WorkspaceClient())
899899

900-
sensor_quality_checks = dq_engine.load_checks_from_workspace_file(
901-
workspace_path=sensor_custom_dq_rules_yaml)
900+
sensor_quality_checks = dq_engine.load_checks(config=WorkspaceFileChecksStorageConfig(location=sensor_custom_dq_rules_yaml))
902901

903902
# Define the custom check
904903
custom_check_functions = {"firmware_version_start_with_v": firmware_version_start_with_v} # list of custom check functions

docs/dqx/docs/guide/data_profiling.mdx

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ Data loaded as a DataFrame can be profiled to generate summary statistics and ca
3030
from databricks.labs.dqx.profiler.profiler import DQProfiler
3131
from databricks.labs.dqx.profiler.generator import DQGenerator
3232
from databricks.labs.dqx.profiler.dlt_generator import DQDltGenerator
33+
from databricks.labs.dqx.config import WorkspaceFileChecksStorageConfig
3334
from databricks.labs.dqx.engine import DQEngine
3435
from databricks.sdk import WorkspaceClient
3536

@@ -47,7 +48,7 @@ Data loaded as a DataFrame can be profiled to generate summary statistics and ca
4748
dq_engine = DQEngine(ws)
4849

4950
# save checks in arbitrary workspace location
50-
dq_engine.save_checks_in_workspace_file(checks, workspace_path="/Shared/App1/checks.yml")
51+
dq_engine.save_checks(checks, config=WorkspaceFileChecksStorageConfig(location="/Shared/App1/checks.yml"))
5152

5253
# generate Lakeflow Pipeline (DLT) expectations
5354
dlt_generator = DQDltGenerator(ws)
@@ -353,6 +354,12 @@ You can save checks defined in code or generated by the profiler to a table or f
353354
<TabItem value="Python" label="Python" default>
354355
```python
355356
from databricks.labs.dqx.engine import DQEngine
357+
from databricks.labs.dqx.config import (
358+
FileChecksStorageConfig,
359+
WorkspaceFileChecksStorageConfig,
360+
InstallationChecksStorageConfig,
361+
TableChecksStorageConfig
362+
)
356363
from databricks.sdk import WorkspaceClient
357364

358365
dq_engine = DQEngine(WorkspaceClient())
@@ -370,29 +377,29 @@ You can save checks defined in code or generated by the profiler to a table or f
370377

371378
# save checks in a local path
372379
# always overwrite the file
373-
dq_engine.save_checks_in_local_file(checks, path="checks.yml")
380+
dq_engine.save_checks(checks, config=FileChecksStorageConfig(location="checks.yml"))
374381

375382
# save checks in arbitrary workspace location
376383
# always overwrite the file
377-
dq_engine.save_checks_in_workspace_file(checks, workspace_path="/Shared/App1/checks.yml")
384+
dq_engine.save_checks(checks, config=WorkspaceFileChecksStorageConfig(location="/Shared/App1/checks.yml"))
378385

379386
# save checks in file defined in 'checks_file' in the run config
380387
# always overwrite the file
381388
# only works if DQX is installed in the workspace
382-
dq_engine.save_checks_in_installation(checks, method="file", assume_user=True, run_config_name="default")
389+
dq_engine.save_checks(checks, config=InstallationChecksStorageConfig(assume_user=True, run_config_name="default"))
383390

384391
# save checks in a Delta table with default run config for filtering
385-
# append checks in the table
386-
dq_engine.save_checks_in_table(checks, table_name="dq.config.checks_table", mode="append")
392+
# append checks in the table for the default run config
393+
dq_engine.save_checks(checks, config=TableChecksStorageConfig(location="dq.config.checks_table", mode="append"))
387394

388395
# save checks in a Delta table with specific run config for filtering
389396
# overwrite checks in the table for the given run config
390-
dq_engine.save_checks_in_table(checks, table_name="dq.config.checks_table", run_config_name="workflow_001", mode="overwrite")
397+
dq_engine.save_checks(checks, config=TableChecksStorageConfig(location="dq.config.checks_table", run_config_name="workflow_001", mode="overwrite"))
391398

392399
# save checks in table defined in 'checks_table' in the run config
393400
# always overwrite checks in the table for the given run config
394401
# only works if DQX is installed in the workspace
395-
dq_engine.save_checks_in_installation(checks, method="table", assume_user=True, run_config_name="default")
402+
dq_engine.save_checks(checks, config=InstallationChecksStorageConfig(assume_user=True, run_config_name="default"))
396403
```
397404
</TabItem>
398405
</Tabs>

0 commit comments

Comments
 (0)