diff --git a/CHANGELOG.md b/CHANGELOG.md index 0382be37e..cbb6853a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ - Add opt-in POST_PARSE invocation telemetry for eligible commands via `connection_parameters.enable_dbt_telemetry` ([#1647](https://github.com/databricks/dbt-databricks/pull/1647)) - Add POST_RUN outcome telemetry to the opt-in invocation telemetry path ([#1648](https://github.com/databricks/dbt-databricks/pull/1648)) +- Add aggregate model configuration statistics to the opt-in POST_PARSE telemetry event ([#1665](https://github.com/databricks/dbt-databricks/pull/1665)) ### Fixes diff --git a/dbt/adapters/databricks/telemetry/builder.py b/dbt/adapters/databricks/telemetry/builder.py index 022b2e496..af3ca41f3 100644 --- a/dbt/adapters/databricks/telemetry/builder.py +++ b/dbt/adapters/databricks/telemetry/builder.py @@ -1,3 +1,5 @@ +from collections import Counter +from dataclasses import dataclass, field from importlib.metadata import version as _pkg_version from typing import Any, Callable, Optional @@ -32,6 +34,121 @@ "run_operation": models.DbtCommand.RUN_OPERATION, } +_MATERIALIZATION_MAP = { + "table": models.Materialization.TABLE, + "view": models.Materialization.VIEW, + "incremental": models.Materialization.INCREMENTAL, + "ephemeral": models.Materialization.EPHEMERAL, + "materialized_view": models.Materialization.MATERIALIZED_VIEW, + "streaming_table": models.Materialization.STREAMING_TABLE, + "metric_view": models.Materialization.METRIC_VIEW, +} + +_LANGUAGE_MAP = { + "sql": models.Language.SQL, + "python": models.Language.PYTHON, +} + +_INCREMENTAL_STRATEGY_MAP = { + "merge": models.IncrementalStrategy.MERGE, + "append": models.IncrementalStrategy.APPEND, + "delete+insert": models.IncrementalStrategy.DELETE_INSERT, + "delete_insert": models.IncrementalStrategy.DELETE_INSERT, + "insert_overwrite": models.IncrementalStrategy.INSERT_OVERWRITE, + "replace_where": models.IncrementalStrategy.REPLACE_WHERE, + "microbatch": models.IncrementalStrategy.MICROBATCH, +} + +_PYTHON_SUBMISSION_METHOD_MAP = { + "serverless_cluster": models.PythonSubmissionMethod.SERVERLESS_CLUSTER, + "job_cluster": models.PythonSubmissionMethod.JOB_CLUSTER, + "all_purpose_cluster": models.PythonSubmissionMethod.ALL_PURPOSE_CLUSTER, + "workflow_job": models.PythonSubmissionMethod.WORKFLOW_JOB, +} + +_CATALOG_TYPE_MAP = { + "unity": models.CatalogType.UNITY_CATALOG, + "unity_catalog": models.CatalogType.UNITY_CATALOG, + "hive_metastore": models.CatalogType.HIVE_METASTORE, +} + +_FILE_FORMAT_MAP = { + "delta": models.EffectiveStorageFormat.DELTA, + "parquet": models.EffectiveStorageFormat.PARQUET, + "hudi": models.EffectiveStorageFormat.HUDI, +} + +_CONSTRAINT_CONFIG_MAP = { + "not_null": models.ModelConfig.NOT_NULL_CONSTRAINT, + "check": models.ModelConfig.CHECK_CONSTRAINT, + "primary_key": models.ModelConfig.PRIMARY_KEY_CONSTRAINT, + "foreign_key": models.ModelConfig.FOREIGN_KEY_CONSTRAINT, + "custom": models.ModelConfig.CUSTOM_CONSTRAINT, +} + +_STORAGE_FORMAT_MATERIALIZATIONS = { + models.Materialization.TABLE, + models.Materialization.INCREMENTAL, +} + +_HMS_CATALOG_NAMES = {"hive_metastore"} + + +@dataclass +class _ModelConfigAccumulator: + scope: models.ModelConfigScope + model_count: int = 0 + materializations: Counter = field(default_factory=Counter) + languages: Counter = field(default_factory=Counter) + incremental_model_count: int = 0 + incremental_strategies: Counter = field(default_factory=Counter) + storage_formats: Counter = field(default_factory=Counter) + catalog_types: Counter = field(default_factory=Counter) + compute_types: Counter = field(default_factory=Counter) + python_model_count: int = 0 + python_submission_methods: Counter = field(default_factory=Counter) + config_usage: Counter = field(default_factory=Counter) + + def to_model(self) -> models.ModelConfigStats: + return models.ModelConfigStats( + scope=self.scope, + model_count=self.model_count, + materialization_counts=_count_rows( + self.materializations, models.Materialization, models.MaterializationCount + ), + language_counts=_count_rows(self.languages, models.Language, models.LanguageCount), + incremental_model_stats=models.IncrementalModelStats( + model_count=self.incremental_model_count, + strategy_counts=_count_rows( + self.incremental_strategies, + models.IncrementalStrategy, + models.IncrementalStrategyCount, + ), + ), + effective_storage_format_counts=_count_rows( + self.storage_formats, + models.EffectiveStorageFormat, + models.EffectiveStorageFormatCount, + ), + catalog_type_counts=_count_rows( + self.catalog_types, models.CatalogType, models.CatalogTypeCount + ), + effective_compute_type_counts=_count_rows( + self.compute_types, models.ComputeType, models.ComputeTypeCount + ), + python_model_stats=models.PythonModelStats( + model_count=self.python_model_count, + submission_method_counts=_count_rows( + self.python_submission_methods, + models.PythonSubmissionMethod, + models.PythonSubmissionMethodCount, + ), + ), + config_usage=_count_rows( + self.config_usage, models.ModelConfig, models.ModelConfigUsage + ), + ) + def classify_compute_type(http_path: Optional[str]) -> models.ComputeType: if not http_path: @@ -148,6 +265,245 @@ def aggregate_manifest(manifest: Any) -> models.ManifestStats: return stats +def _count_rows(counts: Counter, enum_type: Any, row_type: type) -> list: + return [row_type(value, counts[value]) for value in enum_type if counts[value]] + + +def _value(obj: Any, name: str, default: Any = None) -> Any: + if obj is None: + return default + if isinstance(obj, dict): + return obj.get(name, default) + getter = getattr(obj, "get", None) + if getter is not None: + try: + return getter(name, default) + except TypeError: + value = getter(name) + return default if value is None else value + return getattr(obj, name, default) + + +def _normalized(value: Any) -> str: + raw = getattr(value, "value", value) + return str(raw or "").strip().lower().replace("-", "_").replace(" ", "_") + + +def _enabled(value: Any) -> bool: + if isinstance(value, str): + return value.strip().lower() == "true" + return bool(value) + + +def _materialization(config: Any) -> models.Materialization: + value = _normalized(_value(config, "materialized")) + return _MATERIALIZATION_MAP.get(value, models.Materialization.OTHER) + + +def _language(node: Any) -> models.Language: + value = _normalized(getattr(node, "language", "sql")) + return _LANGUAGE_MAP.get(value, models.Language.OTHER) + + +def _incremental_strategy(config: Any) -> models.IncrementalStrategy: + value = _normalized(_value(config, "incremental_strategy") or "merge") + return _INCREMENTAL_STRATEGY_MAP.get(value, models.IncrementalStrategy.OTHER) + + +def _python_submission_method(config: Any) -> models.PythonSubmissionMethod: + value = _normalized(_value(config, "submission_method") or "all_purpose_cluster") + return _PYTHON_SUBMISSION_METHOD_MAP.get(value, models.PythonSubmissionMethod.OTHER) + + +def _physical_catalog_name(catalog_relation: Any, node: Any) -> str: + if catalog_relation is not None: + for attr in ("catalog_database", "catalog_name"): + value = _normalized(getattr(catalog_relation, attr, None)) + if value: + return value + return _normalized(getattr(node, "database", None)) + + +def _catalog_type(catalog_relation: Any, node: Any) -> models.CatalogType: + # Physical hive_metastore is HMS even when the Unity integration supplies + # catalog_type="unity", including the v2 catalog_database override. + if _physical_catalog_name(catalog_relation, node) in _HMS_CATALOG_NAMES: + return models.CatalogType.HIVE_METASTORE + if catalog_relation is None: + return models.CatalogType.TYPE_UNSPECIFIED + value = _normalized(getattr(catalog_relation, "catalog_type", None)) + return _CATALOG_TYPE_MAP.get(value, models.CatalogType.OTHER) + + +def _storage_format( + catalog_relation: Any, use_managed_iceberg: bool +) -> models.EffectiveStorageFormat: + if catalog_relation is None: + return models.EffectiveStorageFormat.TYPE_UNSPECIFIED + if _normalized(getattr(catalog_relation, "table_format", None)) == "iceberg": + if use_managed_iceberg: + return models.EffectiveStorageFormat.MANAGED_ICEBERG + return models.EffectiveStorageFormat.UNIFORM_ICEBERG + file_format = _normalized(getattr(catalog_relation, "file_format", None)) + if not file_format: + return models.EffectiveStorageFormat.TYPE_UNSPECIFIED + return _FILE_FORMAT_MAP.get(file_format, models.EffectiveStorageFormat.OTHER) + + +def _compute_type(config: Any, creds: DatabricksCredentials) -> models.ComputeType: + compute_name = _value(config, "databricks_compute") + if not compute_name: + return classify_compute_type(getattr(creds, "http_path", None)) + compute = getattr(creds, "compute", None) or {} + compute_config = compute.get(compute_name) + return classify_compute_type(_value(compute_config, "http_path")) + + +def _columns(node: Any) -> list[Any]: + columns = getattr(node, "columns", None) or {} + return list(columns.values()) if isinstance(columns, dict) else list(columns) + + +def _column_extra(column: Any) -> dict: + extra = _value(column, "_extra", {}) + return extra if isinstance(extra, dict) else {} + + +def _constraint_name(constraint: Any) -> str: + return _normalized(_value(constraint, "type")) + + +def _constraint_type_key(constraint: Any) -> str: + if isinstance(constraint, str): + return _normalized(constraint) + name = _constraint_name(constraint) + if name: + return name + if _value(constraint, "name") or _value(constraint, "condition"): + return "check" + return "" + + +def _constraint_type_keys(constraints: Any) -> set[str]: + if not isinstance(constraints, (list, tuple)): + return set() + return {key for constraint in constraints if (key := _constraint_type_key(constraint))} + + +def _constraint_configs(node: Any) -> set[models.ModelConfig]: + constraint_names = _constraint_type_keys(getattr(node, "constraints", None)) + meta = getattr(node, "meta", None) or {} + constraint_names.update(_constraint_type_keys(_value(meta, "constraints"))) + for column in _columns(node): + constraint_names.update(_constraint_type_keys(_value(column, "constraints", []))) + legacy_column = _value(_value(column, "meta", {}), "constraint") + if legacy_column: + constraint_names.update(_constraint_type_keys([legacy_column])) + return { + model_config + for name in constraint_names + if (model_config := _CONSTRAINT_CONFIG_MAP.get(name)) is not None + } + + +def _config_usage(node: Any, config: Any) -> set[models.ModelConfig]: + """POST_PARSE measures resolved configuration intent. + + Declarations are counted independently of runtime applicability. + """ + usage = _constraint_configs(node) + auto_liquid_cluster = _enabled(_value(config, "auto_liquid_cluster")) + if _value(config, "liquid_clustered_by") or auto_liquid_cluster: + usage.add(models.ModelConfig.LIQUID_CLUSTERING) + if auto_liquid_cluster: + usage.add(models.ModelConfig.AUTO_LIQUID_CLUSTERING) + if _value(config, "zorder"): + usage.add(models.ModelConfig.ZORDER) + if _value(config, "databricks_tags"): + usage.add(models.ModelConfig.DATABRICKS_RELATION_TAGS) + columns = _columns(node) + if any(_column_extra(column).get("databricks_tags") for column in columns): + usage.add(models.ModelConfig.COLUMN_TAGS) + if any(_column_extra(column).get("column_mask") for column in columns): + usage.add(models.ModelConfig.COLUMN_MASKS) + if _value(config, "row_filter"): + usage.add(models.ModelConfig.ROW_FILTER) + if _value(config, "databricks_compute"): + usage.add(models.ModelConfig.NAMED_COMPUTE_ROUTING) + if _enabled(_value(config, "merge_with_schema_evolution")): + usage.add(models.ModelConfig.MERGE_SCHEMA_EVOLUTION) + if _value(config, "not_matched_by_source_action"): + usage.add(models.ModelConfig.MERGE_NOT_MATCHED_BY_SOURCE) + return usage + + +def _catalog_relation(node: Any, builder: Callable[[Any], Any]) -> Any: + try: + return builder(node) + except Exception: + return None + + +def aggregate_model_configs( + manifest: Any, + creds: DatabricksCredentials, + behavior_flag: Callable[[str], bool], + catalog_relation_builder: Callable[[Any], Any], +) -> list[models.ModelConfigStats]: + root_project = getattr(getattr(manifest, "metadata", None), "project_name", None) + accumulators = { + models.ModelConfigScope.ROOT_PROJECT: _ModelConfigAccumulator( + models.ModelConfigScope.ROOT_PROJECT + ), + models.ModelConfigScope.INSTALLED_PACKAGES: _ModelConfigAccumulator( + models.ModelConfigScope.INSTALLED_PACKAGES + ), + } + use_managed_iceberg = bool(behavior_flag("use_managed_iceberg")) + for node in (getattr(manifest, "nodes", None) or {}).values(): + if _resource_type(node) != "model": + continue + scope = ( + models.ModelConfigScope.ROOT_PROJECT + if getattr(node, "package_name", None) == root_project + else models.ModelConfigScope.INSTALLED_PACKAGES + ) + acc = accumulators[scope] + config = getattr(node, "config", None) + materialization = _materialization(config) + language = _language(node) + + acc.model_count += 1 + acc.materializations[materialization] += 1 + acc.languages[language] += 1 + relation = ( + _catalog_relation(node, catalog_relation_builder) + if materialization != models.Materialization.EPHEMERAL + else None + ) + acc.config_usage.update(_config_usage(node, config)) + + if materialization == models.Materialization.INCREMENTAL: + strategy = _incremental_strategy(config) + acc.incremental_model_count += 1 + acc.incremental_strategies[strategy] += 1 + + if language == models.Language.PYTHON: + acc.python_model_count += 1 + acc.python_submission_methods[_python_submission_method(config)] += 1 + + if materialization != models.Materialization.EPHEMERAL: + acc.catalog_types[_catalog_type(relation, node)] += 1 + acc.compute_types[_compute_type(config, creds)] += 1 + if materialization in _STORAGE_FORMAT_MATERIALIZATIONS: + acc.storage_formats[_storage_format(relation, use_managed_iceberg)] += 1 + + return [ + accumulators[models.ModelConfigScope.ROOT_PROJECT].to_model(), + accumulators[models.ModelConfigScope.INSTALLED_PACKAGES].to_model(), + ] + + def ephemeral_resource_ids(manifest: Any) -> set[str]: """Return ephemeral IDs for local counting only.""" result = set() @@ -223,6 +579,7 @@ def build_post_parse_log( config: Any, creds: DatabricksCredentials, behavior_flag: Callable[[str], bool], + catalog_relation_builder: Callable[[Any], Any], ) -> models.TelemetryLog: invocation_id = _invocation_id(manifest) payload = models.PostParsePayload( @@ -230,6 +587,12 @@ def build_post_parse_log( manifest_stats=aggregate_manifest(manifest), connection_config=build_connection_config(creds), project_config=build_project_config(behavior_flag), + model_config_stats=aggregate_model_configs( + manifest, + creds, + behavior_flag, + catalog_relation_builder, + ), ) return models.TelemetryLog( invocation_id=invocation_id, diff --git a/dbt/adapters/databricks/telemetry/hooks.py b/dbt/adapters/databricks/telemetry/hooks.py index e711967fc..7ad1c9ff7 100644 --- a/dbt/adapters/databricks/telemetry/hooks.py +++ b/dbt/adapters/databricks/telemetry/hooks.py @@ -68,6 +68,7 @@ def on_post_parse(adapter: Any, manifest: Any) -> None: config=config, creds=creds, behavior_flag=adapter.get_behavior_flag_no_warn, + catalog_relation_builder=adapter.build_catalog_relation, ) if not log.invocation_id: return diff --git a/dbt/adapters/databricks/telemetry/models.py b/dbt/adapters/databricks/telemetry/models.py index f84549c00..c1ab7f398 100644 --- a/dbt/adapters/databricks/telemetry/models.py +++ b/dbt/adapters/databricks/telemetry/models.py @@ -82,6 +82,87 @@ class TerminationReason(Enum): INTERNAL_ERROR = "INTERNAL_ERROR" +class ModelConfigScope(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + ROOT_PROJECT = "ROOT_PROJECT" + INSTALLED_PACKAGES = "INSTALLED_PACKAGES" + + +class Materialization(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + TABLE = "TABLE" + VIEW = "VIEW" + INCREMENTAL = "INCREMENTAL" + EPHEMERAL = "EPHEMERAL" + MATERIALIZED_VIEW = "MATERIALIZED_VIEW" + STREAMING_TABLE = "STREAMING_TABLE" + METRIC_VIEW = "METRIC_VIEW" + OTHER = "OTHER" + + +class Language(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + SQL = "SQL" + PYTHON = "PYTHON" + OTHER = "OTHER" + + +class IncrementalStrategy(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + MERGE = "MERGE" + APPEND = "APPEND" + DELETE_INSERT = "DELETE_INSERT" + INSERT_OVERWRITE = "INSERT_OVERWRITE" + REPLACE_WHERE = "REPLACE_WHERE" + MICROBATCH = "MICROBATCH" + OTHER = "OTHER" + + +class EffectiveStorageFormat(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + DELTA = "DELTA" + MANAGED_ICEBERG = "MANAGED_ICEBERG" + UNIFORM_ICEBERG = "UNIFORM_ICEBERG" + PARQUET = "PARQUET" + HUDI = "HUDI" + OTHER = "OTHER" + + +class CatalogType(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + UNITY_CATALOG = "UNITY_CATALOG" + HIVE_METASTORE = "HIVE_METASTORE" + OTHER = "OTHER" + + +class PythonSubmissionMethod(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + SERVERLESS_CLUSTER = "SERVERLESS_CLUSTER" + JOB_CLUSTER = "JOB_CLUSTER" + ALL_PURPOSE_CLUSTER = "ALL_PURPOSE_CLUSTER" + WORKFLOW_JOB = "WORKFLOW_JOB" + OTHER = "OTHER" + + +class ModelConfig(Enum): + TYPE_UNSPECIFIED = "TYPE_UNSPECIFIED" + LIQUID_CLUSTERING = "LIQUID_CLUSTERING" + AUTO_LIQUID_CLUSTERING = "AUTO_LIQUID_CLUSTERING" + ZORDER = "ZORDER" + DATABRICKS_RELATION_TAGS = "DATABRICKS_RELATION_TAGS" + COLUMN_TAGS = "COLUMN_TAGS" + COLUMN_MASKS = "COLUMN_MASKS" + ROW_FILTER = "ROW_FILTER" + NOT_NULL_CONSTRAINT = "NOT_NULL_CONSTRAINT" + CHECK_CONSTRAINT = "CHECK_CONSTRAINT" + PRIMARY_KEY_CONSTRAINT = "PRIMARY_KEY_CONSTRAINT" + FOREIGN_KEY_CONSTRAINT = "FOREIGN_KEY_CONSTRAINT" + CUSTOM_CONSTRAINT = "CUSTOM_CONSTRAINT" + NAMED_COMPUTE_ROUTING = "NAMED_COMPUTE_ROUTING" + MERGE_SCHEMA_EVOLUTION = "MERGE_SCHEMA_EVOLUTION" + MERGE_NOT_MATCHED_BY_SOURCE = "MERGE_NOT_MATCHED_BY_SOURCE" + + @dataclass class ResourceCounts: model_count: int = 0 @@ -133,12 +214,87 @@ class ProjectConfig: use_describe_as_json_for_relation_metadata: bool = False +@dataclass +class ModelConfigUsage: + config: ModelConfig = ModelConfig.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class MaterializationCount: + materialization: Materialization = Materialization.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class LanguageCount: + language: Language = Language.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class IncrementalStrategyCount: + incremental_strategy: IncrementalStrategy = IncrementalStrategy.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class EffectiveStorageFormatCount: + effective_storage_format: EffectiveStorageFormat = EffectiveStorageFormat.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class CatalogTypeCount: + catalog_type: CatalogType = CatalogType.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class ComputeTypeCount: + compute_type: ComputeType = ComputeType.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class PythonSubmissionMethodCount: + submission_method: PythonSubmissionMethod = PythonSubmissionMethod.TYPE_UNSPECIFIED + count: int = 0 + + +@dataclass +class IncrementalModelStats: + model_count: int = 0 + strategy_counts: list[IncrementalStrategyCount] = field(default_factory=list) + + +@dataclass +class PythonModelStats: + model_count: int = 0 + submission_method_counts: list[PythonSubmissionMethodCount] = field(default_factory=list) + + +@dataclass +class ModelConfigStats: + scope: ModelConfigScope = ModelConfigScope.TYPE_UNSPECIFIED + model_count: int = 0 + materialization_counts: list[MaterializationCount] = field(default_factory=list) + language_counts: list[LanguageCount] = field(default_factory=list) + incremental_model_stats: IncrementalModelStats = field(default_factory=IncrementalModelStats) + effective_storage_format_counts: list[EffectiveStorageFormatCount] = field(default_factory=list) + catalog_type_counts: list[CatalogTypeCount] = field(default_factory=list) + effective_compute_type_counts: list[ComputeTypeCount] = field(default_factory=list) + python_model_stats: PythonModelStats = field(default_factory=PythonModelStats) + config_usage: list[ModelConfigUsage] = field(default_factory=list) + + @dataclass class PostParsePayload: invocation_config: InvocationConfig manifest_stats: ManifestStats connection_config: ConnectionConfig project_config: ProjectConfig + model_config_stats: list[ModelConfigStats] = field(default_factory=list) @dataclass diff --git a/tests/unit/telemetry/test_builder.py b/tests/unit/telemetry/test_builder.py index f3cb959f0..a08078a0f 100644 --- a/tests/unit/telemetry/test_builder.py +++ b/tests/unit/telemetry/test_builder.py @@ -3,6 +3,8 @@ import pytest from dbt_common.exceptions import DbtRuntimeError +from dbt.adapters.databricks import constants +from dbt.adapters.databricks.catalogs._unity import UnityCatalogIntegration from dbt.adapters.databricks.telemetry import builder, models @@ -28,46 +30,47 @@ def _node(resource_type, package_name="root", test_metadata=None): ) -class TestReportedClassifications: - @pytest.mark.parametrize( - "http_path, expected", - [ - pytest.param( - "/sql/1.0/warehouses/a?o=9", - models.ComputeType.SQL_WAREHOUSE, - id="warehouse", - ), - pytest.param( - "/sql/1.0/endpoints/a", - models.ComputeType.SQL_WAREHOUSE, - id="legacy_endpoint", - ), - pytest.param( - "/sql/protocolv1/o/1/2", - models.ComputeType.ALL_PURPOSE_CLUSTER, - id="cluster", - ), - pytest.param("/unknown", models.ComputeType.OTHER, id="other"), - pytest.param(None, models.ComputeType.TYPE_UNSPECIFIED, id="missing"), - ], +def _model( + materialized, + *, + package_name="root", + language="sql", + database="main", + columns=None, + constraints=None, + **config, +): + return SimpleNamespace( + resource_type="model", + package_name=package_name, + language=language, + database=database, + config={"materialized": materialized, **config}, + columns=columns or {}, + constraints=constraints or [], + meta={}, ) - def test_compute_type(self, http_path, expected): - assert builder.classify_compute_type(http_path) == expected + +def _unity_delta_relation(_node=None): + return SimpleNamespace(catalog_type="unity", table_format="default", file_format="delta") + + +class TestReportedClassifications: @pytest.mark.parametrize( "creds, expected", [ - pytest.param(_creds(token="dapi"), models.AuthFamily.PAT, id="pat"), + pytest.param( + _creds(token="dapi", azure_client_id="a", azure_client_secret="b"), + models.AuthFamily.PAT, + id="token_wins_over_azure_sp", + ), pytest.param( _creds(azure_client_id="a", azure_client_secret="b"), models.AuthFamily.AZURE_SERVICE_PRINCIPAL, id="azure_sp", ), - pytest.param( - _creds(auth_type="oauth"), - models.AuthFamily.OAUTH_U2M, - id="u2m", - ), + pytest.param(_creds(), models.AuthFamily.OAUTH_U2M, id="no_secret_is_u2m"), pytest.param( _creds(client_id="c", client_secret="s"), models.AuthFamily.LEGACY_CLIENT_SECRET_AMBIGUOUS, @@ -81,19 +84,12 @@ def test_auth_family(self, creds, expected): @pytest.mark.parametrize( "warn_error, options, expected", [ - pytest.param(None, None, models.WarnErrorPolicy.WARN_ERROR_DISABLED, id="disabled"), pytest.param( True, SimpleNamespace(error=[], warn=[], silence=["X"]), models.WarnErrorPolicy.WARN_ERROR_ALL, id="legacy_takes_precedence", ), - pytest.param( - False, - SimpleNamespace(error="all", warn=[], silence=[]), - models.WarnErrorPolicy.WARN_ERROR_ALL, - id="error_all", - ), pytest.param( False, SimpleNamespace(error="all", warn=["SomeWarning"], silence=[]), @@ -130,45 +126,263 @@ def test_named_compute_o_parameter_sets_spog_flag(self): class TestAggregateManifest: - def test_root_installed_and_test_kinds(self): + def test_generic_tests_are_split_from_singular(self): manifest = SimpleNamespace( - metadata=SimpleNamespace(project_name="root", invocation_id="inv-1"), + metadata=SimpleNamespace(project_name="root"), nodes={ - "m1": _node("model", "root"), - "m2": _node("model", "dep_pkg"), - "t_generic": _node("test", "root", test_metadata={"name": "not_null"}), - "t_singular": _node("test", "root"), - "op": _node("operation", "root"), + "t_generic": _node("test", test_metadata={"name": "not_null"}), + "t_singular": _node("test"), }, - sources={}, - exposures={}, - metrics={}, - saved_queries={}, - functions={}, - semantic_models={}, - unit_tests={}, ) ms = builder.aggregate_manifest(manifest) - assert ms.enabled_root_project.model_count == 1 - assert ms.enabled_installed_packages.model_count == 1 assert ms.enabled_total.generic_data_test_count == 1 assert ms.enabled_total.data_test_count == 2 - assert ms.enabled_total.other_count == 1 + + +class TestAggregateModelConfigs: + @pytest.mark.parametrize( + "use_managed, expected_format", + [ + pytest.param(False, models.EffectiveStorageFormat.UNIFORM_ICEBERG, id="uniform"), + pytest.param(True, models.EffectiveStorageFormat.MANAGED_ICEBERG, id="managed"), + ], + ) + def test_iceberg_storage_and_unresolved_named_compute(self, use_managed, expected_format): + manifest = SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={ + "iceberg": _model("table", table_format="iceberg", databricks_compute="missing") + }, + ) + relation = SimpleNamespace(table_format="iceberg", file_format="delta") + + root = builder.aggregate_model_configs( + manifest, + _creds(compute={}), + lambda flag: use_managed if flag == "use_managed_iceberg" else False, + lambda node: relation, + )[0] + + assert root.effective_storage_format_counts == [ + models.EffectiveStorageFormatCount(expected_format, 1) + ] + assert root.effective_compute_type_counts == [ + models.ComputeTypeCount(models.ComputeType.TYPE_UNSPECIFIED, 1) + ] + + def test_config_usage_unions_all_declared_constraint_sources(self): + node = _model( + "view", + columns={ + "id": { + "constraints": [{"type": "foreign_key"}], + "meta": {"constraint": "not_null"}, + } + }, + constraints=[ + {"type": "primary_key"}, + {"type": "check", "expression": "id > 0"}, + ], + ) + node.meta = {"constraints": [{"name": "legacy", "condition": "id > 0"}]} + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={"mixed": node}, + ), + _creds(), + lambda flag: False, + _unity_delta_relation, + )[0] + assert {row.config: row.count for row in root.config_usage} == { + models.ModelConfig.PRIMARY_KEY_CONSTRAINT: 1, + models.ModelConfig.FOREIGN_KEY_CONSTRAINT: 1, + models.ModelConfig.NOT_NULL_CONSTRAINT: 1, + models.ModelConfig.CHECK_CONSTRAINT: 1, + } + + def test_config_usage_ignores_non_sequence_legacy_constraint_metadata(self): + invalid = _model("table") + invalid.meta = {"constraints": True} + declared = _model("table", constraints=[{"type": "not_null"}]) + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={"invalid": invalid, "declared": declared}, + ), + _creds(), + lambda flag: False, + _unity_delta_relation, + )[0] + assert {row.config: row.count for row in root.config_usage} == { + models.ModelConfig.NOT_NULL_CONSTRAINT: 1, + } + + def test_config_usage_counts_conflicting_clustering_declarations(self): + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={ + "both": _model( + "table", + zorder=["id"], + liquid_clustered_by=["id"], + auto_liquid_cluster=True, + ) + }, + ), + _creds(), + lambda flag: False, + _unity_delta_relation, + )[0] + assert {row.config: row.count for row in root.config_usage} == { + models.ModelConfig.ZORDER: 1, + models.ModelConfig.LIQUID_CLUSTERING: 1, + models.ModelConfig.AUTO_LIQUID_CLUSTERING: 1, + } + + def test_config_usage_counts_merge_options_regardless_of_strategy(self): + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={ + "table": _model( + "table", + merge_with_schema_evolution=True, + not_matched_by_source_action="delete", + ), + "append": _model( + "incremental", + incremental_strategy="append", + merge_with_schema_evolution=True, + not_matched_by_source_action="delete", + ), + "invalid": _model( + "incremental", + not_matched_by_source_action="drop", + ), + }, + ), + _creds(), + lambda flag: False, + _unity_delta_relation, + )[0] + assert {row.config: row.count for row in root.config_usage} == { + models.ModelConfig.MERGE_SCHEMA_EVOLUTION: 2, + models.ModelConfig.MERGE_NOT_MATCHED_BY_SOURCE: 3, + } + + def test_config_usage_counts_declarations_on_any_materialization(self): + extra = { + "id": { + "_extra": {"column_mask": {"function": "mask_id"}, "databricks_tags": {"k": "v"}} + } + } + row_filter = {"function": "f", "columns": ["id"]} + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={ + "view": _model( + "view", + zorder=["id"], + liquid_clustered_by=["id"], + columns=extra, + row_filter=row_filter, + databricks_tags={"a": "b"}, + databricks_compute="cluster", + ), + "ephemeral": _model("ephemeral", databricks_compute="cluster", zorder=["id"]), + "py": _model( + "table", + language="python", + auto_liquid_cluster=True, + row_filter=row_filter, + columns=extra, + ), + }, + ), + _creds(compute={"cluster": {"http_path": "/sql/protocolv1/o/1/cluster"}}), + lambda flag: False, + _unity_delta_relation, + )[0] + assert {row.config: row.count for row in root.config_usage} == { + models.ModelConfig.ZORDER: 2, + models.ModelConfig.LIQUID_CLUSTERING: 2, + models.ModelConfig.AUTO_LIQUID_CLUSTERING: 1, + models.ModelConfig.COLUMN_MASKS: 2, + models.ModelConfig.ROW_FILTER: 2, + models.ModelConfig.DATABRICKS_RELATION_TAGS: 1, + models.ModelConfig.COLUMN_TAGS: 2, + models.ModelConfig.NAMED_COMPUTE_ROUTING: 2, + } + assert {row.compute_type: row.count for row in root.effective_compute_type_counts} == { + models.ComputeType.ALL_PURPOSE_CLUSTER: 1, + models.ComputeType.SQL_WAREHOUSE: 1, + } + + def test_physical_hive_metastore_is_classified_as_hms(self): + node = _model("table", database="hive_metastore") + node.schema = "dbt" + node.identifier = "hms" + relation = UnityCatalogIntegration(constants.DEFAULT_UNITY_CATALOG).build_relation(node) + assert relation.catalog_type == "unity" + assert relation.catalog_name == "hive_metastore" + + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={"hms": node}, + ), + _creds(), + lambda flag: False, + lambda _: relation, + )[0] + + assert root.catalog_type_counts == [ + models.CatalogTypeCount(models.CatalogType.HIVE_METASTORE, 1) + ] + + def test_v2_catalog_database_hive_metastore_is_classified_as_hms(self): + node = _model("table", database="main") + node.schema = "analytics" + node.identifier = "model_one" + integration = UnityCatalogIntegration( + SimpleNamespace( + name="v2_routed_catalog", + catalog_type="unity", + catalog_name="logical_catalog_label", + catalog_database="hive_metastore", + table_format="default", + external_volume=None, + file_format="delta", + adapter_properties={}, + ) + ) + relation = integration.build_relation(node) + assert relation.catalog_type == "unity" + assert relation.catalog_name == "logical_catalog_label" + assert relation.catalog_database == "hive_metastore" + + root = builder.aggregate_model_configs( + SimpleNamespace( + metadata=SimpleNamespace(project_name="root"), + nodes={"model": node}, + ), + _creds(), + lambda flag: False, + integration.build_relation, + )[0] + + assert root.catalog_type_counts == [ + models.CatalogTypeCount(models.CatalogType.HIVE_METASTORE, 1) + ] class TestBuildPostRunLog: @pytest.mark.parametrize( "exc_type, results, fail_fast, task_success, status, reason", [ - pytest.param( - None, - [], - False, - None, - models.InvocationStatus.SUCCESS, - models.TerminationReason.NORMAL, - id="success", - ), pytest.param( None, [("model.p.m1", "error")], diff --git a/tests/unit/telemetry/test_config.py b/tests/unit/telemetry/test_config.py index 1c35bcd31..6738e7ca0 100644 --- a/tests/unit/telemetry/test_config.py +++ b/tests/unit/telemetry/test_config.py @@ -30,11 +30,7 @@ class TestCommandEligibility: @pytest.mark.parametrize( "command, eligible", [ - ("build", True), ("run", True), - ("test", True), - ("seed", True), - ("snapshot", True), ("compile", False), ("source freshness", False), ("run-operation", False), diff --git a/tests/unit/telemetry/test_encoder.py b/tests/unit/telemetry/test_encoder.py index 397d36b65..6d9844170 100644 --- a/tests/unit/telemetry/test_encoder.py +++ b/tests/unit/telemetry/test_encoder.py @@ -63,8 +63,6 @@ def test_post_parse_envelope(self): assert "dbt_databricks_telemetry_log" in fe["entry"] assert entry["event_type"] == "POST_PARSE" assert "post_run" not in entry - assert entry["post_parse"]["invocation_config"]["dbt_command"] == "RUN" - assert entry["post_parse"]["connection_config"]["default_compute_type"] == "SQL_WAREHOUSE" @pytest.mark.parametrize( "workspace_id, expected", @@ -86,8 +84,6 @@ def test_post_run_wire_contract(self): rc = entry["post_run"]["result_counts"] assert rc["pass"] == 5 assert "pass_" not in rc - assert entry["post_run"]["run_outcome"]["invocation_status"] == "HANDLED_ERROR" - assert entry["post_run"]["results_by_resource_type"][0]["resource_type"] == "MODEL" def test_unavailable_aggregates_are_omitted(self): log = models.TelemetryLog(