Decouple AI rule generation from Spark session - #1422
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #1422 +/- ##
==========================================
- Coverage 92.97% 92.81% -0.17%
==========================================
Files 103 104 +1
Lines 10819 10897 +78
==========================================
+ Hits 10059 10114 +55
- Misses 760 783 +23
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
✅ 1/1 passed, 19m11s total Running from mcp #249 |
|
❌ 441/442 passed, 1 failed, 7 skipped, 9h6m54s total ❌ test_apply_checks_all_row_checks_as_yaml_with_streaming: pyspark.errors.exceptions.connect.StreamingQueryException: [STREAM_FAILED] Query [id = f63393d6-e91d-438e-912d-9f487c3857a4, runId = 13133a6a-8c09-45c7-9e26-a5c509c22d50] terminated with exception: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS SQLSTATE: XXKST (6m16.945s)Running from acceptance #5471 |
|
❌ 23/55 passed, 32 failed, 8h51m26s total ❌ test_ai_query_explanation_populated_for_anomalous_row: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (12m1.971s)❌ test_ai_query_explanation_degrades_when_endpoint_unavailable: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (12m2.019s)❌ test_apply_anomaly_multiple_checks_by_metadata: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (13m2.648s)❌ test_apply_anomaly_check_by_metadata_criticality_warn: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (19m38.139s)❌ test_apply_anomaly_check_with_contributions: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (20m43.377s)❌ test_ai_query_explanation_references_dominant_feature_within_word_caps: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (20m43.372s)❌ test_ai_query_response_shape_portability[databricks-claude-3-7-sonnet]: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (20m44.529s)❌ test_apply_anomaly_checks_and_split_with_correct_quarantine_structure: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (20m45.269s)❌ test_explicit_columns_no_auto_segment: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (25m4.159s)❌ test_ai_query_explanation_disabled_without_contributions: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (17m14.192s)❌ test_apply_anomaly_check_by_metadata_with_custom_threshold: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (19m13.697s)❌ test_ai_query_explanation_null_for_non_anomalous_row: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (22m14.588s)❌ test_apply_anomaly_check_by_metadata: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (13m33.465s)❌ test_ai_query_response_shape_portability[databricks-llama-4-maverick]: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (14m39.116s)❌ test_anomaly_and_other_checks_combined: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (14m36.995s)❌ test_apply_anomaly_check_by_metadata_with_filter_segmented: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (15m46.357s)❌ test_ai_query_response_shape_portability[databricks-gpt-oss-20b]: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (15m40.112s)❌ test_zero_config_training: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (18m51.323s)❌ test_apply_anomaly_checks: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (11m24.629s)❌ test_apply_anomaly_check_by_metadata_with_contributions: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (10m18.864s)❌ test_apply_anomaly_check_by_metadata_with_columns_autodiscovery: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (8m22.401s)❌ test_ai_query_explanation_redact_columns_filters_output: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (9m19.056s)❌ test_apply_anomaly_check_with_criticality_warn: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (10m25.861s)❌ test_ai_query_response_shape_portability[databricks-gemma-3-12b]: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (11m21.027s)❌ test_ai_query_explanation_on_by_TEST_SCHEMA: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (11m22.956s)❌ test_apply_anomaly_checks_and_split: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (10m3.018s)❌ test_apply_anomaly_check_by_metadata_with_drift_threshold: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (11m12.305s)❌ test_apply_anomaly_check_by_metadata_with_multiple_checks: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (11m9.728s)❌ test_ai_query_explanation_one_call_per_group: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (10m14.986s)❌ test_apply_anomaly_check_info_column_structure: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (9m6.925s)❌ test_ai_query_response_shape_portability[databricks-meta-llama-3-3-70b-instruct]: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (10m14.195s)❌ test_contribution_percentages_sum_to_hundred: pyspark.errors.exceptions.connect.SparkException: [ISOLATION_STARTUP_FAILURE.GENERIC] Failed to start isolated execution environment. Please contact Databricks support. SQLSTATE: XXKSS (11m5.446s)Running from anomaly #1585 |
mwojtyczka
left a comment
There was a problem hiding this comment.
Nice refactor — decoupling rule generation from Spark and reading UC schemas via the SDK is the right shape, and splitting PK detection into DQLLMPrimaryKeyEngine cleanly isolates the one path that genuinely needs Spark. Good test coverage too. A few things worth addressing before merge (most are correctness/consistency around the UC-vs-Spark routing):
- Special-char UC tables silently force Spark (
UC_TABLE_PATTERN). A valid backtick-quoted 3-level table (`my-catalog`.schema.table) is rejected by the pattern and falls through to the Spark path — re-introducing the Spark dependency the PR removes, for exactly the hyphenated/special-char names UC allows. The existingTABLE_PATTERNalready handles backticks. type_textvssimpleString()divergence: the SDK and Spark paths format types differently, so the docstring claim that they're interchangeable isn't quite true — same table, different LLM context.- Routing by omission: non-UC locations (two-level HMS names, views, invalid strings) fall to Spark implicitly rather than being explicitly classified/validated.
DQLLMPrimaryKeyEngine.__init__eagerly builds Spark, inconsistent with the lazy-property pattern introduced forDQGenerator.- Minor: duplicated configurator/dspy-context boilerplate across the two engines; and one PK test asserts an internal
final_statusvalue rather than the documented result contract.
Details inline.
| STORAGE_PATH_PATTERN = re.compile(r"^(/|s3:/|abfss:/|gs:/)") | ||
| # catalog.schema.table or schema.table or database.table (backticks allow special chars like hyphens) | ||
| TABLE_PATTERN = re.compile(r"^(?:(?:`[^`]+`|[a-zA-Z0-9_]+)\.)?(?:`[^`]+`|[a-zA-Z0-9_]+)\.(?:`[^`]+`|[a-zA-Z0-9_]+)$") | ||
| UC_TABLE_PATTERN = re.compile(r"^[a-zA-Z0-9_]+\.[a-zA-Z0-9_]+\.[a-zA-Z0-9_]+$") |
There was a problem hiding this comment.
UC_TABLE_PATTERN only accepts [a-zA-Z0-9_], so it rejects valid backtick-quoted 3-level UC tables (e.g. `my-catalog`.schema.table or main.`my-schema`.tbl). Those are legitimate UC tables the tables.get API can resolve, but they'll fail this match and fall through to the Spark path in generator._get_schema_info — re-introducing the Spark session the PR is trying to avoid. The existing TABLE_PATTERN right above already handles backticks; consider reusing/extending that logic (or stripping backticks) so special-char UC names still take the SDK path. The new test_uc_table_pattern_rejects_non_uc_locations currently codifies this as intended — worth reconsidering.
| NotFound: If the table does not exist or is not accessible. | ||
| """ | ||
| table_info = workspace_client.tables.get(table) | ||
| columns = [{"name": col.name or "", "type": col.type_text or ""} for col in (table_info.columns or [])] |
There was a problem hiding this comment.
col.type_text is the UC DDL type text, which can differ in case/format from Spark's field.dataType.simpleString() used by get_column_metadata (e.g. nested types like ARRAY<INT>/STRUCT<...> vs array<int>/struct<...>). The docstring says the two produce the same shape and are 'interchangeable', but the type values can differ for the same table depending on which path runs, making the LLM prompt context (and thus generated rules) non-deterministic across UC vs Spark. Consider normalizing (e.g. lowercase, or map to Spark's simpleString form) so both paths emit identical type strings.
| Returns: | ||
| A JSON string containing the column metadata with columns wrapped in a "columns" key. | ||
| """ | ||
| if UC_TABLE_PATTERN.match(input_config.location): |
There was a problem hiding this comment.
Routing is by omission: only UC_TABLE_PATTERN is checked, and everything else (two-level HMS name like default.users, a bare view name, or a malformed location) silently goes to the Spark path. For a caller expecting Spark-free operation, a two-level table unexpectedly spins up a session; a truly invalid location defers to read_input_data's generic InvalidConfigError rather than being validated here. Consider explicitly classifying the location (UC 3-level → SDK; storage path / 2-level table → Spark; otherwise raise a clear error) instead of an implicit else-fallback.
| spark: Optional Spark session. If None, a new session is created. | ||
| detector: Optional primary key detector. If None, one is created using *spark*. | ||
| """ | ||
| self.spark = SparkSession.builder.getOrCreate() if spark is None else spark |
There was a problem hiding this comment.
__init__ eagerly calls SparkSession.builder.getOrCreate(), so merely constructing DQLLMPrimaryKeyEngine requires Spark — inconsistent with the lazy spark property you introduced for DQGenerator in this same PR. Since the detector can be injected (and the metadata classification path may not need a live session), consider making Spark lazy here too, so construction doesn't force a session before any data scan actually happens.
| detector: Optional primary key detector. If None, one is created using *spark*. | ||
| """ | ||
| self.spark = SparkSession.builder.getOrCreate() if spark is None else spark | ||
| self._configurator = LLMModelConfigurator(model_config) |
There was a problem hiding this comment.
Minor (DRY): self._configurator = LLMModelConfigurator(model_config) plus the with dspy.settings.context(lm=self._configurator.create_lm()) wrapper are now duplicated between DQLLMEngine and DQLLMPrimaryKeyEngine. The split was to isolate Spark, but the LM-configuration mechanics are copy-pasted; a small shared base/mixin would keep token-handling changes from having to be applied in two places.
|
|
||
| assert result["success"] is False | ||
| assert result["table"] == "catalog.schema.orders" | ||
| assert result["final_status"] == "metadata_error" |
There was a problem hiding this comment.
This asserts result['final_status'] == 'metadata_error', which couples the test to LLMPrimaryKeyDetector's internal failure taxonomy rather than the engine's documented contract (return a failed result with success/table, don't raise). If the detector's error classification is refactored, this breaks even though the engine behavior is unchanged. Consider asserting on the documented keys (success is False, table == ...) and dropping the internal final_status check.
mwojtyczka
left a comment
There was a problem hiding this comment.
Generally looking good, left some comments
Changes
This PR decouples AI-assisted rule generation from Spark to allow use without an active Spark session. A small
get_table_column_metadatagets column metadata for any Unity Catalog tables using the Databricks SDK.Spark sessions are lazily created when required (e.g. for reading data from file paths or inferring primary keys using AI).
Linked issues
Resolves #1095
Tests
Documentation and Demos