Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions datajunction-clients/python/tests/test_deploy.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,6 @@ 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"
"catalog: default\n"
Expand All @@ -69,7 +68,7 @@ def test_reconstruct_deployment_spec_separates_preaggregations(tmp_path):
assert [node["name"] for node in spec["nodes"]] == ["ns.revenue"]
assert len(spec["preaggregations"]) == 1
preagg = spec["preaggregations"][0]
assert preagg["name"] == "revenue_by_day"
assert preagg["table"] == "revenue_agg"
assert preagg["measure_columns"] == {"ns.revenue": "revenue_sum"}
# The discriminator is stripped before reaching the payload.
assert "kind" not in preagg
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -771,7 +771,6 @@ async def register_preaggregations(
session,
query_service_client,
request_headers,
name=data.name,
metrics=data.metrics,
dimensions=data.dimensions,
table=data.table,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -264,8 +264,12 @@ class PreAggregation(Base):
unique=True,
)

# Optional stable, human-supplied handle for externally-registered pre-aggs.
# Used by YAML deploy reconciliation and availability-by-name callbacks.
# Legacy human-supplied handle for externally-registered pre-aggs. Retired:
# registration no longer accepts a name and nothing reads this column. It
# was never load-bearing -- reconciliation matches on revision + grain +
# measure identity, and the availability callback is addressed by id
# (POST /preaggs/{preagg_id}/availability/). Kept nullable so rows written
# before it was retired still load.
name: Mapped[str | None] = mapped_column(String, nullable=True)

