Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
7bb3756
Add semi-additive metric declaration
betodealmeida Sep 4, 2026
19bcbd6
Attach semi-additive rules to metric components
betodealmeida Aug 24, 2026
2dd882e
Support semi-additive metric collapse SQL
betodealmeida Aug 24, 2026
6b37ea6
Support semi-additive derived metrics
betodealmeida Aug 25, 2026
3a8c202
Make pre-agg matching semi-additive aware
betodealmeida Aug 25, 2026
9cae7a6
Guard semi-additive materialization paths
betodealmeida Aug 25, 2026
a2fd915
Fix cube materialization test isolation
betodealmeida Aug 25, 2026
161e255
Make _semi_additive_dimension_requested role-sensitive
betodealmeida Aug 27, 2026
ab7c6d8
Show in UI
betodealmeida Aug 28, 2026
5beb5b4
Create/edit fields
betodealmeida Aug 28, 2026
f0e41cf
Fix current value on edit
betodealmeida Aug 28, 2026
d94c413
Use reaggregate instead
betodealmeida Aug 28, 2026
8764f1b
Prepare semi-additive metrics for review
betodealmeida Sep 3, 2026
f043851
Fix generated schema and reaggregate test expectations
betodealmeida Sep 3, 2026
5b3ba3f
Increase coverage
betodealmeida Sep 3, 2026
01c2103
Fix reaggregate semantic fingerprint expectations
betodealmeida Sep 4, 2026
9ab6f79
Improvements
betodealmeida Sep 4, 2026
17e0ca9
More fixes
betodealmeida Sep 4, 2026
5b698bc
Increase coverage
betodealmeida Sep 8, 2026
6b9c811
Declare a metric's fixed_grain
robinld Sep 10, 2026
297a73a
Broadcast a metric's fixed_grain to the query grain
robinld Sep 10, 2026
9d1b371
Fix two paths that lose a metric-level declaration
robinld Sep 10, 2026
e6dfbfa
Pin what fixed_grain does with experimentation metrics
shangyian Sep 11, 2026
93e959c
Refuse a fixed_grain dimension omitted from the query
robinld Sep 11, 2026
023d957
Update XP metrics tests for the fixed_grain output-dimension refusal
robinld Sep 12, 2026
08d6cec
Cover ratios with matching fixed grains
robinld Sep 12, 2026
61aa5cb
Use readable full SQL assertions for fixed grain
robinld Sep 12, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,7 @@ Same shape for MoM (`month_code`), QoQ (`quarter_code`), YoY (`year`).
| `direction` | ❌ Optional | `higher_is_better` / `lower_is_better` / `neutral` | Indicates performance direction |
| `unit` | ❌ Optional | `dollar` / `unitless` / **⚠️ NOT `count`** | Server rejects `count` — use `unitless` |
| `mode` | ❌ Optional | `draft` / `published` | Default: `published` |
| `fixed_grain` | ❌ Optional | List of dimension names | Grain the aggregate is computed at; omit for query grain, `[]` for global |
| `required_dimensions` | ❌ Optional | List of dimension names | For time-based / windowed metrics |
| `owners` | ❌ Optional but strongly recommended | List of email addresses | Prefer team emails |

Expand Down
510 changes: 280 additions & 230 deletions datajunction-query/uv.lock

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
"""
Add reaggregate column to noderevision

Revision ID: rg0001reaggregate
Revises: cm0003dropowner
Create Date: 2026-08-24 00:00:00.000000+00:00
"""

import sqlalchemy as sa
from alembic import op

# revision identifiers, used by Alembic.
revision = "rg0001reaggregate"
down_revision = "cm0003dropowner"
branch_labels = None
depends_on = None


def upgrade():
op.add_column(
"noderevision",
sa.Column("reaggregate", sa.JSON(), nullable=True),
)


def downgrade():
op.drop_column("noderevision", "reaggregate")
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
"""
Add fixed_grain to node revisions

Revision ID: fg0001fixedgrain
Revises: rg0001reaggregate
Create Date: 2026-09-04 00:00:00.000000+00:00
"""

import sqlalchemy as sa
from alembic import op

# revision identifiers, used by Alembic.
revision = "fg0001fixedgrain"
down_revision = "rg0001reaggregate"
branch_labels = None
depends_on = None


def upgrade():
op.add_column(
"noderevision",
sa.Column("fixed_grain", sa.JSON(), nullable=True),
)


