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