# === Materialization Config ===
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
from datajunction_server.database.namespace import NodeNamespace
from datajunction_server.database.node import MissingParent, NodeRelationship
from datajunction_server.database.partition import Partition
from datajunction_server.database.preaggregation import PreAggregation
from datajunction_server.database.tag import Tag
from datajunction_server.database.user import OAuthProvider, User
from datajunction_server.errors import (
Expand Down Expand Up @@ -89,6 +90,7 @@
MaterializationSpec,
MetricSpec,
NodeSpec,
PreAggSpec,
SourceSpec,
TagSpec,
bump_version,
Expand Down Expand Up @@ -127,6 +129,32 @@
logger = logging.getLogger(__name__)


def preagg_label(table_ref: str, grain_columns: list[str]) -> str:
"""
How a pre-aggregation is named in deployment results: the external table it
adopts, plus the grain it covers.

The table alone would be ambiguous, since one registration can produce
several pre-aggregations on the same table -- one per grain group -- and the
two together are what reconciliation actually matches on. Grain columns are
sorted so the label is stable no matter how the spec ordered them.
"""
return f"{table_ref} [{', '.join(sorted(grain_columns))}]"


def existing_preagg_label(preagg: PreAggregation) -> str:
"""
The same label for a pre-aggregation already in the database. Its table
comes from availability, which a pre-aggregation registered without a
``valid_through_ts`` does not have yet; those fall back to naming the node
they roll up.
"""
return preagg_label(
preagg.materialized_table_ref or preagg.node_revision.name,
preagg.grain_columns,
)


def _extract_dimension_refs_from_filters(
filters: list[str],
) -> list[tuple[str, str]]:
Expand Down Expand Up @@ -476,8 +504,6 @@ async def _plan_and_execute(self) -> DeploymentExecuteResult:
# Pre-aggregations still need reconciling on an otherwise-empty
# deploy: to register specs, or to delete external pre-aggs that were
# dropped from the spec.
from datajunction_server.database.preaggregation import PreAggregation

needs_preagg_reconcile = bool(
self.deployment_spec.preaggregations,
) or bool(
Expand Down Expand Up @@ -1118,10 +1144,11 @@ async def _reconcile_preaggregations(self) -> None:
remove EXTERNAL pre-aggs in this namespace that were dropped from it.

Runs after nodes/links/cubes so the referenced metrics and dimensions
exist. Skipped during dry runs (registration introspects the external
table and mutates state). Reuses the same core as POST /preaggs/register.
exist. A dry run registers nothing (registration introspects the
external table and mutates state) and previews instead, see
``_preview_preaggregations``. Reuses the same core as
POST /preaggs/register.
"""
from datajunction_server.database.preaggregation import PreAggregation
from datajunction_server.internal.preaggregations import (
register_external_preaggregations,
)
Expand Down Expand Up @@ -1175,39 +1202,11 @@ async def _reconcile_preaggregations(self) -> None:
)
return

declared_names = {spec.rendered_name for spec in specs}
existing_ids = {preagg.id for preagg in existing}
existing_names = {preagg.name for preagg in existing}

# Dry run: report the planned register/delete without mutating anything
# (no registration, so no external-table introspection either). Delete
# planning is name-level here; the wet run below deletes by row identity.
# Dry run: report the planned register/delete without mutating anything.
if self.dry_run:
for spec in specs:
self.deployed_results.append(
DeploymentResult(
name=spec.rendered_name,
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=(
DeploymentResult.Operation.UPDATE
if spec.rendered_name in existing_names
else DeploymentResult.Operation.CREATE
),
message="externally-built pre-aggregation",
),
)
for preagg in existing:
if preagg.name not in declared_names:
self.deployed_results.append(
DeploymentResult(
name=preagg.name or f"preaggregation:{preagg.id}",
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=DeploymentResult.Operation.DELETE,
message="dropped from spec",
),
)
await self._preview_preaggregations(specs, existing)
return

query_service_client = self.context.query_service_client
Expand Down Expand Up @@ -1238,7 +1237,6 @@ async def _reconcile_preaggregations(self) -> None:
self.session,
query_service_client,
request_headers,
name=spec.rendered_name,
metrics=spec.rendered_metrics,
dimensions=spec.rendered_dimensions,
table=ExternalPreAggTable(
Expand All @@ -1254,7 +1252,10 @@ async def _reconcile_preaggregations(self) -> None:
for preagg in created:
self.deployed_results.append(
DeploymentResult(
name=spec.rendered_name,
name=preagg_label(
spec.table_ref,
preagg.grain_columns,
),
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=(
Expand All @@ -1270,7 +1271,7 @@ async def _reconcile_preaggregations(self) -> None:
DJError(
code=ErrorCode.INVALID_ARGUMENTS_TO_FUNCTION,
message=(
f"Pre-aggregation '{spec.rendered_name}': {exc.message}"
f"Pre-aggregation on '{spec.table_ref}': {exc.message}"
),
),
)
Expand All @@ -1291,7 +1292,7 @@ async def _reconcile_preaggregations(self) -> None:
if preagg.id not in upserted_ids:
self.deployed_results.append(
DeploymentResult(
name=preagg.name or f"preaggregation:{preagg.id}",
name=existing_preagg_label(preagg),
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=DeploymentResult.Operation.DELETE,
Expand All @@ -1303,6 +1304,96 @@ async def _reconcile_preaggregations(self) -> None:
await self.session.flush()
logger.info("Reconciled %d pre-aggregation spec(s)", len(specs))

async def _preview_preaggregations(
self,
specs: list[PreAggSpec],
existing: list[PreAggregation],
) -> None:
"""
Report what pre-aggregation reconciliation would do, keyed exactly as the
apply keys it: the rows each spec resolves to.

The apply only learns those rows by registering, which a dry run cannot
do -- registration introspects the external table through the query
service. Everything the row identity depends on is pure decomposition,
though, so the preview resolves each spec to its grain groups and their
matching rows itself and reports the same creates and updates the apply
will perform, rather than diffing user-supplied handles (which reported a
single update where the apply did a create plus a delete).

A spec that cannot be resolved yet -- most often because this very deploy
is what creates its metrics -- is reported once as UNKNOWN instead of
being guessed at. Because such a spec might still have claimed an
existing row, the leftover rows are then reported as UNKNOWN too: a
preview must never promise a delete the apply will not perform.
"""
from datajunction_server.internal.preaggregations import (
resolve_external_preagg_targets,
)

claimed_ids: set[int] = set()
unresolved = False
for spec in specs:
try:
targets = await resolve_external_preagg_targets(
self.session,
metrics=spec.rendered_metrics,
dimensions=spec.rendered_dimensions,
)
except DJException as exc:
unresolved = True
self.deployed_results.append(
DeploymentResult(
name=spec.table_ref,
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=DeploymentResult.Operation.UNKNOWN,
message=(
f"cannot be previewed until its metrics and "
f"dimensions exist: {exc.message}"
),
),
)
continue
for target in targets:
if target.existing:
claimed_ids.add(target.existing.id)
self.deployed_results.append(
DeploymentResult(
name=preagg_label(spec.table_ref, target.grain_columns),
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=(
DeploymentResult.Operation.UPDATE
if target.existing
else DeploymentResult.Operation.CREATE
),
message="externally-built pre-aggregation",
),
)

for preagg in existing:
if preagg.id in claimed_ids:
continue
self.deployed_results.append(
DeploymentResult(
name=existing_preagg_label(preagg),
deploy_type=DeploymentResult.Type.PREAGG,
status=DeploymentResult.Status.SUCCESS,
operation=(
DeploymentResult.Operation.UNKNOWN
if unresolved
else DeploymentResult.Operation.DELETE
),
message=(
"may be dropped from spec, depending on the "
"pre-aggregation(s) that could not be previewed"
if unresolved
else "dropped from spec"
),
),
)

async def _setup_owners(self):
"""
Validate that all owners defined in the deployment spec exist.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
direction api -> internal.
"""

from dataclasses import dataclass
from typing import cast

from sqlalchemy.ext.asyncio import AsyncSession
Expand All @@ -28,6 +29,7 @@
compute_grain_group_hash,
compute_preagg_hash,
get_measure_identities,
measure_identity_token,
)
from datajunction_server.errors import DJInvalidInputException
from datajunction_server.models.decompose import PreAggMeasure
Expand Down Expand Up @@ -120,7 +122,6 @@ async def register_external_preaggregations(
query_service_client: QueryServiceClient,
request_headers: dict[str, str],
*,
name: str | None,
metrics: list[str],
dimensions: list[str],
table: ExternalPreAggTable,
Expand Down Expand Up @@ -318,7 +319,6 @@ async def register_external_preaggregations(
existing.columns = columns
existing.sql = grain_group.sql
existing.strategy = MaterializationStrategy.EXTERNAL
existing.name = name
preagg = existing
else:
preagg = PreAggregation(
Expand All @@ -334,7 +334,6 @@ async def register_external_preaggregations(
grain_measures,
),
strategy=MaterializationStrategy.EXTERNAL,
name=name,
)
session.add(preagg)
created_preaggs.append(preagg)
Expand All @@ -354,3 +353,66 @@ async def register_external_preaggregations(

await session.flush()
return created_preaggs


@dataclass
class PreAggTarget:
"""
One grain group a pre-aggregation spec resolves to, together with the
existing row (if any) that registering it would land on.
"""

grain_columns: list[str]
existing: PreAggregation | None


async def resolve_external_preagg_targets(
session: AsyncSession,
*,
metrics: list[str],
dimensions: list[str],
) -> list[PreAggTarget]:
"""
Work out which pre-aggregation rows a spec maps to, without introspecting
the external table.

``register_external_preaggregations`` only discovers this by registering: it
decomposes the metrics into grain groups and, for each, looks for a row that
already covers the same measures at the same grain. Neither step needs the
query service -- only binding the physical columns does -- so a deploy dry
run can reproduce the lookup exactly and preview the very operations the
apply will perform, instead of guessing from a user-supplied handle.
"""
measures_result = await build_measures_sql(
session=session,
metrics=metrics,
dimensions=dimensions,
dialect=Dialect.SPARK,
use_materialized=False,
)
assert_dimension_refs_are_role_qualified(measures_result, dimensions)

targets: list[PreAggTarget] = []
for grain_group in measures_result.grain_groups:
parent_node = measures_result.ctx.nodes.get(grain_group.parent_name)
if not parent_node or not parent_node.current: # pragma: no cover
continue
grain_columns = list(measures_result.requested_dimensions)
targets.append(
PreAggTarget(
grain_columns=grain_columns,
existing=await PreAggregation.find_matching(
session=session,
node_revision_id=parent_node.current.id,
grain_columns=grain_columns,
measure_identities={
measure_identity_token(
compute_expression_hash(component.expression),
component.normalized_aggregation,
)
for component in grain_group.components
},
),
),
)
return targets
Loading
Loading