-
Notifications
You must be signed in to change notification settings - Fork 302
fix(dbt): preserve aggregate expression semantics #437
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -16,17 +16,18 @@ | |||||||||||||||||
| # under the License. | ||||||||||||||||||
|
|
||||||||||||||||||
| from dataclasses import dataclass | ||||||||||||||||||
| from typing import List, Optional, Set | ||||||||||||||||||
| from typing import List, Optional, Set, Tuple | ||||||||||||||||||
|
|
||||||||||||||||||
| from ossie import ( | ||||||||||||||||||
| OssieDataset, | ||||||||||||||||||
| OssieDialect, | ||||||||||||||||||
| OssieDialectExpression, | ||||||||||||||||||
| OssieDocument, | ||||||||||||||||||
| OssieExpression, | ||||||||||||||||||
| OssieField, | ||||||||||||||||||
| OssieSemanticModel, | ||||||||||||||||||
| ) | ||||||||||||||||||
| from ossie_dbt.converter_issues import ConverterResult | ||||||||||||||||||
| from ossie_dbt.converter_issues import ConverterIssue, ConverterIssueType, ConverterResult | ||||||||||||||||||
| from ossie_dbt.expression_utils import ( | ||||||||||||||||||
| _extract_agg_info, | ||||||||||||||||||
| _get_dataset_qualifier, | ||||||||||||||||||
|
|
@@ -59,7 +60,6 @@ | |||||||||||||||||
| PydanticSemanticModel, | ||||||||||||||||||
| ) | ||||||||||||||||||
| from metricflow_semantic_interfaces.type_enums import ( | ||||||||||||||||||
| AggregationType, | ||||||||||||||||||
| DimensionType, | ||||||||||||||||||
| EntityType, | ||||||||||||||||||
| MetricType, | ||||||||||||||||||
|
|
@@ -91,27 +91,27 @@ class OssieToMSIConverter: | |||||||||||||||||
| - single-agg patterns (`SUM(col)`, `COUNT(DISTINCT col)`, …) → SIMPLE | ||||||||||||||||||
| metric with `metric_aggregation_params` (no measure reference needed) | ||||||||||||||||||
| - `(expr_a) / (expr_b)` → RATIO (with auto-generated sub-metrics) | ||||||||||||||||||
| - anything else → SIMPLE with the raw expression stored in `expr` | ||||||||||||||||||
| - expressions that cannot be represented without changing their | ||||||||||||||||||
| aggregation semantics are dropped with a ConverterIssue | ||||||||||||||||||
| """ | ||||||||||||||||||
|
|
||||||||||||||||||
| def __init__(self, dialect: OssieDialect = OssieDialect.ANSI_SQL) -> None: | ||||||||||||||||||
| self._dialect = dialect | ||||||||||||||||||
|
|
||||||||||||||||||
| def convert(self, document: OssieDocument) -> ConverterResult[PydanticSemanticManifest]: | ||||||||||||||||||
| semantic_models: List[PydanticSemanticModel] = [] | ||||||||||||||||||
| metrics: List[PydanticMetric] = [] | ||||||||||||||||||
|
|
||||||||||||||||||
| for dataset in document.datasets: | ||||||||||||||||||
| semantic_models.append(self._convert_dataset(dataset, document)) | ||||||||||||||||||
| metrics.extend(self._convert_metrics(document)) | ||||||||||||||||||
| metrics, issues = self._convert_metrics(document) | ||||||||||||||||||
|
|
||||||||||||||||||
| return ConverterResult( | ||||||||||||||||||
| output=PydanticSemanticManifest( | ||||||||||||||||||
| semantic_models=semantic_models, | ||||||||||||||||||
| metrics=metrics, | ||||||||||||||||||
| project_configuration=PydanticProjectConfiguration(), | ||||||||||||||||||
| ), | ||||||||||||||||||
| issues=[], | ||||||||||||||||||
| issues=issues, | ||||||||||||||||||
| ) | ||||||||||||||||||
|
|
||||||||||||||||||
| # ------------------------------------------------------------------ | ||||||||||||||||||
|
|
@@ -267,20 +267,41 @@ def _classify_field( | |||||||||||||||||
| # Metric conversion | ||||||||||||||||||
| # ------------------------------------------------------------------ | ||||||||||||||||||
|
|
||||||||||||||||||
| def _convert_metrics(self, ossie_sm: OssieSemanticModel) -> List[PydanticMetric]: | ||||||||||||||||||
| def _convert_metrics( | ||||||||||||||||||
| self, ossie_sm: OssieSemanticModel | ||||||||||||||||||
| ) -> Tuple[List[PydanticMetric], List[ConverterIssue]]: | ||||||||||||||||||
| metrics: List[PydanticMetric] = [] | ||||||||||||||||||
| issues: List[ConverterIssue] = [] | ||||||||||||||||||
| for metric in ossie_sm.metrics or []: | ||||||||||||||||||
| expr_str = self._get_expression(metric.expression) | ||||||||||||||||||
| metrics.extend(self._convert_metric(metric.name, expr_str, metric.description, ossie_sm.datasets)) | ||||||||||||||||||
| return metrics | ||||||||||||||||||
| dialect_expr = self._get_dialect_expression(metric.expression) | ||||||||||||||||||
| converted = ( | ||||||||||||||||||
| self._convert_metric( | ||||||||||||||||||
| metric.name, | ||||||||||||||||||
| dialect_expr.expression, | ||||||||||||||||||
| metric.description, | ||||||||||||||||||
| ossie_sm.datasets, | ||||||||||||||||||
| ) | ||||||||||||||||||
| if dialect_expr is not None and dialect_expr.dialect in self._SQL_DIALECTS | ||||||||||||||||||
| else None | ||||||||||||||||||
| ) | ||||||||||||||||||
| if converted is None: | ||||||||||||||||||
| issues.append( | ||||||||||||||||||
| ConverterIssue( | ||||||||||||||||||
| issue_type=ConverterIssueType.UNSUPPORTED_METRIC_EXPRESSION, | ||||||||||||||||||
| element_name=metric.name, | ||||||||||||||||||
| ) | ||||||||||||||||||
| ) | ||||||||||||||||||
| else: | ||||||||||||||||||
| metrics.extend(converted) | ||||||||||||||||||
| return metrics, issues | ||||||||||||||||||
|
|
||||||||||||||||||
| def _convert_metric( | ||||||||||||||||||
| self, | ||||||||||||||||||
| name: str, | ||||||||||||||||||
| expr_str: str, | ||||||||||||||||||
| description: Optional[str], | ||||||||||||||||||
| datasets: List[OssieDataset], | ||||||||||||||||||
| ) -> List[PydanticMetric]: | ||||||||||||||||||
| ) -> Optional[List[PydanticMetric]]: | ||||||||||||||||||
| """Return one or more PydanticMetric objects for the given Ossie expression. | ||||||||||||||||||
|
|
||||||||||||||||||
| Simple metrics use `metric_aggregation_params` to store aggregation type | ||||||||||||||||||
|
|
@@ -331,6 +352,8 @@ def _convert_metric( | |||||||||||||||||
| den_name = f"{name}__denominator" | ||||||||||||||||||
| num_metrics = self._convert_metric(num_name, num_expr, None, datasets) | ||||||||||||||||||
| den_metrics = self._convert_metric(den_name, den_expr, None, datasets) | ||||||||||||||||||
| if num_metrics is None or den_metrics is None: | ||||||||||||||||||
| return None | ||||||||||||||||||
| ratio_metric = PydanticMetric( | ||||||||||||||||||
| name=name, | ||||||||||||||||||
| description=description, | ||||||||||||||||||
|
|
@@ -345,30 +368,7 @@ def _convert_metric( | |||||||||||||||||
| ) | ||||||||||||||||||
| return [*num_metrics, *den_metrics, ratio_metric] | ||||||||||||||||||
|
|
||||||||||||||||||
| # --- Fallback: complex expression that can't be decomposed --- | ||||||||||||||||||
| # Store the raw expression in `expr` with a best-guess aggregation type. | ||||||||||||||||||
| # The caller is responsible for reviewing and correcting these metrics. | ||||||||||||||||||
| fallback_dataset = datasets[0].name if datasets else "" | ||||||||||||||||||
| return [ | ||||||||||||||||||
| PydanticMetric( | ||||||||||||||||||
| name=name, | ||||||||||||||||||
| description=description, | ||||||||||||||||||
| type=MetricType.SIMPLE, | ||||||||||||||||||
| type_params=PydanticMetricTypeParams( | ||||||||||||||||||
| expr=expr_str, | ||||||||||||||||||
| metric_aggregation_params=PydanticMetricAggregationParams( | ||||||||||||||||||
| semantic_model=fallback_dataset, | ||||||||||||||||||
| agg=AggregationType.SUM, | ||||||||||||||||||
| agg_params=None, | ||||||||||||||||||
| agg_time_dimension=None, | ||||||||||||||||||
| non_additive_dimension=None, | ||||||||||||||||||
| ), | ||||||||||||||||||
| ), | ||||||||||||||||||
| filter=None, | ||||||||||||||||||
| metadata=None, | ||||||||||||||||||
| config=None, | ||||||||||||||||||
| ) | ||||||||||||||||||
| ] | ||||||||||||||||||
| return None | ||||||||||||||||||
|
|
||||||||||||||||||
| # ------------------------------------------------------------------ | ||||||||||||||||||
| # Helpers | ||||||||||||||||||
|
|
@@ -406,10 +406,24 @@ def _find_dataset_for_col( | |||||||||||||||||
|
|
||||||||||||||||||
| def _get_expression(self, ossie_expr: OssieExpression) -> str: | ||||||||||||||||||
| """Return the expression string for the preferred dialect (fallback: first available).""" | ||||||||||||||||||
| dialect_expr = self._get_dialect_expression(ossie_expr) | ||||||||||||||||||
| return dialect_expr.expression if dialect_expr is not None else "" | ||||||||||||||||||
|
|
||||||||||||||||||
| def _get_dialect_expression( | ||||||||||||||||||
| self, ossie_expr: OssieExpression | ||||||||||||||||||
| ) -> Optional[OssieDialectExpression]: | ||||||||||||||||||
| """Return the expression for the preferred dialect (fallback: first available).""" | ||||||||||||||||||
| for dialect_expr in ossie_expr.dialects: | ||||||||||||||||||
| if dialect_expr.dialect is self._dialect: | ||||||||||||||||||
| return dialect_expr.expression | ||||||||||||||||||
| return ossie_expr.dialects[0].expression if ossie_expr.dialects else "" | ||||||||||||||||||
| return dialect_expr | ||||||||||||||||||
| return ossie_expr.dialects[0] if ossie_expr.dialects else None | ||||||||||||||||||
|
|
||||||||||||||||||
| _SQL_DIALECTS = { | ||||||||||||||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
On metrics:
- name: revenue
expression:
dialects:
- dialect: OSSIE_SQL_2026
expression: SUM(orders.amount)
The non-SQL dialects ( Suggested fix is one line:
Suggested change
Unrelated to this point, but worth flagging since it is the same function: #464 changes |
||||||||||||||||||
| OssieDialect.ANSI_SQL, | ||||||||||||||||||
| OssieDialect.BIGQUERY, | ||||||||||||||||||
| OssieDialect.DATABRICKS, | ||||||||||||||||||
| OssieDialect.SNOWFLAKE, | ||||||||||||||||||
| } | ||||||||||||||||||
|
|
||||||||||||||||||
| @staticmethod | ||||||||||||||||||
| def _parse_source(source: str) -> PydanticNodeRelation: | ||||||||||||||||||
|
|
||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This classifies "scalar" by blocklisting
exp.AggFunc/exp.Windowrather than allowlisting known-safe node types. Any aggregate sqlglot doesn't specifically model (LISTAGG, vendor/custom UDAFs) parses as genericexp.Anonymousand gets misclassified as scalar, which silently reintroduces the double-aggregation bug this PR is fixing.Maybe also worth covering scalar subqueries (e.g.
SUM((SELECT x FROM other_table LIMIT 1))) pass this check too and get embedded verbatim asexpr.I suggest inverting to an explicit allowlist of node types known to be safe as a scalar expression body, rather than trying to enumerate everything unsafe.