From 6db2d92350e0e45e66029bf84ea8218c51e923b0 Mon Sep 17 00:00:00 2001 From: Juraj Majerik Date: Fri, 31 Jul 2026 11:22:01 +0200 Subject: [PATCH 1/2] feat(experiments): extend metric-events precomputation to count/sum mean metrics --- .../backend/hogql_queries/LAZY_COMPUTATION.md | 2 + .../experiment_mean_query_builder.py | 102 ++++++++++++-- .../hogql_queries/experiment_query_builder.py | 19 ++- .../hogql_queries/experiment_query_runner.py | 56 +++++--- .../__snapshots__/test_mean_metric.ambr | 119 +++++----------- ...iment_mean_metric_events_preaggregation.py | 128 ++++++++++++++++++ 6 files changed, 311 insertions(+), 115 deletions(-) create mode 100644 products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py diff --git a/products/experiments/backend/hogql_queries/LAZY_COMPUTATION.md b/products/experiments/backend/hogql_queries/LAZY_COMPUTATION.md index 483fe62cbefd..48d66ddb0e8a 100644 --- a/products/experiments/backend/hogql_queries/LAZY_COMPUTATION.md +++ b/products/experiments/backend/hogql_queries/LAZY_COMPUTATION.md @@ -265,3 +265,5 @@ GROUP BY variant 5. **Fallback**: If precomputation fails or isn't ready, the query falls back to scanning the events table directly. 6. **Future: GROUP BY (entity_id, variant)**: Currently the query groups by entity_id only and resolves the variant during precomputation (e.g., argMin for first-seen). This bakes the variant handling strategy into the stored data, so switching between "first seen" and "exclude multiple" requires recomputation. A better approach is to GROUP BY (entity_id, variant) so we store one row per user-variant pair, then resolve multiple variants at query time. + +7. **Metric events**: The same lazy computation system also precomputes the metric_events half into `experiment_metric_events_preaggregated`, for ordered funnel metrics (one row per matching event with step indicators in `steps`) and count/sum mean metrics (one row per matching event with its value in `numeric_value`). Eligibility is gated by `_metric_events_precompute_applicable()` in `experiment_query_runner.py`; breakdowns, CUPED, and data warehouse sources fall back to a direct scan of the metric_events CTE. diff --git a/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py b/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py index ab858d337b85..926752c2424a 100644 --- a/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py +++ b/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py @@ -6,6 +6,7 @@ from posthog.hogql.parser import parse_expr, parse_select from products.experiments.backend.hogql_queries.base_query_utils import ( + event_or_action_to_filter, is_session_property_metric, validate_session_property, ) @@ -126,20 +127,45 @@ def get_mean_query_common_ctes(self) -> str: entity_metric_selects = """ {value_agg} AS value""" + if self._b.metric_events_preaggregation_job_ids: + # Read from the precomputed table instead of scanning events. Only reachable + # for eligible metrics (no breakdowns/CUPED/DW/session properties), so the + # branches above never coexist with this one. + # Filter by experiment date range: jobs can cover broader time ranges than + # the experiment for cache reusability, so we must filter on read. The upper + # bound includes the conversion window since metric events can occur after + # experiment end, mirroring build_metric_predicate() on the direct path. + entity_id_cast = "toUUID(t.entity_id)" if self._b.entity_key == "person_id" else "t.entity_id" + metric_events_cte = f""" + metric_events AS ( + SELECT + {entity_id_cast} AS entity_id, + t.timestamp AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE t.job_id IN {{metric_events_job_ids}} + AND t.team_id = {{metric_events_team_id}} + AND t.timestamp >= {{metric_events_date_from}} + AND t.timestamp < {{metric_events_date_to}} + toIntervalSecond({{metric_events_conversion_window_seconds}}) + )""" + else: + metric_events_cte = """ + metric_events AS ( + SELECT + {entity_key} AS entity_id, + {metric_timestamp_field} AS timestamp, + {value_expr} AS value + -- breakdown columns added programmatically below + FROM {metric_table} + WHERE {metric_predicate} + )""" + return f""" exposures AS ( {{exposure_select_query}} ), - metric_events AS ( - SELECT - {{entity_key}} AS entity_id, - {{metric_timestamp_field}} AS timestamp, - {{value_expr}} AS value - -- breakdown columns added programmatically below - FROM {{metric_table}} - WHERE {{metric_predicate}} - ), + {metric_events_cte}, entity_metrics AS ( SELECT @@ -220,6 +246,15 @@ def get_mean_query_common_placeholders(self) -> dict: value_expr=self._b._build_windowed_metric_value_expr(cuped_pre_window_predicate) ) + if self._b.metric_events_preaggregation_job_ids: + placeholders["metric_events_job_ids"] = ast.Constant(value=self._b.metric_events_preaggregation_job_ids) + placeholders["metric_events_team_id"] = ast.Constant(value=self._b.team.id) + placeholders["metric_events_date_from"] = self._b.date_range_query.date_from_as_hogql() + placeholders["metric_events_date_to"] = self._b.date_range_query.date_to_as_hogql() + placeholders["metric_events_conversion_window_seconds"] = ast.Constant( + value=self._b._get_conversion_window_seconds() + ) + # Add join condition for data warehouse if source_info.kind == "datawarehouse": placeholders["join_condition"] = parse_expr( @@ -245,6 +280,55 @@ def get_session_property_placeholders(self) -> dict: "session_conversion_window_predicate": self._b._build_session_conversion_window_predicate(), } + def get_mean_metric_events_query_for_precomputation(self) -> tuple[str, dict[str, ast.Expr]]: + """ + Returns the SELECT query that the lazy computation system wraps in an + INSERT INTO experiment_metric_events_preaggregated. This is the write + path — it scans the events table and stores one row per matching metric + event with its per-event value in numeric_value. The value is already + coalesced to a non-null float by _build_value_expr(), so storing it in + the non-nullable numeric_value column is lossless. + + The query uses {time_window_min} and {time_window_max} placeholders filled + by the lazy computation system for each daily bucket. The experiment date + bounds must stay named placeholders (the caller declares experiment_date_to + a sentinel) rather than reusing _build_metric_predicate(), which bakes the + resolved dates into the AST — a running experiment's moving window end + would then change the job hash and defeat cache reuse. + + Returns: + Tuple of (query_string, placeholders_dict) + """ + assert isinstance(self._b.metric, ExperimentMeanMetric) + source = self._b.metric.source + assert isinstance(source, (ActionsNode, EventsNode)) + + query_string = """ + SELECT + {entity_key} AS entity_id, + timestamp AS timestamp, + uuid AS event_uuid, + `$session_id` AS session_id, + {value_expr} AS numeric_value + FROM events + WHERE timestamp >= {time_window_min} + AND timestamp < {time_window_max} + AND timestamp >= {experiment_date_from} + AND timestamp < {experiment_date_to} + toIntervalSecond({conversion_window_seconds}) + AND {metric_event_filter} + """ + + placeholders: dict[str, ast.Expr] = { + "entity_key": parse_expr(self._b.entity_key), + "value_expr": self._b._build_value_expr(), + "experiment_date_from": self._b.date_range_query.date_from_as_hogql(), + "experiment_date_to": self._b.date_range_query.date_to_as_hogql(), + "conversion_window_seconds": ast.Constant(value=self._b._get_conversion_window_seconds()), + "metric_event_filter": event_or_action_to_filter(self._b.team, source), + } + + return query_string, placeholders + def build_mean_query(self) -> ast.SelectQuery: """ Builds query for mean metrics (count, sum, avg, etc.) diff --git a/products/experiments/backend/hogql_queries/experiment_query_builder.py b/products/experiments/backend/hogql_queries/experiment_query_builder.py index 87782f66e694..adf95f5e3210 100644 --- a/products/experiments/backend/hogql_queries/experiment_query_builder.py +++ b/products/experiments/backend/hogql_queries/experiment_query_builder.py @@ -655,12 +655,13 @@ def get_exposure_query_for_precomputation(self) -> tuple[str, dict[str, ast.Expr """ return self._exposure_query_builder().precomputation_query() - def get_funnel_metric_events_query_for_precomputation(self) -> tuple[str, dict[str, ast.Expr]]: + def get_metric_events_query_for_precomputation(self) -> tuple[str, dict[str, ast.Expr]]: """ Returns the SELECT query that the lazy computation system wraps in an - INSERT INTO experiment_metric_events_preaggregated. This is the write - path — it scans the events table and stores one row per matching event - with step indicators packed into an Array(UInt8). + INSERT INTO experiment_metric_events_preaggregated, dispatched by metric + type. This is the write path — it scans the events table and stores one + row per matching event: funnel metrics pack step indicators into an + Array(UInt8), mean metrics store the per-event value in numeric_value. The query uses {time_window_min} and {time_window_max} placeholders filled by the lazy computation system for each daily bucket. @@ -668,6 +669,16 @@ def get_funnel_metric_events_query_for_precomputation(self) -> tuple[str, dict[s Returns: Tuple of (query_string, placeholders_dict) """ + match self.metric: + case ExperimentFunnelMetric(): + return self._funnel_query_builder().get_funnel_metric_events_query_for_precomputation() + case ExperimentMeanMetric(): + return self._mean_query_builder().get_mean_metric_events_query_for_precomputation() + case _: + raise NotImplementedError(f"Metric-events precomputation is not supported for {type(self.metric)}") + + def get_funnel_metric_events_query_for_precomputation(self) -> tuple[str, dict[str, ast.Expr]]: + """Funnel-specific write query; prefer get_metric_events_query_for_precomputation().""" return self._funnel_query_builder().get_funnel_metric_events_query_for_precomputation() def _build_variant_expr_for_mean(self) -> ast.Expr: diff --git a/products/experiments/backend/hogql_queries/experiment_query_runner.py b/products/experiments/backend/hogql_queries/experiment_query_runner.py index 86a65ff4e3a8..221ec9d19ea3 100644 --- a/products/experiments/backend/hogql_queries/experiment_query_runner.py +++ b/products/experiments/backend/hogql_queries/experiment_query_runner.py @@ -11,12 +11,15 @@ from rest_framework.exceptions import ValidationError from posthog.schema import ( + ActionsNode, CachedExperimentQueryResponse, + EventsNode, ExperimentActorsQuery, ExperimentBreakdownResult, ExperimentDataWarehouseNode, ExperimentFunnelMetric, ExperimentMeanMetric, + ExperimentMetricMathType, ExperimentQuery, ExperimentQueryResponse, ExperimentRatioMetric, @@ -48,7 +51,11 @@ ) from products.cohorts.backend.models.cohort import Cohort from products.experiments.backend.hogql_queries import MULTIPLE_VARIANT_KEY, get_baseline_variant_key -from products.experiments.backend.hogql_queries.base_query_utils import experiment_window, experiment_window_end +from products.experiments.backend.hogql_queries.base_query_utils import ( + experiment_window, + experiment_window_end, + is_session_property_metric, +) from products.experiments.backend.hogql_queries.cuped_config import get_cuped_config from products.experiments.backend.hogql_queries.error_handling import experiment_error_handler from products.experiments.backend.hogql_queries.experiment_query_builder import ( @@ -349,12 +356,12 @@ def _ensure_exposures_precomputed(self, builder: ExperimentQueryBuilder) -> Lazy def _ensure_metric_events_precomputed(self, builder: ExperimentQueryBuilder) -> LazyComputationResult: """ - Ensures lazy-computed funnel metric event data exists for this experiment. + Ensures lazy-computed metric event data exists for this experiment. - Stores one row per matching event with step indicators in the - experiment_metric_events_preaggregated table. + Stores one row per matching event in the experiment_metric_events_preaggregated + table: funnel metrics store step indicators, mean metrics the per-event value. """ - query_string, placeholders = builder.get_funnel_metric_events_query_for_precomputation() + query_string, placeholders = builder.get_metric_events_query_for_precomputation() if not self.experiment.start_date: raise ValidationError("Experiment must have a start date for lazy computation") @@ -362,7 +369,7 @@ def _ensure_metric_events_precomputed(self, builder: ExperimentQueryBuilder) -> date_from = self.experiment.start_date date_to = experiment_window_end(self.experiment, self.as_of) - # Extend time range by conversion window — funnel step events can occur after experiment end + # Extend time range by conversion window — metric events can occur after experiment end conversion_window_seconds = builder._get_conversion_window_seconds() if conversion_window_seconds > 0: date_to = date_to + timedelta(seconds=conversion_window_seconds) @@ -423,14 +430,26 @@ def _precompute_skip_reason(self) -> Optional[str]: return None # precompute was attempted; a direct path means the build failed / wasn't ready def _metric_events_precompute_applicable(self) -> bool: - """Metric-events precompute only supports ordered funnels without breakdowns, CUPED, or data warehouse.""" - return ( - isinstance(self.metric, ExperimentFunnelMetric) - and (self.metric.funnel_order_type or "ordered") == "ordered" - and not self._get_breakdowns_for_builder() - and not self.cuped_config.enabled - and not self.is_data_warehouse_query - ) + """ + Metric-events precompute supports ordered funnels and count/sum-style mean + metrics, in both cases without breakdowns, CUPED, or data warehouse sources. + """ + if self._get_breakdowns_for_builder() or self.cuped_config.enabled or self.is_data_warehouse_query: + return False + if isinstance(self.metric, ExperimentFunnelMetric): + return (self.metric.funnel_order_type or "ordered") == "ordered" + if isinstance(self.metric, ExperimentMeanMetric): + source = self.metric.source + if not isinstance(source, (EventsNode, ActionsNode)): + return False + # Session-property means aggregate via a per-session dedup CTE that the + # precomputed table can't feed; ID-valued math (unique session/DAU/group) + # and HogQL expressions don't fit the Float64 numeric_value column. + if is_session_property_metric(source): + return False + math_type = getattr(source, "math", None) or ExperimentMetricMathType.TOTAL + return math_type in (ExperimentMetricMathType.TOTAL, ExperimentMetricMathType.SUM) + return False def _get_experiment_query(self) -> ast.SelectQuery: """ @@ -495,10 +514,11 @@ def _get_experiment_query(self) -> ast.SelectQuery: }, ) - # Precompute metric events for ordered funnel metrics. CUPED extends the - # funnel scan back by `lookback_days` to source the pre-exposure covariate; - # the precomputed metric_events table only covers the experiment window, so - # skip precomputation here and let the builder issue a fresh scan. + # Precompute metric events for eligible metrics (ordered funnels, count/sum + # means). CUPED extends the metric scan back by `lookback_days` to source the + # pre-exposure covariate; the precomputed metric_events table only covers the + # experiment window, so skip precomputation here and let the builder issue a + # fresh scan. if self._metric_events_precompute_applicable(): try: with tags_context( diff --git a/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr b/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr index 5d2412e84a4d..0524cf6d38a0 100644 --- a/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr +++ b/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr @@ -84,18 +84,11 @@ WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), lessOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-10 00:00:00', 'UTC')))) GROUP BY entity_id), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(replaceRegexpAll(nullIf(nullIf(JSONExtractRaw(events.properties, 'amount'), ''), 'null'), '^"|"$', ''), 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2020-01-10 00:00:00', 'UTC')), toIntervalSecond(86400))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-10 00:00:00', 'UTC')), toIntervalSecond(86400))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -490,18 +483,11 @@ GROUP BY entity_id HAVING lessOrEquals(plus(min(toTimeZone(t.first_exposure_time, 'UTC')), toIntervalSecond(604800)), parseDateTime64BestEffortOrNull('2020-01-15 12:00:00', 6, 'UTC'))), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(replaceRegexpAll(nullIf(nullIf(JSONExtractRaw(events.properties, 'amount'), ''), 'null'), '^"|"$', ''), 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 00:00:00', 'UTC')), toIntervalSecond(604800))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 00:00:00', 'UTC')), toIntervalSecond(604800))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -628,18 +614,11 @@ WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), lessOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')))) GROUP BY entity_id), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(1, 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -776,18 +755,11 @@ WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), lessOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')))) GROUP BY entity_id), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(replaceRegexpAll(nullIf(nullIf(JSONExtractRaw(events.properties, 'amount'), ''), 'null'), '^"|"$', ''), 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -924,18 +896,11 @@ WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), lessOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')))) GROUP BY entity_id), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(replaceRegexpAll(nullIf(nullIf(JSONExtractRaw(events.properties, 'amount'), ''), 'null'), '^"|"$', ''), 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -1446,18 +1411,11 @@ WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), lessOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')))) GROUP BY entity_id), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(replaceRegexpAll(nullIf(nullIf(JSONExtractRaw(events.properties, 'amount'), ''), 'null'), '^"|"$', ''), 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -1574,18 +1532,11 @@ WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), lessOrEquals(t.first_exposure_time, assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')))) GROUP BY entity_id), metric_events AS - (SELECT if(not(empty(events__override.distinct_id)), events__override.person_id, events.person_id) AS entity_id, - toTimeZone(events.timestamp, 'UTC') AS timestamp, - coalesce(accurateCastOrNull(replaceRegexpAll(nullIf(nullIf(JSONExtractRaw(events.properties, 'amount'), ''), 'null'), '^"|"$', ''), 'Float64'), 0) AS value - FROM events - LEFT OUTER JOIN - (SELECT argMax(person_distinct_id_overrides.person_id, person_distinct_id_overrides.version) AS person_id, - person_distinct_id_overrides.distinct_id AS distinct_id - FROM person_distinct_id_overrides - WHERE equals(person_distinct_id_overrides.team_id, 99999) - GROUP BY person_distinct_id_overrides.distinct_id - HAVING equals(argMax(person_distinct_id_overrides.is_deleted, person_distinct_id_overrides.version), 0) SETTINGS optimize_aggregation_in_order=1) AS events__override ON equals(events.distinct_id, events__override.distinct_id) - WHERE and(equals(events.team_id, 99999), greaterOrEquals(events.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(events.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))), equals(events.event, 'purchase'))), + (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, + toTimeZone(t.timestamp, 'UTC') AS timestamp, + t.numeric_value AS value + FROM experiment_metric_events_preaggregated AS t + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, diff --git a/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py b/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py new file mode 100644 index 000000000000..194b05390506 --- /dev/null +++ b/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py @@ -0,0 +1,128 @@ +from datetime import UTC, datetime + +from unittest.mock import patch + +from django.test import override_settings + +from parameterized import parameterized + +from posthog.schema import EventsNode, ExperimentMeanMetric, ExperimentMetricMathType, ExperimentQuery, IntervalType + +from posthog.hogql_queries.utils.query_date_range import QueryDateRange + +from products.experiments.backend.hogql_queries.base_query_utils import experiment_window +from products.experiments.backend.hogql_queries.experiment_query_builder import ( + ExperimentQueryBuilder, + get_exposure_config_params_for_builder, +) +from products.experiments.backend.hogql_queries.experiment_query_runner import ExperimentQueryRunner +from products.experiments.backend.hogql_queries.exposure_query_logic import get_entity_key +from products.experiments.backend.hogql_queries.test.experiment_query_runner.base import ExperimentQueryRunnerBaseTest + + +@override_settings(IN_UNIT_TESTING=True) +class TestExperimentMeanMetricEventsPreaggregation(ExperimentQueryRunnerBaseTest): + def _build_lazy_computation_builder( + self, experiment, feature_flag, metric, as_of: datetime | None = None + ) -> ExperimentQueryBuilder: + exposure_config, multiple_variant_handling, filter_test_accounts = get_exposure_config_params_for_builder( + experiment.exposure_criteria + ) + as_of = as_of if as_of is not None else datetime.now(UTC) + date_range = experiment_window(experiment, self.team, as_of) + return ExperimentQueryBuilder( + team=self.team, + feature_flag_key=feature_flag.key, + exposure_config=exposure_config, + filter_test_accounts=filter_test_accounts, + multiple_variant_handling=multiple_variant_handling, + variants=[v["key"] for v in feature_flag.variants], + date_range_query=QueryDateRange( + date_range=date_range, + team=self.team, + interval=IntervalType.DAY, + now=as_of, + ), + entity_key=get_entity_key(feature_flag.filters.get("aggregation_group_type_index")), + metric=metric, + ) + + def _build_runner(self, experiment, metric: ExperimentMeanMetric, as_of: datetime | None = None): + query = ExperimentQuery(experiment_id=experiment.id, kind="ExperimentQuery", metric=metric) + return ExperimentQueryRunner(query=query, team=self.team, as_of=as_of) + + @patch("products.analytics_platform.backend.lazy_computation.lazy_computation_executor.sync_execute") + def test_mean_metric_events_precomputation_hash_ignores_moving_experiment_end(self, mock_sync_execute): + feature_flag = self.create_feature_flag(key="stable-mean-metric-events-hash") + # Running experiment (no end_date) so the window end follows as_of. The mean build + # query must keep the experiment date bounds as sentinel placeholders — resolving + # them into the query (e.g. via _build_metric_predicate) would change the job hash + # on every load and silently defeat cache reuse. + experiment = self.create_experiment( + feature_flag=feature_flag, + start_date=datetime(2024, 1, 1), + end_date=None, + ) + experiment.end_date = None + experiment.save(update_fields=["end_date"]) + metric = ExperimentMeanMetric( + source=EventsNode(event="purchase", math=ExperimentMetricMathType.SUM, math_property="amount") + ) + + first_as_of = datetime(2024, 1, 5, 12, 0, tzinfo=UTC) + second_as_of = datetime(2024, 1, 5, 12, 30, tzinfo=UTC) + + first_result = self._build_runner(experiment, metric, as_of=first_as_of)._ensure_metric_events_precomputed( + self._build_lazy_computation_builder(experiment, feature_flag, metric, as_of=first_as_of) + ) + second_result = self._build_runner(experiment, metric, as_of=second_as_of)._ensure_metric_events_precomputed( + self._build_lazy_computation_builder(experiment, feature_flag, metric, as_of=second_as_of) + ) + + assert first_result.ready is True + assert second_result.ready is True + assert first_result.job_ids == second_result.job_ids + assert mock_sync_execute.call_count == len(first_result.job_ids) + + @parameterized.expand( + [ + ("count_default_math", EventsNode(event="purchase"), True), + ( + "sum", + EventsNode(event="purchase", math=ExperimentMetricMathType.SUM, math_property="amount"), + True, + ), + ( + "avg_not_yet_allowlisted", + EventsNode(event="purchase", math=ExperimentMetricMathType.AVG, math_property="amount"), + False, + ), + ( + "unique_session_id_valued", + EventsNode(event="purchase", math=ExperimentMetricMathType.UNIQUE_SESSION), + False, + ), + ( + "hogql_user_expression", + EventsNode(event="purchase", math=ExperimentMetricMathType.HOGQL, math_hogql="sum(properties.amount)"), + False, + ), + ( + "session_property", + EventsNode(event="purchase", math=ExperimentMetricMathType.SUM, math_property="$session_duration"), + False, + ), + ] + ) + def test_mean_metric_events_precompute_gate(self, _name, source, applicable): + feature_flag = self.create_feature_flag(key="mean-metric-events-gate") + experiment = self.create_experiment( + feature_flag=feature_flag, + start_date=datetime(2024, 1, 1), + end_date=datetime(2024, 1, 10), + ) + metric = ExperimentMeanMetric(source=source) + + runner = self._build_runner(experiment, metric) + + assert runner._metric_events_precompute_applicable() is applicable From b0c8b96f580ac89e97024986486ff4ce9bd7bda9 Mon Sep 17 00:00:00 2001 From: Juraj Majerik Date: Fri, 31 Jul 2026 15:00:33 +0200 Subject: [PATCH 2/2] fix(experiments): dedup replayed rows in the mean precompute read --- .../experiment_mean_query_builder.py | 7 +- .../__snapshots__/test_mean_metric.ambr | 49 +++++++--- ...iment_mean_metric_events_preaggregation.py | 92 ++++++++++++++++++- 3 files changed, 132 insertions(+), 16 deletions(-) diff --git a/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py b/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py index 926752c2424a..8857fa20ac6d 100644 --- a/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py +++ b/products/experiments/backend/hogql_queries/experiment_mean_query_builder.py @@ -135,18 +135,23 @@ def get_mean_query_common_ctes(self) -> str: # the experiment for cache reusability, so we must filter on read. The upper # bound includes the conversion window since metric events can occur after # experiment end, mirroring build_metric_predicate() on the direct path. + # GROUP BY collapses replayed rows by event identity: ReplacingMergeTree only + # dedups at merge time and this read doesn't use FINAL, so a re-applied build + # INSERT would otherwise double-count sums. Same defense as the exposures read; + # the funnel read skips it because funnel evaluation tolerates duplicate events. entity_id_cast = "toUUID(t.entity_id)" if self._b.entity_key == "person_id" else "t.entity_id" metric_events_cte = f""" metric_events AS ( SELECT {entity_id_cast} AS entity_id, t.timestamp AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t WHERE t.job_id IN {{metric_events_job_ids}} AND t.team_id = {{metric_events_team_id}} AND t.timestamp >= {{metric_events_date_from}} AND t.timestamp < {{metric_events_date_to}} + toIntervalSecond({{metric_events_conversion_window_seconds}}) + GROUP BY t.entity_id, t.timestamp, t.event_uuid )""" else: metric_events_cte = """ diff --git a/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr b/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr index 0524cf6d38a0..b9958c00dfdf 100644 --- a/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr +++ b/products/experiments/backend/hogql_queries/test/experiment_query_runner/__snapshots__/test_mean_metric.ambr @@ -86,9 +86,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-10 00:00:00', 'UTC')), toIntervalSecond(86400))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-10 00:00:00', 'UTC')), toIntervalSecond(86400)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -485,9 +488,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 00:00:00', 'UTC')), toIntervalSecond(604800))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 00:00:00', 'UTC')), toIntervalSecond(604800)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -616,9 +622,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -757,9 +766,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2024-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2024-01-15 12:00:00', 'UTC')), toIntervalSecond(0)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -898,9 +910,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -1413,9 +1428,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, @@ -1534,9 +1552,12 @@ metric_events AS (SELECT accurateCastOrNull(t.entity_id, 'UUID') AS entity_id, toTimeZone(t.timestamp, 'UTC') AS timestamp, - t.numeric_value AS value + any(t.numeric_value) AS value FROM experiment_metric_events_preaggregated AS t - WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0))))), + WHERE and(equals(t.team_id, 99999), in(t.job_id, ['00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000', '00000000-0000-0000-0000-000000000000']), equals(t.team_id, 99999), greaterOrEquals(t.timestamp, assumeNotNull(toDateTime('2020-01-01 00:00:00', 'UTC'))), less(t.timestamp, plus(assumeNotNull(toDateTime('2020-01-15 12:00:00', 'UTC')), toIntervalSecond(0)))) + GROUP BY t.entity_id, + toTimeZone(t.timestamp, 'UTC'), + t.event_uuid), entity_metrics AS (SELECT exposures.entity_id AS entity_id, exposures.variant AS variant, diff --git a/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py b/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py index 194b05390506..f14169c88b53 100644 --- a/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py +++ b/products/experiments/backend/hogql_queries/test/test_experiment_mean_metric_events_preaggregation.py @@ -1,13 +1,24 @@ from datetime import UTC, datetime +from typing import cast +from posthog.test.base import _create_event, _create_person from unittest.mock import patch from django.test import override_settings from parameterized import parameterized -from posthog.schema import EventsNode, ExperimentMeanMetric, ExperimentMetricMathType, ExperimentQuery, IntervalType +from posthog.schema import ( + EventsNode, + ExperimentMeanMetric, + ExperimentMetricMathType, + ExperimentQuery, + ExperimentQueryResponse, + IntervalType, +) +from posthog.clickhouse.client.execute import sync_execute +from posthog.clickhouse.preaggregation.experiment_metric_events_sql import SHARDED_EXPERIMENT_METRIC_EVENTS_TABLE from posthog.hogql_queries.utils.query_date_range import QueryDateRange from products.experiments.backend.hogql_queries.base_query_utils import experiment_window @@ -51,6 +62,85 @@ def _build_runner(self, experiment, metric: ExperimentMeanMetric, as_of: datetim query = ExperimentQuery(experiment_id=experiment.id, kind="ExperimentQuery", metric=metric) return ExperimentQueryRunner(query=query, team=self.team, as_of=as_of) + def _create_exposure_event(self, distinct_id, feature_flag, variant, timestamp): + _create_event( + team=self.team, + event="$feature_flag_called", + distinct_id=distinct_id, + timestamp=timestamp, + properties={ + f"$feature/{feature_flag.key}": variant, + "$feature_flag_response": variant, + "$feature_flag": feature_flag.key, + }, + ) + + def test_precomputed_result_tolerates_replayed_build_rows(self): + feature_flag = self.create_feature_flag(key="mean-metric-events-replay") + experiment = self.create_experiment( + feature_flag=feature_flag, + start_date=datetime(2024, 1, 1), + end_date=datetime(2024, 1, 10), + ) + metric = ExperimentMeanMetric( + source=EventsNode(event="purchase", math=ExperimentMetricMathType.SUM, math_property="amount") + ) + experiment.metrics = [metric.model_dump(mode="json")] + experiment.save() + + for variant, count in (("control", 3), ("test", 3)): + for i in range(count): + _create_person(distinct_ids=[f"{variant}_{i}"], team_id=self.team.pk) + self._create_exposure_event( + f"{variant}_{i}", feature_flag, variant, datetime(2024, 1, 2, 12, 0, tzinfo=UTC) + ) + _create_event( + team=self.team, + event="purchase", + distinct_id=f"{variant}_{i}", + timestamp=datetime(2024, 1, 3, 12, 0, tzinfo=UTC), + properties={"amount": 10 * (i + 1)}, + ) + + self._disable_precomputation() + direct_runner = self._build_runner(experiment, metric) + direct_result = cast(ExperimentQueryResponse, direct_runner.calculate()) + + # First precomputed run builds the jobs whose rows the second run will re-read. + self._enable_precomputation() + first_runner = self._build_runner(experiment, metric) + first_runner.calculate() + assert first_runner._metric_events_precomputed is True + + # Simulate a replayed build INSERT: every stored metric-event row appears twice + # under the same job. ReplacingMergeTree only collapses these at merge time, so + # the read must dedup by event identity or sums double. + table = SHARDED_EXPERIMENT_METRIC_EVENTS_TABLE() + rows_before = sync_execute( + f"SELECT count() FROM {table} WHERE team_id = %(team_id)s", {"team_id": self.team.pk} + )[0][0] + assert rows_before > 0 + sync_execute( + f"INSERT INTO {table} SELECT * FROM {table} WHERE team_id = %(team_id)s", + {"team_id": self.team.pk}, + ) + + precomputed_runner = self._build_runner(experiment, metric) + precomputed_result = cast(ExperimentQueryResponse, precomputed_runner.calculate()) + assert precomputed_runner._metric_events_precomputed is True + + assert direct_result.baseline is not None + assert precomputed_result.baseline is not None + assert precomputed_result.baseline.number_of_samples == direct_result.baseline.number_of_samples + assert precomputed_result.baseline.sum == direct_result.baseline.sum + assert direct_result.variant_results is not None + assert precomputed_result.variant_results is not None + assert precomputed_result.variant_results[0].sum == direct_result.variant_results[0].sum + assert ( + precomputed_result.variant_results[0].number_of_samples + == direct_result.variant_results[0].number_of_samples + ) + @patch("products.analytics_platform.backend.lazy_computation.lazy_computation_executor.sync_execute") def test_mean_metric_events_precomputation_hash_ignores_moving_experiment_end(self, mock_sync_execute): feature_flag = self.create_feature_flag(key="stable-mean-metric-events-hash")