From 3623eea4cd954352981d6476c85a59563a270f17 Mon Sep 17 00:00:00 2001 From: Yian Shang Date: Wed, 12 Aug 2026 08:18:50 -0700 Subject: [PATCH] Bind each metric and dimension to its column on the REST path too POST /preaggs/register took the retired four-field shape -- metrics and dimensions as lists, with the columns in separate measure_columns and dimension_columns maps -- while the kind: preagg YAML spec had moved to declaring each reference together with the physical column that holds it. One concept spelled two ways depending on which surface you reached for. The request body now takes the same two maps, and both models enforce the same three rules from models/preagg_binding: the retired fields are refused, a list where a map belongs is refused, and a key whose column is empty is refused. Each rejection says what was found, what to write instead, and shows the whole shape in the syntax that surface is authored in, so an existing file or API client can be fixed from the message. Since the references are now the map keys on both paths, register_external_preaggregations takes the maps rather than lists plus bindings, which removes the check for a bound dimension that isn't in the grain -- unrepresentable now. Binding every dimension explicitly also means every dimension column is type-checked, which is what the YAML path already did. --- .../python/tests/test_deploy.py | 10 +- .../api/preaggregations.py | 22 +- .../internal/deployment/orchestrator.py | 2 - .../internal/preaggregations.py | 61 +- .../datajunction_server/models/deployment.py | 85 +-- .../models/preagg_binding.py | 158 ++++++ .../models/preaggregation.py | 67 ++- .../tests/api/deployments_test.py | 16 +- .../tests/api/preaggregations_test.py | 232 ++++---- .../build_v3/preagg_substitution_test.py | 535 ++++++++---------- .../tests/models/deployment_test.py | 103 +--- .../tests/models/preagg_binding_test.py | 242 ++++++++ .../docs/dj-concepts/aggregate-awareness.md | 46 +- 13 files changed, 910 insertions(+), 669 deletions(-) create mode 100644 datajunction-server/datajunction_server/models/preagg_binding.py create mode 100644 datajunction-server/tests/models/preagg_binding_test.py diff --git a/datajunction-clients/python/tests/test_deploy.py b/datajunction-clients/python/tests/test_deploy.py index f45c68890..cff930177 100644 --- a/datajunction-clients/python/tests/test_deploy.py +++ b/datajunction-clients/python/tests/test_deploy.py @@ -54,12 +54,11 @@ def test_reconstruct_deployment_spec_separates_preaggregations(tmp_path): (tmp_path / "revenue_by_day.yaml").write_text( "kind: preagg\n" "name: revenue_by_day\n" - "metrics:\n - ns.revenue\n" - "dimensions:\n - ns.date.day\n" + "metrics:\n ns.revenue: revenue_sum\n" + "dimensions:\n ns.date.day: day\n" "catalog: default\n" "schema: agg\n" - "table: revenue_agg\n" - "measure_columns:\n ns.revenue: revenue_sum\n", + "table: revenue_agg\n", ) svc = DeploymentService(MagicMock()) @@ -70,7 +69,8 @@ def test_reconstruct_deployment_spec_separates_preaggregations(tmp_path): assert len(spec["preaggregations"]) == 1 preagg = spec["preaggregations"][0] assert preagg["name"] == "revenue_by_day" - assert preagg["measure_columns"] == {"ns.revenue": "revenue_sum"} + assert preagg["metrics"] == {"ns.revenue": "revenue_sum"} + assert preagg["dimensions"] == {"ns.date.day": "day"} # The discriminator is stripped before reaching the payload. assert "kind" not in preagg diff --git a/datajunction-server/datajunction_server/api/preaggregations.py b/datajunction-server/datajunction_server/api/preaggregations.py index a70f17d8d..3ed49d7ca 100644 --- a/datajunction-server/datajunction_server/api/preaggregations.py +++ b/datajunction-server/datajunction_server/api/preaggregations.py @@ -735,11 +735,13 @@ async def register_preaggregations( Register an externally-built pre-aggregation table. Unlike ``/preaggs/plan`` (where DJ generates and owns the materialization), - this adopts a table built by an external pipeline. DJ decomposes the - requested metrics into component measures, binds each measure to a physical - column via ``measure_columns``, validates them against the table, records - the pre-aggregation, and — when ``valid_through_ts`` is supplied — marks it - available so grain resolution can route queries to it. + this adopts a table built by an external pipeline. Every metric and every + dimension is declared together with the physical column of the external table + that holds it, the same way a ``kind: preagg`` YAML spec declares it. DJ + decomposes the requested metrics into component measures, validates every + declared column against the table, records the pre-aggregation, and — when + ``valid_through_ts`` is supplied — marks it available so grain resolution can + route queries to it. """ request_headers = dict(request.headers) if not query_service_client: @@ -756,10 +758,14 @@ async def register_preaggregations( # inference and row creation). The same result is handed to the registration # below: resolving it twice can pick different parents, which would authorize # one set of nodes and write another. + # The request binds each reference to its physical column inline, so the + # references are the map keys and the bindings are the maps themselves. + metrics = list(data.metrics) + dimensions = list(data.dimensions) measures_result = await build_measures_sql( session=session, - metrics=data.metrics, - dimensions=data.dimensions, + metrics=metrics, + dimensions=dimensions, dialect=Dialect.SPARK, use_materialized=False, ) @@ -775,8 +781,6 @@ async def register_preaggregations( metrics=data.metrics, dimensions=data.dimensions, table=data.table, - measure_columns=data.measure_columns, - dimension_columns=data.dimension_columns, measures_result=measures_result, ) diff --git a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py index 933e659a7..9b8e5a0ac 100644 --- a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py +++ b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py @@ -1247,8 +1247,6 @@ async def _reconcile_preaggregations(self) -> None: table=spec.table, valid_through_ts=spec.valid_through_ts, ), - measure_columns=spec.rendered_measure_columns, - dimension_columns=spec.rendered_dimension_columns, ) upserted_ids.update(preagg.id for preagg in created) for preagg in created: diff --git a/datajunction-server/datajunction_server/internal/preaggregations.py b/datajunction-server/datajunction_server/internal/preaggregations.py index 956053fe6..bfa7c8bde 100644 --- a/datajunction-server/datajunction_server/internal/preaggregations.py +++ b/datajunction-server/datajunction_server/internal/preaggregations.py @@ -169,18 +169,22 @@ async def register_external_preaggregations( request_headers: dict[str, str], *, name: str | None, - metrics: list[str], - dimensions: list[str], + metrics: dict[str, str], + dimensions: dict[str, str], table: ExternalPreAggTable, - measure_columns: dict[str, str], - dimension_columns: dict[str, str] | None = None, measures_result: GeneratedMeasuresSQL | None = None, ) -> list[PreAggregation]: """ Core logic for adopting an externally-built pre-aggregation table. - Decomposes ``metrics`` into component measures, binds each to a physical - column via ``measure_columns``, validates them against ``table``, and upserts + ``metrics`` and ``dimensions`` each map a reference to the physical column of + the external table that holds it — the shape both surfaces declare (a + ``kind: preagg`` YAML spec and the body of ``POST /preaggs/register``). The + references are the keys, so a reference can never be declared without its + binding or bound without being declared. + + Decomposes the metrics into component measures, binds each to its declared + column, validates them against ``table``, and upserts the pre-aggregation(s) marked ``EXTERNAL``. Flushes but does NOT commit — the caller owns the transaction (the endpoint commits; the deploy orchestrator commits its whole plan). Callers must ensure ``query_service_client`` is @@ -196,7 +200,7 @@ async def register_external_preaggregations( # column. Keying on the hash alone would collapse SUM(x) and MAX(x), # silently discarding one metric's declared column. measure_identity_to_column: dict[tuple[str, str], str] = {} - for metric_name, physical_column in measure_columns.items(): + for metric_name, physical_column in metrics.items(): node = await Node.get_by_name( session, metric_name, @@ -209,7 +213,9 @@ async def register_external_preaggregations( ) if not node or node.type != NodeType.METRIC: raise DJInvalidInputException( - message=f"'{metric_name}' in measure_columns is not a metric node.", + message=( + f"'{metric_name}' declared under `metrics` is not a metric node." + ), ) if not node.current.is_measure: raise DJInvalidInputException( @@ -234,12 +240,12 @@ async def register_external_preaggregations( if measures_result is None: measures_result = await build_measures_sql( session=session, - metrics=metrics, - dimensions=dimensions, + metrics=list(metrics), + dimensions=list(dimensions), dialect=Dialect.SPARK, use_materialized=False, ) - assert_dimension_refs_are_role_qualified(measures_result, dimensions) + assert_dimension_refs_are_role_qualified(measures_result, list(dimensions)) # 3. Introspect the external table and confirm the declared columns exist. catalog = await get_catalog_by_name(session=session, name=table.catalog) @@ -250,26 +256,17 @@ async def register_external_preaggregations( request_headers, catalog.engines[0] if catalog.engines else None, ) - dimension_columns = dimension_columns or {} - unknown_dims = sorted(set(dimension_columns) - set(dimensions)) - if unknown_dims: - raise DJInvalidInputException( - message=( - f"dimension_columns references {unknown_dims}, which are not in " - f"the pre-aggregation's dimensions {sorted(dimensions)}." - ), - ) table_columns_by_name = {col.name: col.type for col in table_columns} missing_columns = sorted( column - for column in (*measure_columns.values(), *dimension_columns.values()) + for column in (*metrics.values(), *dimensions.values()) if column not in table_columns_by_name ) if missing_columns: raise DJInvalidInputException( message=( - f"Columns {missing_columns} declared in measure_columns/" - f"dimension_columns were not found in table " + f"Columns {missing_columns} declared under `metrics`/`dimensions` " + f"were not found in table " f"{table.catalog}.{table.schema_}.{table.table}." ), ) @@ -312,14 +309,20 @@ async def register_external_preaggregations( for component in grain_group.components: expr_hash = compute_expression_hash(component.expression) identity = (expr_hash, component.normalized_aggregation) - if identity not in measure_identity_to_column: + # Defensive since the map form landed: the grain groups are built + # from the very metrics `metrics` binds, and each is_measure metric + # contributes exactly one component, so every component is bound. It + # stays as the clear failure a bare KeyError below would not be, in + # case a metric ever decomposes differently in a grain group than it + # does on its own. + if identity not in measure_identity_to_column: # pragma: no cover raise DJInvalidInputException( message=( f"Measure '{component.name}' " f"({component.normalized_aggregation} over " f"'{component.expression}') required by the requested " - f"metrics is not covered by measure_columns. Add the " - f"is_measure metric it corresponds to." + f"metrics is not covered by any column declared under " + f"`metrics`. Add the is_measure metric it corresponds to." ), ) physical_column = measure_identity_to_column[identity] @@ -346,13 +349,13 @@ async def register_external_preaggregations( ), ) - # Bind dimension columns to physical columns via dimension_columns, keyed - # by dimension reference. Measures carry their binding on PreAggMeasure. + # Bind each grain dimension to its declared physical column, keyed by + # dimension reference. Measures carry their binding on PreAggMeasure. columns = [] for col in grain_group.columns: source_column = None if col.semantic_type == "dimension": - source_column = dimension_columns.get(col.semantic_name) + source_column = dimensions.get(col.semantic_name) if source_column is not None: assert_column_type_compatible( subject=f"dimension '{col.semantic_name}'", diff --git a/datajunction-server/datajunction_server/models/deployment.py b/datajunction-server/datajunction_server/models/deployment.py index eafc7053e..00b034796 100644 --- a/datajunction-server/datajunction_server/models/deployment.py +++ b/datajunction-server/datajunction_server/models/deployment.py @@ -37,6 +37,10 @@ NodeType, ) from datajunction_server.models.partition import Granularity, PartitionType +from datajunction_server.models.preagg_binding import ( + SPEC_SURFACE, + validate_column_bindings, +) from datajunction_server.models.unit import ( Unit, legacy_unit_to_structured, @@ -208,6 +212,10 @@ class PreAggSpec(NamespacedSpec): separate ``measure_columns``/``dimension_columns`` maps -- is no longer accepted, and a spec still using it is rejected with a message describing what to write instead. + + ``POST /preaggs/register`` takes the same declaration in JSON, and both + surfaces enforce these rules from + :mod:`datajunction_server.models.preagg_binding` so they cannot drift apart. """ # Metric/dimension reference -> the physical column of the external table @@ -222,89 +230,30 @@ class PreAggSpec(NamespacedSpec): model_config = ConfigDict(populate_by_name=True) - # Each removed field and the map that replaced it. - _REPLACED_FIELDS: ClassVar[dict[str, str]] = { - "measure_columns": "metrics", - "dimension_columns": "dimensions", - } - - # Why a missing column can't be defaulted, per axis. - _UNBOUND_REASONS: ClassVar[dict[str, str]] = { - "metrics": ( - "A measure's DJ-side name is auto-generated with an expression-hash " - "suffix, so there is no name to fall back on." - ), - "dimensions": ( - "Write the column out even when it matches the DJ column name, so " - "the file says what the table actually holds." - ), - } - @model_validator(mode="before") @classmethod def require_column_bindings(cls, data: Any) -> Any: """ Hold the author to the map form: every metric and dimension names the physical column that holds it, and the fields that used to carry those - bindings separately are gone. + bindings separately are gone. Shared with ``POST /preaggs/register``, + which enforces the identical rules against a JSON body. """ - if not isinstance(data, dict): - return data - - name = data.get("name") - for removed, replacement in cls._REPLACED_FIELDS.items(): - if removed in data: - raise DJInvalidDeploymentConfig( - message=( - f"Pre-aggregation '{name}' declares `{removed}`, which is no " - f"longer a pre-aggregation field. Declare the physical column " - f"alongside what it holds instead, as `{replacement}: " - f"{{: }}`, and drop the `{removed}` block." - ), - ) - - for field, reason in cls._UNBOUND_REASONS.items(): - value = data.get(field) - if value is None: - continue - if not isinstance(value, dict): - raise DJInvalidDeploymentConfig( - message=( - f"Pre-aggregation '{name}' declares `{field}` as a " - f"{type(value).__name__}. `{field}` is a map from each " - f"reference to the physical column of the external table " - f"that holds it, e.g. `{field}: {{: }}`." - ), - ) - unbound = [ - reference for reference, column in value.items() if column is None - ] - if unbound: - raise DJInvalidDeploymentConfig( - message=( - f"Pre-aggregation '{name}' leaves the physical column empty " - f"under `{field}` for {unbound}. {reason}" - ), - ) - return data - - @property - def rendered_metrics(self) -> list[str]: - return [render_prefixes(metric, self.namespace) for metric in self.metrics] - - @property - def rendered_dimensions(self) -> list[str]: - return [render_prefixes(dim, self.namespace) for dim in self.dimensions] + return validate_column_bindings(data, SPEC_SURFACE) @property - def rendered_measure_columns(self) -> dict[str, str]: + def rendered_metrics(self) -> dict[str, str]: + """Each metric reference rendered against the namespace, still bound to + its physical column.""" return { render_prefixes(metric, self.namespace): column for metric, column in self.metrics.items() } @property - def rendered_dimension_columns(self) -> dict[str, str]: + def rendered_dimensions(self) -> dict[str, str]: + """Each dimension reference rendered against the namespace, still bound to + its physical column.""" return { render_prefixes(dimension, self.namespace): column for dimension, column in self.dimensions.items() diff --git a/datajunction-server/datajunction_server/models/preagg_binding.py b/datajunction-server/datajunction_server/models/preagg_binding.py new file mode 100644 index 000000000..e3b5fe073 --- /dev/null +++ b/datajunction-server/datajunction_server/models/preagg_binding.py @@ -0,0 +1,158 @@ +""" +Shared validation for the inline column bindings a pre-aggregation declares. + +Two surfaces adopt an externally-built pre-aggregation table: the ``kind: preagg`` +YAML spec (``PreAggSpec``) and the JSON body of ``POST /preaggs/register`` +(``RegisterPreAggregationsRequest``). Both spell the declaration the same way -- +every metric and every dimension is named together with the physical column of the +external table that holds it, as a map -- so both enforce the rules from here +rather than keeping a copy each. One concept, one spelling, one implementation. + +The three rules are the same on both surfaces: + +1. The retired four-field form (``metrics``/``dimensions`` as lists alongside + separate ``measure_columns``/``dimension_columns`` maps) is refused. +2. A list where a map belongs is refused. +3. A key whose column is empty is refused -- including a dimension whose physical + column happens to match its DJ column name, which is written out rather than + left blank. + +The surfaces differ only in wording: a message has to name what the author is +looking at (a spec file, which has a ``name``, versus a request body) and show the +shape in the syntax they are writing (YAML versus JSON). That is all +:class:`BindingSurface` carries. The messages are the migration path for anyone +holding an existing file or an existing API client, so each one says what was +found, what to write instead, and shows the shape concretely. +""" + +from collections.abc import Callable +from dataclasses import dataclass +from typing import Any + +from datajunction_server.errors import ( + DJException, + DJInvalidDeploymentConfig, + DJInvalidInputException, +) + +# Each retired field and the map that replaced it. +REPLACED_FIELDS: dict[str, str] = { + "measure_columns": "metrics", + "dimension_columns": "dimensions", +} + +# One entry of each map, singular, for prose. +AXIS_NOUN: dict[str, str] = { + "metrics": "metric", + "dimensions": "dimension", +} + +# Why a missing column can't be defaulted, per axis. +UNBOUND_REASONS: dict[str, str] = { + "metrics": ( + "A measure's DJ-side name is auto-generated with an expression-hash " + "suffix, so there is no name to fall back on." + ), + "dimensions": ( + "Write the column out even when it matches the DJ column name, so the " + "declaration says what the table actually holds." + ), +} + +# The same two metrics and two dimensions, spelled for each surface. A dimension +# stored under a different name (page) and one stored under the name DJ uses +# (country_iso_code) are both shown, since the second is the one people leave out. +YAML_EXAMPLE = """metrics: + default.view_secs: view_secs_sum + default.session_count: session_cnt +dimensions: + default.page_d.page_id: page + default.geo_country_d.country_iso_code: country_iso_code""" + +JSON_EXAMPLE = """{ + "metrics": { + "default.view_secs": "view_secs_sum", + "default.session_count": "session_cnt" + }, + "dimensions": { + "default.page_d.page_id": "page", + "default.geo_country_d.country_iso_code": "country_iso_code" + } +}""" + + +@dataclass(frozen=True) +class BindingSurface: + """ + How the rules below address one surface: what to call the declaration, the + example that shows it the way the author writes it, and the exception to raise. + """ + + # Built from the raw input, since a spec names itself and a request body doesn't. + subject: Callable[[dict[str, Any]], str] + example: str + exception: type[DJException] + + +SPEC_SURFACE = BindingSurface( + subject=lambda data: f"Pre-aggregation '{data.get('name')}'", + example=YAML_EXAMPLE, + exception=DJInvalidDeploymentConfig, +) + +REQUEST_SURFACE = BindingSurface( + subject=lambda data: "The POST /preaggs/register request body", + example=JSON_EXAMPLE, + exception=DJInvalidInputException, +) + + +def validate_column_bindings(data: Any, surface: BindingSurface) -> Any: + """ + Hold a pre-aggregation declaration to the map form, for use as a Pydantic + ``model_validator(mode="before")`` on either surface. Returns ``data`` + untouched when it passes; raises ``surface.exception`` otherwise. + """ + if not isinstance(data, dict): + return data + + subject = surface.subject(data) + + for removed, replacement in REPLACED_FIELDS.items(): + if removed in data: + raise surface.exception( + message=( + f"{subject} declares `{removed}`, which is no longer a " + f"pre-aggregation field. Declare each " + f"{AXIS_NOUN[replacement]} together with the physical column " + f"of the external table that holds it, under `{replacement}`, " + f"and drop the `{removed}` block. The map form looks like " + f"this:\n\n{surface.example}" + ), + ) + + for field, reason in UNBOUND_REASONS.items(): + value = data.get(field) + if value is None: + continue + if not isinstance(value, dict): + raise surface.exception( + message=( + f"{subject} declares `{field}` as a {type(value).__name__}, " + f"not a map. `{field}` maps each {AXIS_NOUN[field]} reference " + f"to the physical column of the external table that holds it. " + f"The map form looks like this:\n\n{surface.example}" + ), + ) + unbound = [ + reference for reference, column in value.items() if column is None + ] + if unbound: + raise surface.exception( + message=( + f"{subject} leaves the physical column empty under `{field}` " + f"for {unbound}. {reason} The map form looks like " + f"this:\n\n{surface.example}" + ), + ) + return data diff --git a/datajunction-server/datajunction_server/models/preaggregation.py b/datajunction-server/datajunction_server/models/preaggregation.py index b6454a232..cdc78ee57 100644 --- a/datajunction-server/datajunction_server/models/preaggregation.py +++ b/datajunction-server/datajunction_server/models/preaggregation.py @@ -5,7 +5,7 @@ from datetime import date, datetime from typing import Any -from pydantic import BaseModel, Field, field_validator +from pydantic import BaseModel, Field, field_validator, model_validator from datajunction_server.enum import StrEnum from datajunction_server.models.decompose import PreAggMeasure @@ -13,6 +13,10 @@ from datajunction_server.models.node import PartitionAvailability from datajunction_server.models.node_type import NodeNameVersion from datajunction_server.models.partition import Granularity +from datajunction_server.models.preagg_binding import ( + REQUEST_SURFACE, + validate_column_bindings, +) from datajunction_server.models.query import ColumnMetadata, V3ColumnMetadata @@ -112,9 +116,28 @@ class RegisterPreAggregationsRequest(BaseModel): Unlike /preaggs/plan (where DJ generates and owns the materialization), this adopts a table already built by an external pipeline. DJ decomposes the - metrics to determine the required measures, binds each to a physical column - via ``measure_columns``, validates them against the table, and records the - pre-aggregation so grain resolution can route queries to it. + metrics to determine the required measures, validates each declared column + against the table, and records the pre-aggregation so grain resolution can + route queries to it. + + Every metric and every dimension is declared together with the physical + column of the external table that holds it, as a map:: + + { + "metrics": {"default.view_secs": "view_secs_sum"}, + "dimensions": {"default.page_d.page_id": "page"} + } + + This is the same declaration a ``kind: preagg`` YAML spec makes, so the two + surfaces share their validation -- see + :mod:`datajunction_server.models.preagg_binding`. Both maps require a value + for every key, including a dimension whose physical column matches its DJ + column name. + + What goes under ``metrics`` are the measures the table stores. A derived + metric has no column of its own, so it is not listed and cannot be; it is + served anyway, since registration and query-time matching both work on + decomposed measure identities rather than metric names. """ name: str | None = Field( @@ -124,29 +147,31 @@ class RegisterPreAggregationsRequest(BaseModel): "reconciliation and availability-by-name callbacks." ), ) - metrics: list[str] = Field( - description="Metric node names the external table should serve", - ) - dimensions: list[str] = Field( - description="Dimension references defining the table's grain", - ) - table: ExternalPreAggTable = Field(description="The externally-built table") - measure_columns: dict[str, str] = Field( - default_factory=dict, + metrics: dict[str, str] = Field( description=( - "Maps each is_measure metric name to the physical column in the " - "external table holding its pre-aggregated value. Every component " - "measure of the requested metrics must be covered." + "Maps each is_measure metric name the external table serves to the " + "physical column holding its pre-aggregated value. A derived metric " + "has no column and is not listed; it is covered once the measures it " + "decomposes into are." ), ) - dimension_columns: dict[str, str] = Field( - default_factory=dict, + dimensions: dict[str, str] = Field( description=( - "Optional map of dimension reference to the physical column in the " - "external table backing it, when the column name differs from the " - "dimension's. Unmapped dimensions default to their DJ column name." + "Maps each dimension reference defining the table's grain to the " + "physical column in the external table backing it. Required for every " + "dimension, including one stored under the DJ column name." ), ) + table: ExternalPreAggTable = Field(description="The externally-built table") + + @model_validator(mode="before") + @classmethod + def require_column_bindings(cls, data: Any) -> Any: + """ + Hold the caller to the map form. Shared with the ``kind: preagg`` YAML + spec, which enforces the identical rules against a spec file. + """ + return validate_column_bindings(data, REQUEST_SURFACE) class UpdatePreAggregationAvailabilityRequest(BaseModel): diff --git a/datajunction-server/tests/api/deployments_test.py b/datajunction-server/tests/api/deployments_test.py index 331705b08..bc5220d09 100644 --- a/datajunction-server/tests/api/deployments_test.py +++ b/datajunction-server/tests/api/deployments_test.py @@ -1905,12 +1905,12 @@ async def test_deploy_preagg_applies_dimension_columns( default_us_state, ): """ - A pre-agg declared in a deployment spec with ``dimension_columns`` binds - each grain dimension to its physical column, and the generated measures - SQL reads that physical column (aliased back to the DJ name) instead of - the DJ dimension name. + A pre-agg declared in a deployment spec binds each grain dimension to + its physical column under ``dimensions``, and the generated measures SQL + reads that physical column (aliased back to the DJ name) instead of the + DJ dimension name. - Regression: the ``dimension_columns`` feature was only covered through the + Regression: the dimension binding was only covered through the POST /preaggs/register API path; the deployment/orchestrator path (what ``dj push`` uses) had no coverage. This exercises the two shapes that path must get right: a joined dimension's *key* satisfied by a differently @@ -7684,8 +7684,8 @@ async def test_deploy_preagg_invalid_metric_fails( default_us_states, default_us_state, ): - """A pre-agg whose measure_columns key is not a valid measure fails the - whole deploy and rolls it back -- no nodes or pre-aggs are left behind.""" + """A pre-agg whose `metrics` key is not a valid measure fails the whole + deploy and rolls it back -- no nodes or pre-aggs are left behind.""" _override_query_service( client, {"hard_hat_count": "bigint", "state_name": "string"}, @@ -8110,7 +8110,7 @@ async def test_deploy_preagg_with_dimension_columns( default_us_states, default_us_state, ): - """dimension_columns on a deployed pre-agg flows through reconcile so the + """A deployed pre-agg's dimension binding flows through reconcile so the renamed physical dimension column is read at query time.""" _override_query_service( client, diff --git a/datajunction-server/tests/api/preaggregations_test.py b/datajunction-server/tests/api/preaggregations_test.py index db1dd0bcb..c823dc56b 100644 --- a/datajunction-server/tests/api/preaggregations_test.py +++ b/datajunction-server/tests/api/preaggregations_test.py @@ -2351,14 +2351,25 @@ def _mock_query_service(client_with_build_v3, columns): ``columns`` may be a list of column names (each defaulting to a numeric type) or a mapping of column name -> SQL type string, for exercising the measure/column type-compatibility check. + + Registration binds every dimension to a physical column and type-checks the + binding, so the list form always reports the dimension columns these tests + grain on (``category``, ``status``) as strings, whether or not the caller + listed them. A test that cares about a dimension column's type passes the + mapping form instead. """ from types import SimpleNamespace - items = ( - columns.items() - if isinstance(columns, dict) - else [(name, "double") for name in columns] - ) + grain_columns = {"category": "string", "status": "string"} + if isinstance(columns, dict): + items = list(columns.items()) + else: + items = [(name, grain_columns.get(name, "double")) for name in columns] + items += [ + (name, type_str) + for name, type_str in grain_columns.items() + if name not in columns + ] table_columns = [ SimpleNamespace(name=name, type=type_str) for name, type_str in items ] @@ -2387,18 +2398,17 @@ async def test_register_external_preagg(self, client_with_build_v3: AsyncClient) response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "avg_order_value_agg", "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 201, response.text @@ -2416,20 +2426,19 @@ async def test_register_rejects_non_measure( self, client_with_build_v3: AsyncClient, ): - """A derived/ratio metric cannot be used as a measure_columns key.""" + """A derived/ratio metric cannot be declared under `metrics`.""" _mock_query_service(client_with_build_v3, ["some_col"]) try: response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": {"v3.avg_order_value": "some_col"}, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "bad_agg", }, - "measure_columns": {"v3.avg_order_value": "some_col"}, }, ) assert response.status_code == 422 @@ -2448,17 +2457,16 @@ async def test_register_rejects_missing_column( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "avg_agg", }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 422 @@ -2484,17 +2492,16 @@ async def test_register_rejects_incompatible_column_type( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "bad_types_agg", }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 422 @@ -2525,17 +2532,16 @@ async def test_register_accepts_compatible_but_different_type( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "compatible_types_agg", }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 201, response.text @@ -2557,17 +2563,16 @@ async def test_register_without_valid_through_ts_is_pending( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "agg_pending", }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 201, response.text @@ -2580,18 +2585,17 @@ async def test_register_is_idempotent(self, client_with_build_v3: AsyncClient): """Registering the same table twice updates the existing pre-agg.""" _mock_query_service(client_with_build_v3, ["revenue_total", "order_cnt"]) payload = { - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "agg_idem", "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, } try: first = await client_with_build_v3.post("/preaggs/register", json=payload) @@ -2603,29 +2607,57 @@ async def test_register_is_idempotent(self, client_with_build_v3: AsyncClient): del client_with_build_v3.app.dependency_overrides[get_query_service_client] @pytest.mark.asyncio - async def test_register_rejects_uncovered_measure( + async def test_register_covers_a_derived_metric_not_listed( self, client_with_build_v3: AsyncClient, ): - """Every component measure of the requested metrics must be covered.""" - _mock_query_service(client_with_build_v3, ["revenue_total"]) + """ + A derived metric has no column of its own, so it cannot be declared under + `metrics` — and doesn't need to be. Declaring the measures it decomposes + into is what covers it, because matching works on measure identities + rather than metric names. + """ + _mock_query_service( + client_with_build_v3, + ["revenue_total", "order_cnt", "category"], + ) try: response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + # avg_order_value is deliberately absent. + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", - "table": "agg_partial", + "table": "agg_derived", + "valid_through_ts": 1700000000, }, - # order_count's measure is missing. - "measure_columns": {"v3.total_revenue": "revenue_total"}, }, ) - assert response.status_code == 422 - assert "not covered by measure_columns" in response.text + assert response.status_code == 201, response.text + preagg = response.json()["preaggs"][0] + # Every metric that decomposes into the two stored measures, derived + # ones included -- none of which had to be listed to be served. + assert sorted(preagg["related_metrics"]) == [ + "v3.aov_growth_index", + "v3.avg_order_value", + "v3.efficiency_ratio", + "v3.mom_revenue_change", + "v3.revenue_per_customer", + "v3.revenue_per_page_view", + "v3.revenue_per_visitor", + "v3.total_revenue", + "v3.trailing_7d_revenue", + "v3.trailing_7d_revenue_inferred_dim", + "v3.trailing_wow_revenue_change", + "v3.wow_aov_change", + "v3.wow_revenue_change", + ] finally: del client_with_build_v3.app.dependency_overrides[get_query_service_client] @@ -2634,24 +2666,27 @@ async def test_register_rejects_unknown_metric( self, client_with_build_v3: AsyncClient, ): - """A measure_columns key that is not a metric node is rejected.""" + """ + A `metrics` key that names no node is rejected. The keys are the metrics + the table serves, so decomposition sees the bad name first and reports it + before any binding is looked at. + """ _mock_query_service(client_with_build_v3, ["c"]) try: response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": {"v3.does_not_exist": "c"}, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "t", }, - "measure_columns": {"v3.does_not_exist": "c"}, }, ) assert response.status_code == 422 - assert "not a metric node" in response.text + assert response.json()["message"] == "Metric not found: v3.does_not_exist" finally: del client_with_build_v3.app.dependency_overrides[get_query_service_client] @@ -2668,14 +2703,13 @@ async def test_register_requires_query_service( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": {"v3.total_revenue": "revenue_total"}, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "t", }, - "measure_columns": {"v3.total_revenue": "revenue_total"}, }, ) assert response.status_code != 201 @@ -2695,18 +2729,17 @@ async def test_register_sets_external_strategy_and_name( "/preaggs/register", json={ "name": "aov_by_category", - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "aov_agg", "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 201, response.text @@ -2721,18 +2754,17 @@ async def _register_external(self, client, table_name): response = await client.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": table_name, "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 201, response.text @@ -2804,18 +2836,17 @@ async def test_measures_are_built_once( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "built_once", "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 201, response.text @@ -2901,18 +2932,17 @@ async def test_register_denies_without_write( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.order_count": "order_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "denied_agg", "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.order_count": "order_cnt", - }, }, ) assert response.status_code == 403, response.text @@ -2965,13 +2995,12 @@ async def _register_aov( response = await client.post( "/preaggs/register", json={ - "metrics": ["v3.avg_order_value"], - "dimensions": ["v3.product.category"], - "table": table_spec, - "measure_columns": { + "metrics": { "v3.total_revenue": "revenue_total", "v3.order_count": "order_cnt", }, + "dimensions": {"v3.product.category": "category"}, + "table": table_spec, }, ) assert response.status_code == 201, response.text @@ -3079,14 +3108,13 @@ async def test_inheritance_skips_other_declarations_on_the_node( other_measures = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.max_unit_price"], - "dimensions": ["v3.product.category"], + "metrics": {"v3.max_unit_price": "up_max"}, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "scan_other_measures_agg", }, - "measure_columns": {"v3.max_unit_price": "up_max"}, }, ) assert other_measures.status_code == 201, other_measures.text @@ -3095,14 +3123,13 @@ async def test_inheritance_skips_other_declarations_on_the_node( other_grain = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.total_revenue"], - "dimensions": ["v3.order_details.status"], + "metrics": {"v3.total_revenue": "revenue_total"}, + "dimensions": {"v3.order_details.status": "status"}, "table": { "catalog": "default", "schema": "analytics", "table": "scan_other_grain_agg", }, - "measure_columns": {"v3.total_revenue": "revenue_total"}, }, ) assert other_grain.status_code == 201, other_grain.text @@ -3213,18 +3240,17 @@ async def test_one_availability_state_per_registration( response = await client_with_build_v3.post( "/preaggs/register", json={ - "metrics": ["v3.total_revenue", "v3.page_view_count"], - "dimensions": ["v3.product.category"], + "metrics": { + "v3.total_revenue": "revenue_total", + "v3.page_view_count": "view_cnt", + }, + "dimensions": {"v3.product.category": "category"}, "table": { "catalog": "default", "schema": "analytics", "table": "cross_fact_agg", "valid_through_ts": 1700000000, }, - "measure_columns": { - "v3.total_revenue": "revenue_total", - "v3.page_view_count": "view_cnt", - }, }, ) assert response.status_code == 201, response.text diff --git a/datajunction-server/tests/construction/build_v3/preagg_substitution_test.py b/datajunction-server/tests/construction/build_v3/preagg_substitution_test.py index e2ec14858..681be295f 100644 --- a/datajunction-server/tests/construction/build_v3/preagg_substitution_test.py +++ b/datajunction-server/tests/construction/build_v3/preagg_substitution_test.py @@ -47,22 +47,32 @@ async def _register_external_preagg( metrics, dimensions, table_ref, - measure_columns, table_columns, - dimension_columns=None, expected_status=201, ): """ Register an externally-built pre-aggregation via /preaggs/register with a mocked query service that reports ``table_columns`` for the external table. - Returns the response (asserts ``expected_status``). + + ``metrics`` and ``dimensions`` are the request's maps from reference to the + physical column of the external table that holds it. Returns the response + (asserts ``expected_status``). + + ``table_columns`` is either a name -> type map, or a bare list of names whose + types are inferred from what each column backs: `double` for a measure and + `string` for a dimension, which is what most of these fixtures hold. Every + binding is type-checked at registration, so a test whose dimension is numeric + (a date or an id) passes the map form and spells the type out. """ - items = ( - table_columns.items() - if isinstance(table_columns, dict) - else [(name, "double") for name in table_columns] - ) + if isinstance(table_columns, dict): + items = list(table_columns.items()) + else: + bound_to_a_dimension = set(dimensions.values()) + items = [ + (name, "string" if name in bound_to_a_dimension else "double") + for name in table_columns + ] async def _fake_columns(*args, **kwargs): return [SimpleNamespace(name=name, type=type_str) for name, type_str in items] @@ -75,10 +85,7 @@ async def _fake_columns(*args, **kwargs): "metrics": metrics, "dimensions": dimensions, "table": table_ref, - "measure_columns": measure_columns, } - if dimension_columns is not None: - payload["dimension_columns"] = dimension_columns response = await client.post("/preaggs/register", json=payload) assert response.status_code == expected_status, response.text return response @@ -94,15 +101,14 @@ async def test_external_preagg_used_at_exact_grain(self, client_with_build_v3): """An exact-grain query reads the external table's physical source column.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns=["status", "revenue_sum"], ) # Measures SQL: reads the external table, applying SUM over the @@ -154,15 +160,14 @@ async def test_external_preagg_rolls_up_additive(self, client_with_build_v3): """An additive measure rolls up from an external pre-agg at a coarser grain.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns=["status", "revenue_sum"], ) # Query at a coarser grain (no dimensions) -> roll up the additive sum @@ -204,15 +209,14 @@ async def test_external_non_additive_not_rolled_up(self, client_with_build_v3): grain; the query falls back to raw sources.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.order_count"], - dimensions=["v3.order_details.status"], + metrics={"v3.order_count": "order_cnt"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "orders_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.order_count": "order_cnt"}, table_columns=["status", "order_cnt"], ) # Coarser grain than the pre-agg -> a distinct count cannot be summed. @@ -243,18 +247,17 @@ async def test_external_preagg_serves_unregistered_derived_metric( """ response = await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_quantity", "v3.order_count"], - dimensions=["v3.order_details.status"], + metrics={ + "v3.total_quantity": "qty_sum", + "v3.order_count": "order_id_col", + }, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "qty_orders_by_status", "valid_through_ts": 20250101, }, - measure_columns={ - "v3.total_quantity": "qty_sum", - "v3.order_count": "order_id_col", - }, table_columns={ "status": "string", "qty_sum": "double", @@ -303,15 +306,14 @@ async def test_external_preagg_not_used_for_incompatible_aggregation( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_unit_price"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_unit_price": "unit_price_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "unit_price_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_unit_price": "unit_price_sum"}, table_columns=["status", "unit_price_sum"], ) # SUM metric routes to the agg (control). @@ -369,18 +371,17 @@ async def test_external_preagg_same_expression_two_aggregations( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_unit_price", "v3.max_unit_price"], - dimensions=["v3.order_details.status"], + metrics={ + "v3.total_unit_price": "unit_price_sum", + "v3.max_unit_price": "unit_price_max", + }, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "unit_price_both", "valid_through_ts": 20250101, }, - measure_columns={ - "v3.total_unit_price": "unit_price_sum", - "v3.max_unit_price": "unit_price_max", - }, table_columns={ "status": "string", "unit_price_sum": "double", @@ -438,15 +439,14 @@ async def test_separate_preaggs_same_grain_different_aggregation_coexist( ): await _register_external_preagg( client_with_build_v3, - metrics=[metric], - dimensions=["v3.order_details.status"], + metrics={metric: column}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": table, "valid_through_ts": 20250101, }, - measure_columns={metric: column}, table_columns={"status": "string", column: "double"}, ) @@ -500,15 +500,14 @@ async def test_external_preagg_stranded_by_node_revision( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"status": "string", "rev_sum": "double"}, ) params = { @@ -563,31 +562,28 @@ async def test_external_preagg_same_grain_tie_is_first_registered( displace an earlier candidate. Pinning this makes the tie-break explicit rather than incidental. """ - for table, metrics, measure_columns, table_columns in ( + for table, metric_bindings, table_columns in ( ( "revenue_only", - ["v3.total_revenue"], {"v3.total_revenue": "rev_sum"}, {"status": "string", "rev_sum": "double"}, ), ( "revenue_and_quantity", - ["v3.total_revenue", "v3.total_quantity"], {"v3.total_revenue": "rev_sum", "v3.total_quantity": "qty_sum"}, {"status": "string", "rev_sum": "double", "qty_sum": "double"}, ), ): await _register_external_preagg( client_with_build_v3, - metrics=metrics, - dimensions=["v3.order_details.status"], + metrics=metric_bindings, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": table, "valid_through_ts": 20250101, }, - measure_columns=measure_columns, table_columns=table_columns, ) @@ -630,16 +626,14 @@ async def test_external_preagg_multi_hop_role_dimension( dimension = "v3.location.country[customer->home]" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[dimension], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={dimension: "home_country"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_home_country", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={dimension: "home_country"}, table_columns={"home_country": "string", "rev_sum": "double"}, ) response = await client_with_build_v3.get( @@ -662,14 +656,13 @@ async def test_external_preagg_pending_not_used(self, client_with_build_v3): used to answer queries.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_pending", }, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns=["status", "revenue_sum"], ) response = await client_with_build_v3.get( @@ -692,18 +685,17 @@ async def test_external_preagg_multiple_measures_covered( from it at the exact grain.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue", "v3.total_quantity"], - dimensions=["v3.order_details.status"], + metrics={ + "v3.total_revenue": "revenue_sum", + "v3.total_quantity": "qty_sum", + }, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_qty_by_status", "valid_through_ts": 20250101, }, - measure_columns={ - "v3.total_revenue": "revenue_sum", - "v3.total_quantity": "qty_sum", - }, table_columns=["status", "revenue_sum", "qty_sum"], ) measures_response = await client_with_build_v3.get( @@ -763,18 +755,17 @@ async def test_external_preagg_partial_metric_coverage_falls_back( even the covered total_revenue/total_quantity come from source).""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue", "v3.total_quantity"], - dimensions=["v3.order_details.status"], + metrics={ + "v3.total_revenue": "revenue_sum", + "v3.total_quantity": "qty_sum", + }, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_qty_by_status", "valid_through_ts": 20250101, }, - measure_columns={ - "v3.total_revenue": "revenue_sum", - "v3.total_quantity": "qty_sum", - }, table_columns=["status", "revenue_sum", "qty_sum"], ) # total_unit_price is NOT covered by the pre-agg. @@ -825,19 +816,22 @@ async def test_external_preagg_rollup_over_covered_dimension( rolling the additive measure up over the dropped dimension.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.order_details.status", - "v3.order_details.product_id", - ], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={ + "v3.order_details.status": "status", + "v3.order_details.product_id": "product_id", + }, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status_product", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, - table_columns=["status", "product_id", "revenue_sum"], + table_columns={ + "status": "string", + "product_id": "int", + "revenue_sum": "double", + }, ) measures_response = await client_with_build_v3.get( "/sql/measures/v3/", @@ -866,15 +860,14 @@ async def test_external_preagg_extra_dimension_falls_back( cannot add a grain), so it computes from source.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns=["status", "revenue_sum"], ) # product_id is not in the pre-agg grain -> cannot be served by it. @@ -918,15 +911,14 @@ async def test_external_preagg_filter_on_covered_dimension( pre-agg read.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns=["status", "revenue_sum"], ) measures_response = await client_with_build_v3.get( @@ -985,15 +977,14 @@ async def test_external_preagg_filter_on_uncovered_column( to source (the pre-agg has already aggregated that column away).""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns=["status", "revenue_sum"], ) measures_response = await client_with_build_v3.get( @@ -1036,16 +1027,17 @@ async def test_external_preagg_filter_on_rolled_up_dimension( Previously the predicate was dropped and the result over-counted.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status", "v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.order_details.status": "status", + "v3.product.category": "cat", + }, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_status_cat", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.product.category": "cat"}, table_columns={"status": "string", "cat": "string", "rev_sum": "double"}, ) params = { @@ -1104,16 +1096,17 @@ async def test_external_preagg_projected_and_rolled_up_filters( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status", "v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.order_details.status": "status", + "v3.product.category": "cat", + }, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_status_cat_both", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.product.category": "cat"}, table_columns={"status": "string", "cat": "string", "rev_sum": "double"}, ) params = { @@ -1172,21 +1165,19 @@ async def test_external_preagg_coarsest_covering_agg_wins( only the finer one carries reads the finer one -- each through its own column mapping. """ - for table, dimensions, dimension_columns, table_columns in ( + for table, dimensions, table_columns in ( ( "rev_by_status_cat", - ["v3.order_details.status", "v3.product.category"], - {"v3.product.category": "cat"}, + { + "v3.order_details.status": "status", + "v3.product.category": "cat", + }, {"status": "string", "cat": "string", "rev_sum": "double"}, ), ( "rev_by_status_cat_day", - [ - "v3.order_details.status", - "v3.product.category", - "v3.date.date_id[order]", - ], { + "v3.order_details.status": "status", "v3.product.category": "cat", "v3.date.date_id[order]": "day_key", }, @@ -1200,7 +1191,7 @@ async def test_external_preagg_coarsest_covering_agg_wins( ): await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], + metrics={"v3.total_revenue": "rev_sum"}, dimensions=dimensions, table_ref={ "catalog": "default", @@ -1208,8 +1199,6 @@ async def test_external_preagg_coarsest_covering_agg_wins( "table": table, "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns=dimension_columns, table_columns=table_columns, ) @@ -1281,13 +1270,16 @@ async def test_external_preagg_two_roles_of_one_dimension( ``country[from]`` and ``country[to]`` both resolve to the column ``country``, so they must be disambiguated to ``country_from`` / - ``country_to``; mapping only the ``to`` role checks the two are bound - independently rather than sharing one column. + ``country_to``; storing the ``to`` role under a different physical name + checks the two are bound independently rather than sharing one column. """ - dimensions = ["v3.location.country[from]", "v3.location.country[to]"] + dimensions = { + "v3.location.country[from]": "country_from", + "v3.location.country[to]": "dest_country", + } await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], + metrics={"v3.total_revenue": "rev_sum"}, dimensions=dimensions, table_ref={ "catalog": "default", @@ -1295,15 +1287,13 @@ async def test_external_preagg_two_roles_of_one_dimension( "table": "rev_by_route", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.location.country[to]": "dest_country"}, table_columns={ "country_from": "string", "dest_country": "string", "rev_sum": "double", }, ) - params = {"metrics": ["v3.total_revenue"], "dimensions": dimensions} + params = {"metrics": ["v3.total_revenue"], "dimensions": list(dimensions)} measures_response = await client_with_build_v3.get( "/sql/measures/v3/", params=params, @@ -1349,16 +1339,14 @@ async def test_external_preagg_order_by_mapped_dimension( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.date.date_id[order]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.date.date_id[order]": "day_key"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_day", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.date.date_id[order]": "day_key"}, table_columns={"day_key": "int", "rev_sum": "double"}, ) metrics_response = await client_with_build_v3.get( @@ -1400,16 +1388,15 @@ async def test_external_preagg_cross_fact_partial_substitution( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.customer.customer_id[customer]"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.customer.customer_id[customer]": "customer_id_customer"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_customer", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, - table_columns=["customer_id", "revenue_sum"], + table_columns={"customer_id_customer": "int", "revenue_sum": "double"}, ) metrics_response = await client_with_build_v3.get( "/sql/metrics/v3/", @@ -1463,19 +1450,17 @@ async def test_external_preagg_renamed_dimension_column( client_with_build_v3, ): """A grain dimension stored under a different physical column name is read - via dimension_columns and aliased back to the DJ name.""" + under `dimensions` and aliased back to the DJ name.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "order_status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "revenue_sum"}, - dimension_columns={"v3.order_details.status": "order_status"}, table_columns={"order_status": "string", "revenue_sum": "double"}, ) measures_response = await client_with_build_v3.get( @@ -1510,16 +1495,14 @@ async def test_external_preagg_joined_key_via_foreign_key_column( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.date.date_id[order]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.date.date_id[order]": "day_key"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_day", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.date.date_id[order]": "day_key"}, table_columns={"day_key": "int", "rev_sum": "double"}, ) # Measures SQL: reads the mapped physical column (day_key) and aliases it @@ -1577,14 +1560,13 @@ async def test_external_preagg_rejects_unqualified_role_dimension( """ rejected = await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.location.country"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.location.country": "country"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_country", }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"country": "string", "rev_sum": "double"}, expected_status=422, ) @@ -1596,44 +1578,48 @@ async def test_external_preagg_rejects_unqualified_role_dimension( # References that are legal stay legal: a role-qualified dimension, a # locally-owned column, the FK column behind a join, and a role-free link. - for dimension, table, columns in ( + for dimension, column, table, columns in ( ( "v3.location.country[from]", + "country_from", "revenue_by_from_country", {"country_from": "string", "rev_sum": "double"}, ), ( "v3.order_details.status", + "status", "revenue_by_status_local", {"status": "string", "rev_sum": "double"}, ), ( "v3.order_details.order_date", + "order_date", "revenue_by_order_date", {"order_date": "int", "rev_sum": "double"}, ), ( "v3.product.category", + "category", "revenue_by_category_plain", {"category": "string", "rev_sum": "double"}, ), ( "v3.customer.customer_id", + "customer_id", "revenue_by_customer_bare", {"customer_id": "int", "rev_sum": "double"}, ), ): await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[dimension], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={dimension: column}, table_ref={ "catalog": "default", "schema": "analytics", "table": table, "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns=columns, ) @@ -1648,16 +1634,14 @@ async def test_external_preagg_bare_and_role_qualified_refs_interoperate( """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.customer.customer_id"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.customer.customer_id": "cust_key"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_customer_key", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.customer.customer_id": "cust_key"}, table_columns={"cust_key": "int", "rev_sum": "double"}, ) # The output alias follows the spelling that was requested; only matching @@ -1685,46 +1669,34 @@ async def test_external_preagg_dimension_column_must_exist( self, client_with_build_v3, ): - """A dimension_columns mapping to a column absent from the table is - rejected; an unknown dimension key is rejected too.""" + """A dimension bound to a column the table doesn't have is rejected, and + so is one bound to a column of an incompatible type.""" missing = await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "nope"}, table_ref={"catalog": "default", "schema": "analytics", "table": "t"}, - measure_columns={"v3.total_revenue": "revenue_sum"}, - dimension_columns={"v3.order_details.status": "nope"}, table_columns=["order_status", "revenue_sum"], expected_status=422, ) - assert "not found in table" in missing.json()["message"] - - unknown = await _register_external_preagg( - client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], - table_ref={"catalog": "default", "schema": "analytics", "table": "t"}, - measure_columns={"v3.total_revenue": "revenue_sum"}, - dimension_columns={"v3.order_details.not_a_dim": "x"}, - table_columns=["order_status", "revenue_sum", "x"], - expected_status=422, + assert missing.json()["message"] == ( + "Columns ['nope'] declared under `metrics`/`dimensions` were not " + "found in table default.analytics.t." ) - assert "not in the pre-aggregation's dimensions" in unknown.json()["message"] # A string dimension bound to a numeric column is type-incompatible. bad_type = await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.order_details.status": "status_num"}, table_ref={"catalog": "default", "schema": "analytics", "table": "t"}, - measure_columns={"v3.total_revenue": "revenue_sum"}, - dimension_columns={"v3.order_details.status": "status_num"}, table_columns={"status_num": "bigint", "revenue_sum": "double"}, expected_status=422, ) - assert ( - "not type-compatible with dimension 'v3.order_details.status'" - in bad_type.json()["message"] + assert bad_type.json()["message"] == ( + "Column 'status_num' (type bigint) in table default.analytics.t is " + "not type-compatible with dimension 'v3.order_details.status' " + "(expected string)." ) @pytest.mark.asyncio @@ -1733,20 +1705,18 @@ async def test_external_preagg_joined_attribute_read_directly( client_with_build_v3, ): """A joined dimension attribute stored (denormalized) in the external - table is read directly from it via dimension_columns -- no join back to - the dimension node.""" + table is read directly from its declared column -- no join back to the + dimension node.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.product.category": "cat"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_category", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.product.category": "cat"}, table_columns={"cat": "string", "rev_sum": "double"}, ) measures_response = await client_with_build_v3.get( @@ -1778,19 +1748,14 @@ async def test_external_preagg_renamed_measure_and_dimensions( joined attribute) coexist on one external table.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status", "v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "st", "v3.product.category": "cat"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "revenue_by_status_category", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={ - "v3.order_details.status": "st", - "v3.product.category": "cat", - }, table_columns={"st": "string", "cat": "string", "rev_sum": "double"}, ) measures_response = await client_with_build_v3.get( @@ -1822,28 +1787,29 @@ async def test_external_preagg_filter_only_dimension_selects_finer_agg( ranking candidates picked the coarse agg and dropped the predicate.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"status": "string", "rev_sum": "double"}, ) await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status", "v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.order_details.status": "status", + "v3.product.category": "category", + }, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_status_cat", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={ "status": "string", "category": "string", @@ -1900,16 +1866,17 @@ async def test_external_preagg_filter_pushdown_uses_mapped_column( would reference a column the external table does not have.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status", "v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.order_details.status": "st", + "v3.product.category": "category", + }, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_st_cat", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.order_details.status": "st"}, table_columns={"st": "string", "category": "string", "rev_sum": "double"}, ) params = { @@ -1962,29 +1929,26 @@ async def test_external_preagg_dimension_mapping_is_per_agg( another, depending on which agg the metric routes to.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "order_status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_st", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.order_details.status": "order_status"}, table_columns={"order_status": "string", "rev_sum": "double"}, ) await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_quantity"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_quantity": "qty_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "qty_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_quantity": "qty_sum"}, table_columns={"status": "string", "qty_sum": "double"}, ) revenue_response = await client_with_build_v3.get( @@ -2030,15 +1994,14 @@ async def test_external_preagg_metric_filter_becomes_having( exposes).""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"status": "string", "rev_sum": "double"}, ) params = { @@ -2090,15 +2053,14 @@ async def test_external_preagg_order_by_metric_with_limit( leaks into the pre-agg CTE, which would truncate rows before the roll-up.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"status": "string", "rev_sum": "double"}, ) params = { @@ -2149,29 +2111,29 @@ async def test_external_preagg_degenerate_grains(self, client_with_build_v3): key still rolls up (on the mapped column).""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_total", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"rev_sum": "double"}, ) await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status", "v3.product.category"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.order_details.status": "status", + "v3.product.category": "cat", + }, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_st_cat", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.product.category": "cat"}, table_columns={"status": "string", "cat": "string", "rev_sum": "double"}, ) # (a) No dimensions at all -> the empty-grain agg, scanned without GROUP BY. @@ -2252,15 +2214,14 @@ async def test_external_preagg_uncovered_distinct_measure_falls_back( carries no column for the distinct component.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": "rev_by_status", "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"status": "string", "rev_sum": "double"}, ) params = { @@ -3267,10 +3228,9 @@ async def _register(self, client) -> int: assert partition_response.status_code == 201 response = await _register_external_preagg( client, - metrics=["v3.total_revenue"], - dimensions=["v3.date.date_id[order]"], + metrics={"v3.total_revenue": "revenue_sum"}, + dimensions={"v3.date.date_id[order]": "date_id_order"}, table_ref={**self.TABLE_REF, "valid_through_ts": 20250101}, - measure_columns={"v3.total_revenue": "revenue_sum"}, table_columns={"date_id_order": "int", "revenue_sum": "double"}, ) return response.json()["preaggs"][0]["id"] @@ -3496,14 +3456,12 @@ async def test_joins_back_to_a_retained_key(self, client_with_build_v3): """ await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.customer_id[customer]", - "v3.order_details.status", - ], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.customer.customer_id[customer]": "cust_key", + "v3.order_details.status": "status", + }, table_ref=self.TABLE, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.customer.customer_id[customer]": "cust_key"}, table_columns={"cust_key": "int", "status": "string", "rev_sum": "double"}, ) params = { @@ -3555,23 +3513,20 @@ async def test_direct_match_beats_a_joining_one(self, client_with_build_v3): larger grain -- joins rank ahead of grain size, so no join is emitted.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.customer.customer_id[customer]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.customer.customer_id[customer]": "cust_key"}, table_ref={**self.TABLE, "table": "rev_by_cust_key"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.customer.customer_id[customer]": "cust_key"}, table_columns={"cust_key": "int", "rev_sum": "double"}, ) await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.name[customer]", - "v3.order_details.status", - "v3.product.category", - ], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.customer.name[customer]": "name_customer", + "v3.order_details.status": "status", + "v3.product.category": "category", + }, table_ref={**self.TABLE, "table": "rev_by_name_status_cat"}, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={ "name_customer": "string", "status": "string", @@ -3615,11 +3570,9 @@ async def test_non_additive_measure_builds_from_source( join back is by definition a roll-up -- so the query builds from source.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.order_count"], - dimensions=["v3.customer.customer_id[customer]"], + metrics={"v3.order_count": "order_cnt"}, + dimensions={"v3.customer.customer_id[customer]": "cust_key"}, table_ref={**self.TABLE, "table": "orders_by_cust"}, - measure_columns={"v3.order_count": "order_cnt"}, - dimension_columns={"v3.customer.customer_id[customer]": "cust_key"}, table_columns={"cust_key": "int", "order_cnt": "bigint"}, ) params = { @@ -3674,14 +3627,12 @@ async def test_filter_on_a_joined_back_attribute(self, client_with_build_v3): that doesn't exist.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.customer_id[customer]", - "v3.order_details.status", - ], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.customer.customer_id[customer]": "cust_key", + "v3.order_details.status": "status", + }, table_ref=self.TABLE, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.customer.customer_id[customer]": "cust_key"}, table_columns={"cust_key": "int", "status": "string", "rev_sum": "double"}, ) params = { @@ -3739,18 +3690,13 @@ async def test_two_joined_back_attributes_are_both_filtered( into the same WHERE against their own dimension aliases.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.customer_id[customer]", - "v3.product.product_id", - "v3.order_details.status", - ], - table_ref={**self.TABLE, "table": "rev_cust_prod_status"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={ + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ "v3.customer.customer_id[customer]": "cust_key", "v3.product.product_id": "prod_key", + "v3.order_details.status": "status", }, + table_ref={**self.TABLE, "table": "rev_cust_prod_status"}, table_columns={ "cust_key": "int", "prod_key": "int", @@ -3828,10 +3774,9 @@ async def test_retained_column_that_is_not_the_key(self, client_with_build_v3): and multiply the revenue, so the query builds from source instead.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.customer.location_id[customer]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.customer.location_id[customer]": "location_id_customer"}, table_ref={**self.TABLE, "table": "rev_by_cust_loc"}, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"location_id_customer": "int", "rev_sum": "double"}, ) params = { @@ -3887,10 +3832,9 @@ async def test_key_retained_at_another_role(self, client_with_build_v3): asks for the ``to`` role: a different join path, so no match.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.location.location_id[from]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.location.location_id[from]": "location_id_from"}, table_ref={**self.TABLE, "table": "rev_by_from_loc"}, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"location_id_from": "int", "rev_sum": "double"}, ) params = { @@ -3974,17 +3918,12 @@ async def test_composite_key_fully_retained(self, client_with_build_v3): await self._add_composite_key_dimension(client_with_build_v3) await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.order_line.order_id[line]", - "v3.order_line.line_number[line]", - ], - table_ref={**self.TABLE, "table": "rev_by_line"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={ + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ "v3.order_line.order_id[line]": "ord", "v3.order_line.line_number[line]": "ln", }, + table_ref={**self.TABLE, "table": "rev_by_line"}, table_columns={"ord": "int", "ln": "int", "rev_sum": "double"}, ) params = { @@ -4042,10 +3981,9 @@ async def test_composite_key_only_partly_retained(self, client_with_build_v3): await self._add_composite_key_dimension(client_with_build_v3) await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.order_line.order_id[line]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_line.order_id[line]": "order_id_line"}, table_ref={**self.TABLE, "table": "rev_by_order"}, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"order_id_line": "int", "rev_sum": "double"}, ) params = { @@ -4109,17 +4047,12 @@ async def test_two_dimensions_joined_back(self, client_with_build_v3): """Two retained keys, two dimensions joined onto the same scan.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.customer_id[customer]", - "v3.product.product_id", - ], - table_ref={**self.TABLE, "table": "rev_by_cust_prod"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={ + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ "v3.customer.customer_id[customer]": "cust_key", "v3.product.product_id": "prod_key", }, + table_ref={**self.TABLE, "table": "rev_by_cust_prod"}, table_columns={"cust_key": "int", "prod_key": "int", "rev_sum": "double"}, ) params = { @@ -4180,17 +4113,12 @@ async def test_one_dimension_joined_back_at_two_roles(self, client_with_build_v3 under distinct aliases -- the two ``city`` columns don't collide.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.location.location_id[from]", - "v3.location.location_id[to]", - ], - table_ref={**self.TABLE, "table": "rev_by_lanes"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={ + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ "v3.location.location_id[from]": "orig", "v3.location.location_id[to]": "dest", }, + table_ref={**self.TABLE, "table": "rev_by_lanes"}, table_columns={"orig": "int", "dest": "int", "rev_sum": "double"}, ) params = { @@ -4250,17 +4178,12 @@ async def test_joined_attribute_shares_a_name_with_a_preagg_column( than the pre-agg's.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.customer_id[customer]", - "v3.order_details.status", - ], - table_ref={**self.TABLE, "table": "rev_shadowed"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={ + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ "v3.customer.customer_id[customer]": "cust_key", "v3.order_details.status": "name", }, + table_ref={**self.TABLE, "table": "rev_shadowed"}, table_columns={"cust_key": "int", "name": "string", "rev_sum": "double"}, ) params = { @@ -4324,11 +4247,9 @@ async def test_joins_back_to_a_materialized_dimension(self, client_with_build_v3 assert response.status_code == 200, response.text await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=["v3.customer.customer_id[customer]"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.customer.customer_id[customer]": "cust_key"}, table_ref={**self.TABLE, "table": "rev_by_cust_key"}, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.customer.customer_id[customer]": "cust_key"}, table_columns={"cust_key": "int", "rev_sum": "double"}, ) params = { @@ -4378,14 +4299,12 @@ async def test_retained_key_is_also_requested(self, client_with_build_v3): collide.""" await _register_external_preagg( client_with_build_v3, - metrics=["v3.total_revenue"], - dimensions=[ - "v3.customer.customer_id[customer]", - "v3.order_details.status", - ], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={ + "v3.customer.customer_id[customer]": "cust_key", + "v3.order_details.status": "status", + }, table_ref=self.TABLE, - measure_columns={"v3.total_revenue": "rev_sum"}, - dimension_columns={"v3.customer.customer_id[customer]": "cust_key"}, table_columns={"cust_key": "int", "status": "string", "rev_sum": "double"}, ) params = { @@ -4462,18 +4381,17 @@ async def _register_wide(client, table="wide_by_status"): """A pre-agg covering both revenue and quantity, at status grain.""" return await _register_external_preagg( client, - metrics=["v3.total_revenue", "v3.total_quantity"], - dimensions=["v3.order_details.status"], + metrics={ + "v3.total_revenue": "rev_sum", + "v3.total_quantity": "qty_sum", + }, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": table, "valid_through_ts": 20250101, }, - measure_columns={ - "v3.total_revenue": "rev_sum", - "v3.total_quantity": "qty_sum", - }, table_columns={ "status": "string", "rev_sum": "double", @@ -4486,15 +4404,14 @@ async def _register_narrow(client, table="narrow_by_status"): """A pre-agg covering only revenue, at the same grain.""" return await _register_external_preagg( client, - metrics=["v3.total_revenue"], - dimensions=["v3.order_details.status"], + metrics={"v3.total_revenue": "rev_sum"}, + dimensions={"v3.order_details.status": "status"}, table_ref={ "catalog": "default", "schema": "analytics", "table": table, "valid_through_ts": 20250101, }, - measure_columns={"v3.total_revenue": "rev_sum"}, table_columns={"status": "string", "rev_sum": "double"}, ) diff --git a/datajunction-server/tests/models/deployment_test.py b/datajunction-server/tests/models/deployment_test.py index 3dc7068f8..3bf5cb152 100644 --- a/datajunction-server/tests/models/deployment_test.py +++ b/datajunction-server/tests/models/deployment_test.py @@ -1,5 +1,4 @@ import pytest -from pydantic import ValidationError from datajunction_server.errors import DJInvalidDeploymentConfig from datajunction_server.models.deployment import ( @@ -657,15 +656,12 @@ def _preagg_spec_dict() -> dict: def test_preagg_spec_renders_column_bindings(): """ - The maps are split back into the references-plus-bindings shape the - registration internals take, with every reference prefix-rendered against - the namespace and the physical columns left as-is. + Every reference is prefix-rendered against the namespace and stays bound to + the physical column that holds it, which is left as-is. """ spec = PreAggSpec.model_validate(_preagg_spec_dict()) - assert spec.rendered_metrics == ["ns.count"] - assert spec.rendered_dimensions == ["ns.d.attr", "ns.d.same"] - assert spec.rendered_measure_columns == {"ns.count": "cnt"} - assert spec.rendered_dimension_columns == { + assert spec.rendered_metrics == {"ns.count": "cnt"} + assert spec.rendered_dimensions == { "ns.d.attr": "phys_attr", "ns.d.same": "same", } @@ -693,94 +689,9 @@ def test_preagg_spec_round_trips_through_model_dump(): } -def test_preagg_spec_rejects_measure_columns_block(): - """ - The four-field form is gone: `measure_columns` would otherwise be ignored as - an unknown key, silently dropping every binding it declared. - """ - spec = _preagg_spec_dict() - spec["metrics"] = ["${prefix}count"] - spec["measure_columns"] = {"${prefix}count": "cnt"} - with pytest.raises(DJInvalidDeploymentConfig) as exc_info: - PreAggSpec.model_validate(spec) - assert exc_info.value.message == ( - "Pre-aggregation 'p' declares `measure_columns`, which is no longer a " - "pre-aggregation field. Declare the physical column alongside what it " - "holds instead, as `metrics: {: }`, and drop the " - "`measure_columns` block." - ) - - -def test_preagg_spec_rejects_dimension_columns_block(): - """Same for `dimension_columns`, the other half of the retired form.""" - spec = _preagg_spec_dict() - spec["dimensions"] = ["${prefix}d.attr"] - spec["dimension_columns"] = {"${prefix}d.attr": "phys_attr"} - with pytest.raises(DJInvalidDeploymentConfig) as exc_info: - PreAggSpec.model_validate(spec) - assert exc_info.value.message == ( - "Pre-aggregation 'p' declares `dimension_columns`, which is no longer a " - "pre-aggregation field. Declare the physical column alongside what it " - "holds instead, as `dimensions: {: }`, and drop the " - "`dimension_columns` block." - ) - - -def test_preagg_spec_rejects_list_of_references(): - """ - A bare list of references carries no bindings at all, so it is refused with - the shape to write instead rather than a pydantic type error. - """ - spec = _preagg_spec_dict() - spec["metrics"] = ["${prefix}count"] - with pytest.raises(DJInvalidDeploymentConfig) as exc_info: - PreAggSpec.model_validate(spec) - assert exc_info.value.message == ( - "Pre-aggregation 'p' declares `metrics` as a list. `metrics` is a map " - "from each reference to the physical column of the external table that " - "holds it, e.g. `metrics: {: }`." - ) - - -def test_preagg_spec_rejects_metric_without_column(): - """A measure's DJ-side name is hashed, so there is nothing to default to.""" - spec = _preagg_spec_dict() - spec["metrics"] = {"${prefix}count": None, "${prefix}total": "total_sum"} - with pytest.raises(DJInvalidDeploymentConfig) as exc_info: - PreAggSpec.model_validate(spec) - assert exc_info.value.message == ( - "Pre-aggregation 'p' leaves the physical column empty under `metrics` " - "for ['${prefix}count']. A measure's DJ-side name is auto-generated " - "with an expression-hash suffix, so there is no name to fall back on." - ) - - -def test_preagg_spec_rejects_dimension_without_column(): - """ - A dimension is held to the same rule, even though its DJ column name would - be a plausible default: a trailing colon is too easy to write by accident, - and the written-out name documents the table. - """ - spec = _preagg_spec_dict() - spec["dimensions"] = {"${prefix}d.attr": "phys_attr", "${prefix}d.same": None} - with pytest.raises(DJInvalidDeploymentConfig) as exc_info: - PreAggSpec.model_validate(spec) - assert exc_info.value.message == ( - "Pre-aggregation 'p' leaves the physical column empty under " - "`dimensions` for ['${prefix}d.same']. Write the column out even when " - "it matches the DJ column name, so the file says what the table " - "actually holds." - ) - - -def test_preagg_spec_rejects_non_mapping_input(): - """ - Input that isn't a mapping at all falls through the binding checks and gets - pydantic's own error, rather than blowing up inside the validator. - """ - with pytest.raises(ValidationError) as exc_info: - PreAggSpec.model_validate(["not", "a", "spec"]) - assert exc_info.value.errors()[0]["type"] == "model_type" +# The three rejections that hold a spec to the map form are shared with +# POST /preaggs/register and asserted for both surfaces in +# tests/models/preagg_binding_test.py. def test_tag_spec_metadata_aliases(): diff --git a/datajunction-server/tests/models/preagg_binding_test.py b/datajunction-server/tests/models/preagg_binding_test.py new file mode 100644 index 000000000..fb621bd4e --- /dev/null +++ b/datajunction-server/tests/models/preagg_binding_test.py @@ -0,0 +1,242 @@ +""" +Tests for the column bindings a pre-aggregation declares inline. + +Two surfaces declare an externally-built pre-aggregation table: the ``kind: preagg`` +YAML spec (``PreAggSpec``) and the body of ``POST /preaggs/register`` +(``RegisterPreAggregationsRequest``). Both spell it as maps from reference to +physical column, and both enforce that through +``datajunction_server.models.preagg_binding``, so every rule below is asserted +against both models rather than once against one of them. + +The messages are pinned in full. They are the migration path for anyone holding a +file or an API client written against the retired four-field form -- someone should +be able to fix their payload from the message alone -- so a reworded message is a +change worth noticing. +""" + +import pytest +from pydantic import ValidationError + +from datajunction_server.errors import ( + DJInvalidDeploymentConfig, + DJInvalidInputException, +) +from datajunction_server.models.deployment import PreAggSpec +from datajunction_server.models.preaggregation import RegisterPreAggregationsRequest + +# The shape each surface shows, in the syntax that surface is authored in. Spelled +# out here rather than imported, so a change to the example is a change to a test. +YAML_SHAPE = """metrics: + default.view_secs: view_secs_sum + default.session_count: session_cnt +dimensions: + default.page_d.page_id: page + default.geo_country_d.country_iso_code: country_iso_code""" + +JSON_SHAPE = """{ + "metrics": { + "default.view_secs": "view_secs_sum", + "default.session_count": "session_cnt" + }, + "dimensions": { + "default.page_d.page_id": "page", + "default.geo_country_d.country_iso_code": "country_iso_code" + } +}""" + + +def _spec_payload(**overrides) -> dict: + """A valid `kind: preagg` spec, before any override under test.""" + payload = { + "name": "p", + "namespace": "ns", + "metrics": {"${prefix}count": "cnt"}, + "dimensions": {"${prefix}d.attr": "phys_attr"}, + "catalog": "c", + "schema": "s", + "table": "t", + } + payload.update(overrides) + return payload + + +def _request_payload(**overrides) -> dict: + """A valid POST /preaggs/register body, before any override under test.""" + payload = { + "metrics": {"ns.count": "cnt"}, + "dimensions": {"ns.d.attr": "phys_attr"}, + "table": {"catalog": "c", "schema": "s", "table": "t"}, + } + payload.update(overrides) + return payload + + +# Each surface: how to build a payload, how to validate it, the exception it +# raises, how its messages name the declaration, and the shape they show. +SURFACES = [ + pytest.param( + _spec_payload, + PreAggSpec.model_validate, + DJInvalidDeploymentConfig, + "Pre-aggregation 'p'", + YAML_SHAPE, + id="yaml-spec", + ), + pytest.param( + _request_payload, + RegisterPreAggregationsRequest.model_validate, + DJInvalidInputException, + "The POST /preaggs/register request body", + JSON_SHAPE, + id="register-request", + ), +] + +SURFACE_ARGS = "payload_for,validate,exception,subject,shape" + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +def test_map_form_is_accepted(payload_for, validate, exception, subject, shape): + """The map form validates on both surfaces.""" + declaration = validate(payload_for()) + assert declaration.metrics == {list(payload_for()["metrics"])[0]: "cnt"} + assert declaration.dimensions == { + list(payload_for()["dimensions"])[0]: "phys_attr", + } + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +@pytest.mark.parametrize( + "retired,replacement,noun", + [ + ("measure_columns", "metrics", "metric"), + ("dimension_columns", "dimensions", "dimension"), + ], +) +def test_rejects_retired_column_blocks( + retired, + replacement, + noun, + payload_for, + validate, + exception, + subject, + shape, +): + """ + The retired four-field form is refused rather than silently ignored: an + unknown key would drop every binding it declared. + """ + with pytest.raises(exception) as exc_info: + validate(payload_for(**{retired: {"ns.count": "cnt"}})) + assert exc_info.value.message == ( + f"{subject} declares `{retired}`, which is no longer a pre-aggregation " + f"field. Declare each {noun} together with the physical column of the " + f"external table that holds it, under `{replacement}`, and drop the " + f"`{retired}` block. The map form looks like this:\n\n{shape}" + ) + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +@pytest.mark.parametrize( + "field,noun", + [("metrics", "metric"), ("dimensions", "dimension")], +) +def test_rejects_list_of_references( + field, + noun, + payload_for, + validate, + exception, + subject, + shape, +): + """ + A bare list of references carries no bindings at all, so it earns the shape + to write instead rather than a pydantic type error. + """ + with pytest.raises(exception) as exc_info: + validate(payload_for(**{field: ["ns.count"]})) + assert exc_info.value.message == ( + f"{subject} declares `{field}` as a list, not a map. `{field}` maps each " + f"{noun} reference to the physical column of the external table that " + f"holds it. The map form looks like this:\n\n{shape}" + ) + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +def test_rejects_metric_without_column( + payload_for, + validate, + exception, + subject, + shape, +): + """A measure's DJ-side name is hashed, so there is nothing to default to.""" + with pytest.raises(exception) as exc_info: + validate(payload_for(metrics={"ns.count": None, "ns.total": "total_sum"})) + assert exc_info.value.message == ( + f"{subject} leaves the physical column empty under `metrics` for " + f"['ns.count']. A measure's DJ-side name is auto-generated with an " + f"expression-hash suffix, so there is no name to fall back on. The map " + f"form looks like this:\n\n{shape}" + ) + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +def test_rejects_dimension_without_column( + payload_for, + validate, + exception, + subject, + shape, +): + """ + A dimension is held to the same rule even though its DJ column name would be + a plausible default: a trailing colon is too easy to write by accident, and + the written-out name documents the table. + """ + with pytest.raises(exception) as exc_info: + validate(payload_for(dimensions={"ns.d.attr": "phys_attr", "ns.d.same": None})) + assert exc_info.value.message == ( + f"{subject} leaves the physical column empty under `dimensions` for " + f"['ns.d.same']. Write the column out even when it matches the DJ column " + f"name, so the declaration says what the table actually holds. The map " + f"form looks like this:\n\n{shape}" + ) + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +def test_matching_column_still_has_to_be_written_out( + payload_for, + validate, + exception, + subject, + shape, +): + """ + A dimension stored under the name DJ uses is written out like any other. The + binding is kept as declared rather than collapsed to "same as default", so + what the file says is what the table holds. + """ + dimension = list(payload_for()["dimensions"])[0] + attribute = dimension.rsplit(".", 1)[-1] + declaration = validate(payload_for(dimensions={dimension: attribute})) + assert declaration.dimensions == {dimension: attribute} + + +@pytest.mark.parametrize(SURFACE_ARGS, SURFACES) +def test_non_mapping_input_falls_through_to_pydantic( + payload_for, + validate, + exception, + subject, + shape, +): + """ + Input that isn't a mapping at all falls through the binding checks and gets + pydantic's own error, rather than blowing up inside the validator. + """ + with pytest.raises(ValidationError) as exc_info: + validate(["not", "a", "declaration"]) + assert exc_info.value.errors()[0]["type"] == "model_type" diff --git a/docs/content/0.1.0/docs/dj-concepts/aggregate-awareness.md b/docs/content/0.1.0/docs/dj-concepts/aggregate-awareness.md index 42d401370..8eee3d094 100644 --- a/docs/content/0.1.0/docs/dj-concepts/aggregate-awareness.md +++ b/docs/content/0.1.0/docs/dj-concepts/aggregate-awareness.md @@ -248,36 +248,44 @@ gets reconciled on every deploy. #### REST: `POST /preaggs/register` -Send DJ the metrics and dimensions the table covers, where the table lives, and which physical column -backs each measure metric: +Send DJ where the table lives, plus the metrics and dimensions it covers — each one named together with +the physical column that holds it: ```json { - "metrics": ["default.view_secs", "default.session_count", "default.view_rate"], - "dimensions": ["default.page_d.page_id"], + "metrics": { + "default.view_secs": "view_secs_sum", + "default.session_count": "session_cnt" + }, + "dimensions": { + "default.page_d.page_id": "page", + "default.geo_country_d.country_iso_code": "country_iso_code" + }, "table": { "catalog": "warehouse", "schema": "analytics", "table": "views_by_page_daily" - }, - "measure_columns": { - "default.view_secs": "view_secs_sum", - "default.session_count": "session_cnt" } } ``` -Notice that `view_rate`, a ratio metric, is listed under `metrics` but doesn't appear in -`measure_columns` — it doesn't need a column, because it's derived from `view_secs` and -`session_count`, which are the two measures that do have columns. - -On registration DJ decomposes every metric you listed, then validates the binding: it confirms each key -in `measure_columns` really is a measure, checks (via query-service introspection) that every column you -named actually exists in the table **and is type-compatible with the measure it backs** — a numeric -`SUM` can't bind to a string column, for instance (the check is category-level, so `int` vs `bigint` -vs `decimal` are all fine) — and confirms that every measure any of your metrics decomposes into is -covered by some column. If all of that checks out, DJ records the pre-aggregation. If you also pass a -`valid_through_ts`, DJ marks it available immediately so routing can start using it right away. +Both maps require a value for every key, so `page_d.page_id` says it is stored as `page`, and +`geo_country_d.country_iso_code`, stored under the very name DJ uses, says that too rather than leaving +the value off. This is the same declaration the YAML spec below makes, in the same shape — the two paths +differ in where the file lives and when it is reconciled, not in what you write. + +What goes under `metrics` are the *measures* the table physically stores. A ratio metric such as +`view_rate` has no column of its own, so it can't be named here — and doesn't need to be, for the reason +[the YAML section](#yaml-a-kind-preagg-file) spells out: declaring `view_secs` and `session_count` is +what covers everything built on them. + +On registration DJ decomposes the metrics you named, then validates every binding: it confirms each key +under `metrics` really is a measure, and checks (via query-service introspection) that every column you +named actually exists in the table **and is type-compatible with the metric or dimension it backs** — a +numeric `SUM` can't bind to a string column, for instance (the check is category-level, so `int` vs +`bigint` vs `decimal` are all fine). If all of that checks out, DJ records the pre-aggregation. If you +also pass a `valid_through_ts`, DJ marks it available immediately so routing can start using it right +away. #### YAML: a `kind: preagg` file