def downgrade():
op.drop_column("noderevision", "fixed_grain")
31 changes: 31 additions & 0 deletions datajunction-server/datajunction_server/api/cubes.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@
_reorder_partition_column_last,
build_combiner_sql_from_preaggs,
)
from datajunction_server.construction.build_v3.cube_matcher import (
validate_cube_reaggregate_materialization,
)
from datajunction_server.construction.build_v3.cte import strip_role_suffix
from datajunction_server.construction.dimensions import build_dimensions_from_cube_query
from datajunction_server.database.materialization import Materialization
Expand Down Expand Up @@ -201,6 +204,32 @@ def _build_metrics_spec(
return metrics


async def _validate_cube_reaggregate_materialization(
session: AsyncSession,
cube: Node,
) -> None:
"""
Validate materialization safety using full metric decomposition.
"""
if not cube.current: # pragma: no cover
return

from datajunction_server.construction.build_v3.builder import setup_build_context

ctx = await setup_build_context(
session=session,
metrics=cube.current.cube_node_metrics,
dimensions=cube.current.cube_node_dimensions,
filters=cube.current.cube_filters or None,
dialect=Dialect.SPARK,
use_materialized=False,
)
validate_cube_reaggregate_materialization(
cube.current,
decomposed_metrics=ctx.decomposed_metrics,
)


@router.get("/cubes", name="Get all Cubes")
async def get_all_cubes(
*,
Expand Down Expand Up @@ -288,6 +317,7 @@ async def cube_materialization_info(
message=f"Cube node `{name}` does not exist.",
http_status_code=404,
)
await _validate_cube_reaggregate_materialization(session, node)
temporal_partitions = node.current.temporal_partition_columns() # type: ignore
if len(temporal_partitions) != 1:
raise DJInvalidInputException(
Expand Down Expand Up @@ -520,6 +550,7 @@ async def materialize_cube(
message=f"Cube '{name}' has no current revision",
http_status_code=HTTPStatus.NOT_FOUND,
)
await _validate_cube_reaggregate_materialization(session, node)

cube_tps = cube_revision.temporal_partition_columns()

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -353,9 +353,8 @@ async def batch_load_extracted_measures(
list(nr_to_name.values()),
)

# 3) Bulk-load Node objects for the full ancestor set. We only need
# name, type, and current.query — everything else is noloaded so this
# is a narrow query. noload(Node.created_by/Node.tags) are safe because
# 3) Bulk-load Node objects for the full ancestor set.
# noload(Node.created_by/Node.tags) are safe because
# MetricComponentExtractor never reads them.
nodes_cache: dict[str, DBNode] = {}
if all_names:
Expand All @@ -372,10 +371,16 @@ async def batch_load_extracted_measures(
noload(DBNode.tags),
joinedload(DBNode.current).options(
noload(DBNodeRevision.created_by),
# Must name every field `_build_metric_data_from_cache`
# reads. An omitted one is deferred, and reading a
# deferred field here raises MissingGreenlet, which the
# loop below turns into a null `extractedMeasures`.
load_only(
DBNodeRevision.fixed_grain,
DBNodeRevision.id,
DBNodeRevision.name,
DBNodeRevision.query,
DBNodeRevision.reaggregate,
),
),
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,15 @@
MetricComponent as MetricComponent_,
)
from datajunction_server.models.node import MetricDirection as MetricDirection_
from datajunction_server.models.reaggregate import (
DimensionReaggregateRule as DimensionReaggregateRule_,
ReaggregateSpec as ReaggregateSpec_,
ReaggregationFunction as ReaggregationFunction_,
)

MetricDirection = strawberry.enum(MetricDirection_)
Aggregability = strawberry.enum(Aggregability_)
ReaggregationFunction = strawberry.enum(ReaggregationFunction_)


@strawberry.type
Expand All @@ -32,6 +38,17 @@ class Unit:
abbreviation: str | None


@strawberry.experimental.pydantic.type(
model=DimensionReaggregateRule_,
all_fields=True,
)
class DimensionReaggregateRule: ...


@strawberry.experimental.pydantic.type(model=ReaggregateSpec_, all_fields=True)
class ReaggregateSpec: ...


@strawberry.experimental.pydantic.type(model=AggregationRule_, all_fields=True)
class AggregationRule: ...

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,9 @@
)
from datajunction_server.api.graphql.scalars.metricmetadata import (
DecomposedMetric,
DimensionReaggregateRule,
MetricMetadata,
ReaggregateSpec,
)
from datajunction_server.api.graphql.scalars.user import User
from datajunction_server.api.graphql.utils import extract_fields
Expand All @@ -45,6 +47,7 @@
from datajunction_server.models.node import NodeMode as NodeMode_
from datajunction_server.models.node import NodeStatus as NodeStatus_
from datajunction_server.models.node import NodeType as NodeType_
from datajunction_server.models.reaggregate import parse_reaggregate_spec
from datajunction_server.sql.parsing.backends.antlr4 import ast, parse

NodeType = strawberry.enum(NodeType_)
Expand Down Expand Up @@ -410,6 +413,37 @@ def materializations(
# Only metrics will have these fields
required_dimensions: list[Column] | None = None

@strawberry.field
def fixed_grain(self, root: DBNodeRevision) -> list[str] | None:
"""
The grain this metric's aggregate is computed at.
"""
if root.type != NodeType.METRIC:
return None
return root.fixed_grain

@strawberry.field
def reaggregate(self, root: DBNodeRevision) -> ReaggregateSpec | None:
"""
Metric reaggregation declaration.
"""
if root.type != NodeType.METRIC:
return None
spec = parse_reaggregate_spec(root.reaggregate)
if not spec:
return None
return ReaggregateSpec(
fn=spec.fn, # type: ignore
weight=spec.weight,
rules=[
DimensionReaggregateRule(
dimension=rule.dimension,
fn=rule.fn, # type: ignore
)
for rule in spec.rules
],
)

@strawberry.field
def primary_key(self, root: DBNodeRevision) -> list[str]:
"""
Expand Down
27 changes: 27 additions & 0 deletions datajunction-server/datajunction_server/api/graphql/schema.graphql
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ enum Aggregability {
type AggregationRule {
type: Aggregability!
level: [String!]
reaggregate: DimensionReaggregateRule
fixedGrain: [String!]
}

type Attribute {
Expand Down Expand Up @@ -149,6 +151,11 @@ type DimensionLink {
defaultValue: String
}

type DimensionReaggregateRule {
dimension: String!
fn: ReaggregationFunction!
}

type Engine {
name: String!
version: String!
Expand Down Expand Up @@ -397,6 +404,8 @@ type NodeRevision {
dimensionLinks: [DimensionLink!]!
availability: AvailabilityState
materializations: [MaterializationConfig!]
fixedGrain: [String!]
reaggregate: ReaggregateSpec
primaryKey: [String!]!
metricMetadata: MetricMetadata
isDerivedMetric: Boolean!
Expand Down Expand Up @@ -716,6 +725,24 @@ type Query {
listNamespaces: [Namespace!]!
}

type ReaggregateSpec {
fn: ReaggregationFunction
weight: String
rules: [DimensionReaggregateRule!]!
}

enum ReaggregationFunction {
AUTO
NONE
SUM
AVG
WEIGHTED_AVG
LAST_VALUE
FIRST_VALUE
MIN
MAX
}

type SemanticEntity {
name: String!

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
)
from datajunction_server.construction.build_v3.decomposition import (
decompose_and_group_metrics,
missing_fixed_grain_dimensions,
missing_reaggregate_dimensions,
)
from datajunction_server.construction.build_v3.dimensions import parse_dimension_ref
from datajunction_server.construction.build_v3.filters import (
Expand Down Expand Up @@ -320,6 +322,20 @@ async def setup_build_context(

# Add dimensions referenced in metric expressions (e.g., LAG ORDER BY)
add_dimensions_from_metric_expressions(ctx, ctx.decomposed_metrics)
output_dimensions_after_expression_scan = list(ctx.dimensions)
internal_reaggregate_dimensions = missing_reaggregate_dimensions(
ctx.decomposed_metrics.values(),
output_dimensions_after_expression_scan,
)
# Appended only so `load_nodes` pulls in the join path: `ctx.dimensions` is
# reset below and the dimension travels on the grain group instead.
internal_fixed_grain_dimensions = missing_fixed_grain_dimensions(
ctx.decomposed_metrics.values(),
output_dimensions_after_expression_scan,
)
for dimension in internal_reaggregate_dimensions + internal_fixed_grain_dimensions:
if dimension not in ctx.dimensions:
ctx.dimensions.append(dimension)

# A second load_nodes pass is needed when either:
# 1. metric expressions introduced dimension nodes not yet in ctx.nodes, OR
Expand All @@ -333,8 +349,16 @@ async def setup_build_context(
}
missing_dim_nodes = dim_roots_after - ctx.nodes.keys()
internally_added_roots = dim_roots_after - dim_roots_before_load
if missing_dim_nodes or internally_added_roots:
await load_nodes(ctx)
try:
if (
missing_dim_nodes
or internally_added_roots
or internal_reaggregate_dimensions
or internal_fixed_grain_dimensions
):
await load_nodes(ctx)
finally:
ctx.dimensions = output_dimensions_after_expression_scan

# Classify filters into dimension filters (WHERE) and metric filters (HAVING)
# This MUST happen AFTER all nodes are loaded so we can correctly identify
Expand Down
Loading
Loading