From 142216b23df86e0bc7d766cc45af1038ce52210e Mon Sep 17 00:00:00 2001 From: Asher Cohen Date: Mon, 1 Jun 2026 15:22:48 -0700 Subject: [PATCH 1/3] Fix ClickHouse activity dedupe build --- analytics/macros/bounded_activity_graph.sql | 214 ------------------ .../read_models/activity_duplicate_groups.sql | 120 ++++++++++ .../activity_duplicate_matches.sql | 101 +++++++++ .../read_models/activity_source_records.sql | 140 ++++++++++++ .../models/read_models/deduped_activities.sql | 129 +++++------ .../read_model_microbatch.sql.test.ts | 82 ++++--- docs/production-incident-baseline.md | 42 ++++ entrypoint.sh | 2 +- 8 files changed, 515 insertions(+), 315 deletions(-) delete mode 100644 analytics/macros/bounded_activity_graph.sql create mode 100644 analytics/models/read_models/activity_duplicate_groups.sql create mode 100644 analytics/models/read_models/activity_duplicate_matches.sql create mode 100644 analytics/models/read_models/activity_source_records.sql diff --git a/analytics/macros/bounded_activity_graph.sql b/analytics/macros/bounded_activity_graph.sql deleted file mode 100644 index 30b44b9b80..0000000000 --- a/analytics/macros/bounded_activity_graph.sql +++ /dev/null @@ -1,214 +0,0 @@ -{% macro bounded_activity_graph() %} -{% set batch_start_value = config.get("__dbt_internal_microbatch_event_time_start") %} -{% set batch_end_value = config.get("__dbt_internal_microbatch_event_time_end") %} -{% if batch_start_value and batch_end_value %} -{% set batch_start = batch_start_value.strftime("%Y-%m-%d %H:%M:%S") %} -{% set batch_end = batch_end_value.strftime("%Y-%m-%d %H:%M:%S") %} -{% endif %} -active_activity AS ( - SELECT * - FROM {{ source('postgres_fitness', 'activity') }} FINAL - WHERE _peerdb_is_deleted = 0 - {% if batch_start_value and batch_end_value %} - AND started_at < toDateTime64('{{ batch_end }}', 6, 'UTC') - AND coalesce(ended_at, started_at + INTERVAL 12 HOUR) >= toDateTime64('{{ batch_start }}', 6, 'UTC') - {% endif %} -), - -{{ activity_dedup_graph() }} -{% endmacro %} - -{% macro activity_dedup_graph() %} -active_provider_priority AS ( - SELECT * - FROM {{ source('postgres_fitness', 'provider_priority') }} FINAL - WHERE _peerdb_is_deleted = 0 -), - -active_device_priority AS ( - SELECT * - FROM {{ source('postgres_fitness', 'device_priority') }} FINAL - WHERE _peerdb_is_deleted = 0 -), - -device_priority_match AS ( - SELECT - activity_id, - priority - FROM ( - SELECT - active_activity.id AS activity_id, - active_device_priority.priority AS priority, - row_number() OVER ( - PARTITION BY active_activity.id - ORDER BY length(active_device_priority.source_name_pattern) DESC - ) AS row_number - FROM active_activity - INNER JOIN active_device_priority - ON active_device_priority.provider_id = active_activity.provider_id - AND active_activity.source_name LIKE active_device_priority.source_name_pattern - ) - WHERE row_number = 1 -), - -ranked AS ( - SELECT - active_activity.id AS id, - active_activity.provider_id AS provider_id, - active_activity.user_id AS user_id, - active_activity.external_id AS external_id, - active_activity.activity_type AS activity_type, - active_activity.started_at AS started_at, - active_activity.ended_at AS ended_at, - active_activity.source_name AS source_name, - active_activity.name AS name, - active_activity.notes AS notes, - active_activity.timezone AS timezone, - active_activity.raw AS raw, - active_activity._peerdb_synced_at AS _peerdb_synced_at, - coalesce(device_priority_match.priority, active_provider_priority.priority, 100) AS priority - FROM active_activity - LEFT JOIN active_provider_priority - ON active_provider_priority.provider_id = active_activity.provider_id - LEFT JOIN device_priority_match - ON device_priority_match.activity_id = active_activity.id -), - -pairs AS ( - SELECT - left_activity.id AS id1, - right_activity.id AS id2 - FROM ranked AS left_activity - INNER JOIN ranked AS right_activity - ON left_activity.user_id = right_activity.user_id - AND toString(left_activity.id) < toString(right_activity.id) - AND dateDiff( - 'second', - greatest(left_activity.started_at, right_activity.started_at), - least( - coalesce(left_activity.ended_at, left_activity.started_at + INTERVAL 12 HOUR), - coalesce(right_activity.ended_at, right_activity.started_at + INTERVAL 12 HOUR) - ) - ) / nullIf(dateDiff( - 'second', - least(left_activity.started_at, right_activity.started_at), - greatest( - coalesce(left_activity.ended_at, left_activity.started_at + INTERVAL 12 HOUR), - coalesce(right_activity.ended_at, right_activity.started_at + INTERVAL 12 HOUR) - ) - ), 0) > 0.8 -), - -graph_edges AS ( - SELECT - id1 AS from_id, - id2 AS to_id - FROM pairs - UNION ALL - SELECT - id2 AS from_id, - id1 AS to_id - FROM pairs -), - -connected_components AS ( - SELECT - id AS activity_id, - id AS connected_activity_id, - [toString(id)] AS visited_activity_ids - FROM ranked - UNION ALL - SELECT - connected_components.activity_id AS activity_id, - graph_edges.to_id AS connected_activity_id, - arrayConcat(connected_components.visited_activity_ids, [toString(graph_edges.to_id)]) AS visited_activity_ids - FROM connected_components - INNER JOIN graph_edges - ON graph_edges.from_id = connected_components.connected_activity_id - WHERE NOT has(connected_components.visited_activity_ids, toString(graph_edges.to_id)) -), - -final_groups AS ( - SELECT - activity_id, - min(toString(connected_activity_id)) AS group_id - FROM connected_components - GROUP BY activity_id -), - -best AS ( - SELECT * - FROM ( - SELECT - final_groups.group_id AS group_id, - ranked.id AS canonical_id, - ranked.provider_id AS provider_id, - ranked.user_id AS user_id, - ranked.activity_type AS activity_type, - ranked.started_at AS started_at, - ranked.ended_at AS ended_at, - ranked.source_name AS source_name, - ranked.priority AS priority, - row_number() OVER ( - PARTITION BY final_groups.group_id - ORDER BY ranked.priority ASC, toString(ranked.id) ASC - ) AS row_number - FROM final_groups - INNER JOIN ranked - ON ranked.id = final_groups.activity_id - ) - WHERE row_number = 1 -), - -merged AS ( - SELECT - best.group_id AS group_id, - best.canonical_id AS id, - any(best.provider_id) AS provider_id, - any(best.user_id) AS user_id, - any(best.activity_type) AS activity_type, - min(ranked.started_at) AS started_at, - max(coalesce(ranked.ended_at, ranked.started_at + INTERVAL 12 HOUR)) AS ended_at, - any(best.source_name) AS source_name, - argMinIf(ranked.name, ranked.priority, ranked.name IS NOT NULL) AS name, - argMinIf(ranked.notes, ranked.priority, ranked.notes IS NOT NULL) AS notes, - argMinIf(ranked.timezone, ranked.priority, ranked.timezone IS NOT NULL) AS timezone, - argMinIf(ranked.raw, ranked.priority, ranked.raw IS NOT NULL) AS raw, - max(ranked._peerdb_synced_at) AS source_synced_at, - arraySort(groupUniqArray(ranked.provider_id)) AS source_providers, - groupArrayIf( - map('providerId', ranked.provider_id, 'externalId', ranked.external_id), - ranked.external_id IS NOT NULL AND ranked.external_id != '' - ) AS source_external_ids, - groupArray(ranked.id) AS member_activity_ids - FROM best - INNER JOIN final_groups - ON final_groups.group_id = best.group_id - INNER JOIN ranked - ON ranked.id = final_groups.activity_id - GROUP BY best.group_id, best.canonical_id -), - -current_activity AS ( - SELECT - id AS activity_id, - user_id, - activity_type, - name, - started_at, - ended_at, - source_synced_at - FROM merged -), - -activity_members AS ( - SELECT - id AS activity_id, - user_id, - started_at, - ended_at, - source_synced_at, - arrayJoin(member_activity_ids) AS member_activity_id - FROM merged -) -{% endmacro %} diff --git a/analytics/models/read_models/activity_duplicate_groups.sql b/analytics/models/read_models/activity_duplicate_groups.sql new file mode 100644 index 0000000000..f4faf7a110 --- /dev/null +++ b/analytics/models/read_models/activity_duplicate_groups.sql @@ -0,0 +1,120 @@ +{{ config( + materialized='incremental', + incremental_strategy='append', + engine='ReplacingMergeTree(refresh_version)', + order_by='activity_id', + query_settings={ + 'max_threads': 1 + } +) }} + +WITH source_records AS ( + SELECT activity_id + FROM {{ ref('activity_source_records') }} FINAL + WHERE is_deleted = 0 +), + +duplicate_links AS ( + SELECT + duplicate_matches.activity_id AS activity_id, + duplicate_matches.duplicate_activity_id AS linked_activity_id + FROM {{ ref('activity_duplicate_matches') }} AS duplicate_matches FINAL + WHERE duplicate_matches.is_deleted = 0 + + UNION ALL + + SELECT + duplicate_matches.duplicate_activity_id AS activity_id, + duplicate_matches.activity_id AS linked_activity_id + FROM {{ ref('activity_duplicate_matches') }} AS duplicate_matches FINAL + WHERE duplicate_matches.is_deleted = 0 +), + +duplicate_walk_rows AS ( + -- One bridge collapses provider chains such as Apple Health -> WHOOP -> Strava + -- into a single real-world activity group without path enumeration. + SELECT + activity_id, + activity_id AS connected_activity_id + FROM source_records + + UNION ALL + + SELECT + source_records.activity_id, + duplicate_links.linked_activity_id AS connected_activity_id + FROM source_records + INNER JOIN duplicate_links + ON duplicate_links.activity_id = source_records.activity_id + + UNION ALL + + SELECT + source_records.activity_id, + second_link.linked_activity_id AS connected_activity_id + FROM source_records + INNER JOIN duplicate_links AS first_link + ON first_link.activity_id = source_records.activity_id + INNER JOIN duplicate_links AS second_link + ON second_link.activity_id = first_link.linked_activity_id +), + +duplicate_walk AS ( + SELECT DISTINCT + activity_id, + connected_activity_id + FROM duplicate_walk_rows +), + +current_duplicate_groups AS ( + SELECT + activity_id, + min(toString(connected_activity_id)) AS group_id + FROM duplicate_walk + GROUP BY activity_id +), + +existing_duplicate_groups AS ( + {% if is_incremental() %} + SELECT activity_id + FROM {{ this }} FINAL + WHERE is_deleted = 0 + {% else %} + SELECT CAST(null, 'Nullable(UUID)') AS activity_id + WHERE 1 = 0 + {% endif %} +), + +stale_duplicate_groups AS ( + SELECT existing_duplicate_groups.activity_id + FROM existing_duplicate_groups + LEFT JOIN current_duplicate_groups + ON current_duplicate_groups.activity_id = existing_duplicate_groups.activity_id + WHERE current_duplicate_groups.activity_id IS null +), + +refresh_clock AS ( + SELECT + toUInt64(toUnixTimestamp64Nano(now64(9))) AS refresh_version, + now64(9) AS refreshed_at +) + +SELECT + activity_id, + group_id, + refresh_clock.refresh_version AS refresh_version, + 0 AS is_deleted, + refresh_clock.refreshed_at AS refreshed_at +FROM current_duplicate_groups +CROSS JOIN refresh_clock + +UNION ALL + +SELECT + assumeNotNull(activity_id) AS activity_id, + CAST(null, 'Nullable(String)') AS group_id, + refresh_clock.refresh_version AS refresh_version, + 1 AS is_deleted, + refresh_clock.refreshed_at AS refreshed_at +FROM stale_duplicate_groups +CROSS JOIN refresh_clock diff --git a/analytics/models/read_models/activity_duplicate_matches.sql b/analytics/models/read_models/activity_duplicate_matches.sql new file mode 100644 index 0000000000..700f13cec5 --- /dev/null +++ b/analytics/models/read_models/activity_duplicate_matches.sql @@ -0,0 +1,101 @@ +{{ config( + materialized='incremental', + incremental_strategy='append', + engine='ReplacingMergeTree(refresh_version)', + order_by='(activity_id, duplicate_activity_id)', + query_settings={ + 'max_threads': 1 + } +) }} + +WITH source_records AS ( + SELECT + activity_id, + user_id, + started_at, + coalesce(ended_at, started_at + INTERVAL 12 HOUR) AS ended_at + FROM {{ ref('activity_source_records') }} FINAL + WHERE is_deleted = 0 +), + +current_duplicate_matches AS ( + SELECT + left_activity.activity_id AS activity_id, + right_activity.activity_id AS duplicate_activity_id, + dateDiff( + 'second', + greatest(left_activity.started_at, right_activity.started_at), + least(left_activity.ended_at, right_activity.ended_at) + ) / nullIf(dateDiff( + 'second', + least(left_activity.started_at, right_activity.started_at), + greatest(left_activity.ended_at, right_activity.ended_at) + ), 0) AS overlap_ratio + FROM source_records AS left_activity + INNER JOIN source_records AS right_activity + ON left_activity.user_id = right_activity.user_id + AND toString(left_activity.activity_id) < toString(right_activity.activity_id) + AND dateDiff( + 'second', + greatest(left_activity.started_at, right_activity.started_at), + least(left_activity.ended_at, right_activity.ended_at) + ) / nullIf(dateDiff( + 'second', + least(left_activity.started_at, right_activity.started_at), + greatest(left_activity.ended_at, right_activity.ended_at) + ), 0) > 0.8 +), + +existing_duplicate_matches AS ( + {% if is_incremental() %} + SELECT + activity_id, + duplicate_activity_id + FROM {{ this }} FINAL + WHERE is_deleted = 0 + {% else %} + SELECT + CAST(null, 'Nullable(UUID)') AS activity_id, + CAST(null, 'Nullable(UUID)') AS duplicate_activity_id + WHERE 1 = 0 + {% endif %} +), + +stale_duplicate_matches AS ( + SELECT + existing_duplicate_matches.activity_id, + existing_duplicate_matches.duplicate_activity_id + FROM existing_duplicate_matches + LEFT JOIN current_duplicate_matches + ON current_duplicate_matches.activity_id = existing_duplicate_matches.activity_id + AND current_duplicate_matches.duplicate_activity_id = existing_duplicate_matches.duplicate_activity_id + WHERE current_duplicate_matches.activity_id IS null +), + +refresh_clock AS ( + SELECT + toUInt64(toUnixTimestamp64Nano(now64(9))) AS refresh_version, + now64(9) AS refreshed_at +) + +SELECT + activity_id, + duplicate_activity_id, + overlap_ratio, + refresh_clock.refresh_version AS refresh_version, + 0 AS is_deleted, + refresh_clock.refreshed_at AS refreshed_at +FROM current_duplicate_matches +CROSS JOIN refresh_clock + +UNION ALL + +SELECT + assumeNotNull(activity_id) AS activity_id, + assumeNotNull(duplicate_activity_id) AS duplicate_activity_id, + CAST(null, 'Nullable(Float64)') AS overlap_ratio, + refresh_clock.refresh_version AS refresh_version, + 1 AS is_deleted, + refresh_clock.refreshed_at AS refreshed_at +FROM stale_duplicate_matches +CROSS JOIN refresh_clock diff --git a/analytics/models/read_models/activity_source_records.sql b/analytics/models/read_models/activity_source_records.sql new file mode 100644 index 0000000000..7cae1513fc --- /dev/null +++ b/analytics/models/read_models/activity_source_records.sql @@ -0,0 +1,140 @@ +{{ config( + materialized='incremental', + incremental_strategy='append', + engine='ReplacingMergeTree(refresh_version)', + order_by='activity_id', + query_settings={ + 'max_threads': 1, + 'join_use_nulls': 1 + } +) }} + +WITH active_activity AS ( + SELECT * + FROM {{ source('postgres_fitness', 'activity') }} FINAL + WHERE _peerdb_is_deleted = 0 +), + +active_provider_priority AS ( + SELECT * + FROM {{ source('postgres_fitness', 'provider_priority') }} FINAL + WHERE _peerdb_is_deleted = 0 +), + +active_device_priority AS ( + SELECT * + FROM {{ source('postgres_fitness', 'device_priority') }} FINAL + WHERE _peerdb_is_deleted = 0 +), + +device_priority_match AS ( + SELECT + activity_id, + priority + FROM ( + SELECT + active_activity.id AS activity_id, + active_device_priority.priority AS priority, + row_number() OVER ( + PARTITION BY active_activity.id + ORDER BY length(active_device_priority.source_name_pattern) DESC + ) AS row_number + FROM active_activity + INNER JOIN active_device_priority + ON active_device_priority.provider_id = active_activity.provider_id + AND active_activity.source_name LIKE active_device_priority.source_name_pattern + ) + WHERE row_number = 1 +), + +current_source_records AS ( + SELECT + active_activity.id AS activity_id, + active_activity.provider_id AS provider_id, + active_activity.user_id AS user_id, + active_activity.external_id AS external_id, + active_activity.activity_type AS activity_type, + active_activity.started_at AS started_at, + active_activity.ended_at AS ended_at, + active_activity.source_name AS source_name, + active_activity.name AS name, + active_activity.notes AS notes, + active_activity.timezone AS timezone, + active_activity.raw AS raw, + active_activity._peerdb_synced_at AS source_synced_at, + coalesce(device_priority_match.priority, active_provider_priority.priority, 100) AS priority + FROM active_activity + LEFT JOIN active_provider_priority + ON active_provider_priority.provider_id = active_activity.provider_id + LEFT JOIN device_priority_match + ON device_priority_match.activity_id = active_activity.id +), + +existing_source_records AS ( + {% if is_incremental() %} + SELECT activity_id + FROM {{ this }} FINAL + WHERE is_deleted = 0 + {% else %} + SELECT CAST(null, 'Nullable(UUID)') AS activity_id + WHERE 1 = 0 + {% endif %} +), + +stale_source_records AS ( + SELECT existing_source_records.activity_id + FROM existing_source_records + LEFT JOIN current_source_records + ON current_source_records.activity_id = existing_source_records.activity_id + WHERE current_source_records.activity_id IS null +), + +refresh_clock AS ( + SELECT + toUInt64(toUnixTimestamp64Nano(now64(9))) AS refresh_version, + now64(9) AS refreshed_at +) + +SELECT + activity_id, + provider_id, + user_id, + external_id, + activity_type, + started_at, + ended_at, + source_name, + name, + notes, + timezone, + raw, + source_synced_at, + priority, + refresh_clock.refresh_version AS refresh_version, + 0 AS is_deleted, + refresh_clock.refreshed_at AS refreshed_at +FROM current_source_records +CROSS JOIN refresh_clock + +UNION ALL + +SELECT + assumeNotNull(activity_id) AS activity_id, + CAST(null, 'Nullable(String)') AS provider_id, + CAST(null, 'Nullable(UUID)') AS user_id, + CAST(null, 'Nullable(String)') AS external_id, + CAST(null, 'Nullable(String)') AS activity_type, + CAST(null, 'Nullable(DateTime64(6, ''UTC''))') AS started_at, + CAST(null, 'Nullable(DateTime64(6, ''UTC''))') AS ended_at, + CAST(null, 'Nullable(String)') AS source_name, + CAST(null, 'Nullable(String)') AS name, + CAST(null, 'Nullable(String)') AS notes, + CAST(null, 'Nullable(String)') AS timezone, + CAST(null, 'Nullable(String)') AS raw, + CAST(null, 'Nullable(DateTime64(9, ''UTC''))') AS source_synced_at, + CAST(null, 'Nullable(Int32)') AS priority, + refresh_clock.refresh_version AS refresh_version, + 1 AS is_deleted, + refresh_clock.refreshed_at AS refreshed_at +FROM stale_source_records +CROSS JOIN refresh_clock diff --git a/analytics/models/read_models/deduped_activities.sql b/analytics/models/read_models/deduped_activities.sql index 315b77d007..a37898f84e 100644 --- a/analytics/models/read_models/deduped_activities.sql +++ b/analytics/models/read_models/deduped_activities.sql @@ -9,66 +9,73 @@ } ) }} -WITH RECURSIVE target_state AS ( +WITH ranked AS ( + SELECT * + FROM {{ ref('activity_source_records') }} FINAL + WHERE is_deleted = 0 +), + +final_groups AS ( SELECT - coalesce( - max(refreshed_at), - toDateTime64('1970-01-01 00:00:00', 9, 'UTC') - ) AS last_refreshed_at, - {% if is_incremental() %}(count() = 0){% else %}1{% endif %} AS is_empty - FROM {% if is_incremental() %}{{ this }}{% else %}(SELECT CAST(null, 'Nullable(DateTime64(9, ''UTC''))') AS refreshed_at){% endif %} + activity_id, + group_id + FROM {{ ref('activity_duplicate_groups') }} FINAL + WHERE is_deleted = 0 ), -priority_changes AS ( - SELECT count() > 0 AS has_changes +best AS ( + SELECT * FROM ( - SELECT _peerdb_synced_at - FROM {{ source('postgres_fitness', 'provider_priority') }} FINAL - UNION ALL - SELECT _peerdb_synced_at - FROM {{ source('postgres_fitness', 'device_priority') }} FINAL + SELECT + final_groups.group_id AS group_id, + ranked.activity_id AS canonical_id, + ranked.provider_id AS provider_id, + ranked.user_id AS user_id, + ranked.activity_type AS activity_type, + ranked.started_at AS started_at, + ranked.ended_at AS ended_at, + ranked.source_name AS source_name, + ranked.priority AS priority, + row_number() OVER ( + PARTITION BY final_groups.group_id + ORDER BY ranked.priority ASC, toString(ranked.activity_id) ASC + ) AS row_number + FROM final_groups + INNER JOIN ranked + ON ranked.activity_id = final_groups.activity_id ) - WHERE NOT (SELECT is_empty FROM target_state) - AND _peerdb_synced_at > (SELECT last_refreshed_at FROM target_state) + WHERE row_number = 1 ), -changed_activity_windows AS ( +merged AS ( SELECT - activity.id AS activity_id, - activity.user_id AS user_id, - activity.started_at AS started_at, - coalesce(activity.ended_at, activity.started_at + INTERVAL 12 HOUR) AS ended_at - FROM {{ source('postgres_fitness', 'activity') }} AS activity FINAL - WHERE NOT (SELECT is_empty FROM target_state) - AND activity._peerdb_synced_at > (SELECT last_refreshed_at FROM target_state) -), - -changed_user_windows AS ( - SELECT - user_id, - min(started_at) - INTERVAL 12 HOUR AS started_at, - max(ended_at) + INTERVAL 12 HOUR AS ended_at - FROM changed_activity_windows - GROUP BY user_id + best.group_id AS group_id, + best.canonical_id AS id, + any(best.provider_id) AS provider_id, + any(best.user_id) AS user_id, + any(best.activity_type) AS activity_type, + min(ranked.started_at) AS started_at, + max(coalesce(ranked.ended_at, ranked.started_at + INTERVAL 12 HOUR)) AS ended_at, + any(best.source_name) AS source_name, + argMinIf(ranked.name, ranked.priority, ranked.name IS NOT NULL) AS name, + argMinIf(ranked.notes, ranked.priority, ranked.notes IS NOT NULL) AS notes, + argMinIf(ranked.timezone, ranked.priority, ranked.timezone IS NOT NULL) AS timezone, + argMinIf(ranked.raw, ranked.priority, ranked.raw IS NOT NULL) AS raw, + max(ranked.source_synced_at) AS source_synced_at, + arraySort(groupUniqArray(ranked.provider_id)) AS source_providers, + groupArrayIf( + map('providerId', ranked.provider_id, 'externalId', ranked.external_id), + ranked.external_id IS NOT NULL AND ranked.external_id != '' + ) AS source_external_ids, + groupArray(ranked.activity_id) AS member_activity_ids + FROM best + INNER JOIN final_groups + ON final_groups.group_id = best.group_id + INNER JOIN ranked + ON ranked.activity_id = final_groups.activity_id + GROUP BY best.group_id, best.canonical_id ), -active_activity AS ( - SELECT activity.* - FROM {{ source('postgres_fitness', 'activity') }} AS activity FINAL - LEFT JOIN changed_user_windows - ON changed_user_windows.user_id = activity.user_id - AND activity.started_at < changed_user_windows.ended_at - AND coalesce(activity.ended_at, activity.started_at + INTERVAL 12 HOUR) >= changed_user_windows.started_at - WHERE activity._peerdb_is_deleted = 0 - AND ( - (SELECT is_empty FROM target_state) - OR (SELECT has_changes FROM priority_changes) - OR changed_user_windows.user_id IS NOT null - ) -), - -{{ activity_dedup_graph() }}, - current_deduped_activities AS ( SELECT id AS activity_id, @@ -109,35 +116,23 @@ existing_deduped_activities AS ( deduped.source_external_ids, deduped.member_activity_ids FROM {{ this }} AS deduped FINAL - LEFT JOIN changed_user_windows - ON changed_user_windows.user_id = deduped.user_id - AND deduped.started_at < changed_user_windows.ended_at - AND coalesce(deduped.ended_at, deduped.started_at + INTERVAL 12 HOUR) >= changed_user_windows.started_at WHERE deduped.is_deleted = 0 - AND ( - (SELECT has_changes FROM priority_changes) - OR changed_user_windows.user_id IS NOT null - ) ), stale_deduped_activities AS ( - SELECT * + SELECT existing_deduped_activities.* FROM existing_deduped_activities + LEFT JOIN current_deduped_activities + ON current_deduped_activities.activity_id = existing_deduped_activities.activity_id + WHERE current_deduped_activities.activity_id IS null ) {% endif %} , -changed_user_window_count AS ( - SELECT count() AS changed_window_count - FROM changed_user_windows -), - refresh_clock AS ( SELECT toUInt64(toUnixTimestamp64Nano(now64(9))) AS refresh_version, now64(9) AS refreshed_at - FROM priority_changes - CROSS JOIN changed_user_window_count ) SELECT @@ -183,7 +178,7 @@ SELECT source_providers, source_external_ids, member_activity_ids, - refresh_clock.refresh_version - 1 AS refresh_version, + refresh_clock.refresh_version AS refresh_version, 1 AS is_deleted, refresh_clock.refreshed_at AS refreshed_at FROM stale_deduped_activities diff --git a/analytics/models/read_models/read_model_microbatch.sql.test.ts b/analytics/models/read_models/read_model_microbatch.sql.test.ts index 394b282bae..c85fd2bda4 100644 --- a/analytics/models/read_models/read_model_microbatch.sql.test.ts +++ b/analytics/models/read_models/read_model_microbatch.sql.test.ts @@ -46,6 +46,9 @@ describe("production analytics read-model build", () => { expect(safeModelMatch?.[1]?.split(" ")).toEqual([ "sensor_scalar_sample", "deduped_sensor", + "activity_source_records", + "activity_duplicate_matches", + "activity_duplicate_groups", "deduped_activities", "deduped_activity_members", "sleep_heart_rate_sample", @@ -81,30 +84,56 @@ describe("production analytics read-model build", () => { expect(normalizedSql).not.toContain("ref('sleep_heart_rate_sample') }} FINAL"); }); - it("materializes deduped activities once from the bounded activity graph", () => { + it("materializes deduped activities from domain activity dedupe models", () => { expect(existsSync(new URL("./deduped_activities.sql", import.meta.url))).toBe(true); const sql = readModel("deduped_activities"); const normalizedSql = compactWhitespace(sql); expect(sql).toContain("materialized='incremental'"); expect(sql).toContain("engine='ReplacingMergeTree(refresh_version)'"); - expect(sql).toContain("WITH RECURSIVE target_state AS"); - expect(sql).toContain("{{ activity_dedup_graph() }}"); - expect(sql).toContain("priority_changes AS"); - expect(sql).toContain("changed_user_windows AS"); - expect(sql).toContain("activity._peerdb_synced_at > (SELECT last_refreshed_at FROM target_state)"); + expect(sql).toContain("ref('activity_source_records')"); + expect(sql).toContain("ref('activity_duplicate_groups')"); + expect(sql).not.toContain("{{ activity_dedup_graph() }}"); + expect(sql).not.toContain("connected_components AS"); + expect(sql).not.toContain("visited_activity_ids"); expect(sql).toContain("current_deduped_activities AS"); expect(sql).toContain("member_activity_ids"); expect(sql).toContain("stale_deduped_activities AS"); expect(sql).toContain("{% if is_incremental() %}"); expect(sql).toContain("'join_use_nulls': 1"); expect(normalizedSql).toContain("FROM existing_deduped_activities"); - expect(sql).toContain("refresh_clock.refresh_version - 1 AS refresh_version"); - expect(normalizedSql).toContain("LEFT JOIN changed_user_windows"); expect(normalizedSql).toContain("FROM {{ this }} AS deduped FINAL"); expect(normalizedSql).toContain("WHERE deduped.is_deleted = 0"); }); + it("breaks activity deduplication into conceptual domain stages", () => { + const sourceRecordsSql = readModel("activity_source_records"); + const matchesSql = readModel("activity_duplicate_matches"); + const groupsSql = readModel("activity_duplicate_groups"); + + expect(sourceRecordsSql).toContain("materialized='incremental'"); + expect(sourceRecordsSql).toContain("engine='ReplacingMergeTree(refresh_version)'"); + expect(sourceRecordsSql).toContain("active_provider_priority AS"); + expect(sourceRecordsSql).toContain("device_priority_match AS"); + expect(sourceRecordsSql).toContain("current_source_records AS"); + + expect(matchesSql).toContain("materialized='incremental'"); + expect(matchesSql).toContain("ref('activity_source_records')"); + expect(matchesSql).toContain("current_duplicate_matches AS"); + expect(matchesSql).toContain("overlap_ratio"); + + expect(groupsSql).toContain("materialized='incremental'"); + expect(groupsSql).toContain("ref('activity_source_records')"); + expect(groupsSql).toContain("ref('activity_duplicate_matches')"); + expect(groupsSql).toContain("duplicate_links AS"); + expect(groupsSql).toContain("duplicate_walk AS"); + expect(groupsSql).toContain("current_duplicate_groups AS"); + expect(groupsSql).toContain("GROUP BY activity_id"); + expect(groupsSql).not.toContain("WITH RECURSIVE"); + expect(groupsSql).not.toContain("visited_activity_ids"); + expect(groupsSql).not.toContain("activity_reachability_"); + }); + it("materializes deduped activity member aliases from deduped activities", () => { expect(existsSync(new URL("./deduped_activity_members.sql", import.meta.url))).toBe(true); const sql = readModel("deduped_activity_members"); @@ -164,33 +193,20 @@ describe("production analytics read-model build", () => { expect(sql).toContain("trim(BOTH '()' FROM location_rows.point_text)"); }); - it("bounds activity graph construction to the active microbatch window", () => { - const sql = readProjectFile("analytics/macros/bounded_activity_graph.sql"); - - expect(sql).toContain("__dbt_internal_microbatch_event_time_start"); - expect(sql).toContain("__dbt_internal_microbatch_event_time_end"); - expect(sql).toContain("started_at < toDateTime64('{{ batch_end }}', 6, 'UTC')"); - expect(sql).toContain( - "coalesce(ended_at, started_at + INTERVAL 12 HOUR) >= toDateTime64('{{ batch_start }}', 6, 'UTC')", - ); - }); - - it("uses the same null-ended activity window for overlap matching", () => { - const sql = readProjectFile("analytics/macros/bounded_activity_graph.sql"); - - expect(sql).toContain("coalesce(left_activity.ended_at, left_activity.started_at + INTERVAL 12 HOUR)"); - expect(sql).toContain("coalesce(right_activity.ended_at, right_activity.started_at + INTERVAL 12 HOUR)"); - expect(sql).not.toContain("INTERVAL 1 HOUR"); - }); + it("uses the same null-ended activity window for duplicate matches and merged activities", () => { + const matchesSql = readModel("activity_duplicate_matches"); + const dedupedActivitiesSql = readModel("deduped_activities"); - it("uses merged group time bounds for current activity sensor membership", () => { - const sql = readProjectFile("analytics/macros/bounded_activity_graph.sql"); + expect(matchesSql).toContain("coalesce(ended_at, started_at + INTERVAL 12 HOUR) AS ended_at"); + expect(matchesSql).toContain("greatest(left_activity.started_at, right_activity.started_at)"); + expect(matchesSql).toContain("least(left_activity.ended_at, right_activity.ended_at)"); + expect(matchesSql).not.toContain("INTERVAL 1 HOUR"); - expect(sql).toContain("min(ranked.started_at) AS started_at"); - expect(sql).toContain("max(coalesce(ranked.ended_at, ranked.started_at + INTERVAL 12 HOUR)) AS ended_at"); - expect(sql).toContain("max(ranked._peerdb_synced_at) AS source_synced_at"); - expect(sql).not.toContain("any(best.started_at) AS started_at"); - expect(sql).not.toContain("any(best.ended_at) AS ended_at"); + expect(dedupedActivitiesSql).toContain("min(ranked.started_at) AS started_at"); + expect(dedupedActivitiesSql).toContain("max(coalesce(ranked.ended_at, ranked.started_at + INTERVAL 12 HOUR)) AS ended_at"); + expect(dedupedActivitiesSql).toContain("max(ranked.source_synced_at) AS source_synced_at"); + expect(dedupedActivitiesSql).not.toContain("any(best.started_at) AS started_at"); + expect(dedupedActivitiesSql).not.toContain("any(best.ended_at) AS ended_at"); }); it("carries upstream source freshness through lookback microbatch intermediaries", () => { diff --git a/docs/production-incident-baseline.md b/docs/production-incident-baseline.md index b5512950c6..34ef355cce 100644 --- a/docs/production-incident-baseline.md +++ b/docs/production-incident-baseline.md @@ -9405,3 +9405,45 @@ new incremental tables are populated. - Remaining risk: Local full-stack e2e validation was blocked by Docker network address-pool exhaustion; local single-model dbt first-build and incremental runs reproduced and validated the failing model path. + +### 2026-06-01 activities empty-state update + +- Symptoms: `https://dofek.asherlc.com/activities` showed "No activities in the + last 4 weeks" even though recent activities should exist. +- User impact: The activity calendar list and overview could temporarily hide + recent activity data after sync/import and analytics catch-up. +- Evidence: Production Postgres `fitness.v_activity` had 73 recent completed + activities for user `f923fed7-d934-4cd9-8cb9-8e83020d0e69` since + `2026-05-04`, latest `2026-06-01 03:39:00.46+00`. Mirrored ClickHouse + `postgres_fitness.activity FINAL` had 214 recent raw completed rows and + `analytics.activity_summary` had 211 recent rows, but + `analytics.deduped_activities FINAL` had zero rows total. The + `analytics-worker` first fatal log line was ClickHouse + `MEMORY_LIMIT_EXCEEDED` while executing `RecursiveCTESource` in + `deduped_activities`. +- Root cause: The first build of dbt-owned `analytics.deduped_activities` ran + the recursive activity-overlap graph over all mirrored historical activities + because the target table was empty. That unbounded recursive CTE exceeded the + ClickHouse memory limit, leaving the table empty. The Activities page reads + `analytics.deduped_activities`, so it returned no recent activity cards even + though canonical and summary data existed. +- Fix / mitigation: Production was manually populated through ClickHouse using + materialized intermediate source-record, duplicate-match, duplicate-group, + and canonical-activity stages. After the manual insert, + `analytics.deduped_activities FINAL` had 828 active rows; the + page-equivalent recent query returned 73 rows with matching + `analytics.activity_summary` rows and latest start + `2026-06-01 03:39:00`. The repo fix replaces the path-enumerating recursive + graph macro with dbt-owned domain read models: + `activity_source_records`, `activity_duplicate_matches`, + `activity_duplicate_groups`, and `deduped_activities`. A read-only + production performance check of the monolithic domain CTE returned the right + 828 groups but took 9.4s; the dbt-style materialized component check returned + the same 828 groups with the duplicate-group stage taking 27ms and about 10MB + peak memory. +- Remaining risk: Production is manually mitigated, but the analytics worker + still needs this repo fix deployed before scheduled dbt builds stop retrying + the old recursive model. Local dbt-templated SQL lint compiled the project but + could not complete because local ClickHouse at `127.0.0.1:8123` was not + running; starting the local compose dependency was blocked by Docker address + pool exhaustion in this workspace. diff --git a/entrypoint.sh b/entrypoint.sh index 1bd74c116e..a593e807d5 100755 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -18,7 +18,7 @@ fi # Node 22+ natively handles TypeScript — transform-types also rewrites .ts imports NODE="node --experimental-transform-types --enable-source-maps --disable-warning=ExperimentalWarning --import ./src/opentelemetry-hook.mjs --import ./src/instrumentation.ts" -DBT_SAFE_MODELS="sensor_scalar_sample deduped_sensor deduped_activities deduped_activity_members sleep_heart_rate_sample resting_heart_rate_sleep_window activity_sensor_sample activity_location_sample activity_sensor_summary_rows activity_location_summary_rows activity_summary_rows activity_vo2max_estimate" +DBT_SAFE_MODELS="sensor_scalar_sample deduped_sensor activity_source_records activity_duplicate_matches activity_duplicate_groups deduped_activities deduped_activity_members sleep_heart_rate_sample resting_heart_rate_sleep_window activity_sensor_sample activity_location_sample activity_sensor_summary_rows activity_location_summary_rows activity_summary_rows activity_vo2max_estimate" case "${1:-sync}" in web) From 690a896c777d6ba4b1433748d39091375e473233 Mon Sep 17 00:00:00 2001 From: Asher Cohen Date: Mon, 1 Jun 2026 19:03:51 -0700 Subject: [PATCH 2/3] Fix deduped activity nullable sort key --- .../models/read_models/deduped_activities.sql | 4 ++-- .../read_model_microbatch.sql.test.ts | 1 + docs/production-incident-baseline.md | 22 +++++++++++++++++++ 3 files changed, 25 insertions(+), 2 deletions(-) diff --git a/analytics/models/read_models/deduped_activities.sql b/analytics/models/read_models/deduped_activities.sql index a37898f84e..54fcce8296 100644 --- a/analytics/models/read_models/deduped_activities.sql +++ b/analytics/models/read_models/deduped_activities.sql @@ -138,7 +138,7 @@ refresh_clock AS ( SELECT activity_id, provider_id, - user_id, + assumeNotNull(user_id) AS user_id, activity_id AS primary_activity_id, activity_type, started_at, @@ -164,7 +164,7 @@ UNION ALL SELECT activity_id, provider_id, - user_id, + assumeNotNull(user_id) AS user_id, activity_id AS primary_activity_id, activity_type, started_at, diff --git a/analytics/models/read_models/read_model_microbatch.sql.test.ts b/analytics/models/read_models/read_model_microbatch.sql.test.ts index c85fd2bda4..ebbc1fd6d3 100644 --- a/analytics/models/read_models/read_model_microbatch.sql.test.ts +++ b/analytics/models/read_models/read_model_microbatch.sql.test.ts @@ -98,6 +98,7 @@ describe("production analytics read-model build", () => { expect(sql).not.toContain("visited_activity_ids"); expect(sql).toContain("current_deduped_activities AS"); expect(sql).toContain("member_activity_ids"); + expect(sql).toContain("assumeNotNull(user_id) AS user_id"); expect(sql).toContain("stale_deduped_activities AS"); expect(sql).toContain("{% if is_incremental() %}"); expect(sql).toContain("'join_use_nulls': 1"); diff --git a/docs/production-incident-baseline.md b/docs/production-incident-baseline.md index 34ef355cce..0fc25594cf 100644 --- a/docs/production-incident-baseline.md +++ b/docs/production-incident-baseline.md @@ -9447,3 +9447,25 @@ new incremental tables are populated. could not complete because local ClickHouse at `127.0.0.1:8123` was not running; starting the local compose dependency was blocked by Docker address pool exhaustion in this workspace. + +### 2026-06-02 PR CI nullable sort-key update + +- Symptoms: PR `Test / E2E Tests (Web)` failed during the tracked e2e + analytics build, causing `Test / Test Gate` and `CI Gate` to fail. +- User impact: The activity dedupe PR could not pass required CI despite the + production deploy succeeding against an existing `deduped_activities` table. +- Evidence: The failing e2e step was + `docker compose -f docker-compose.e2e.yml up -d --no-build analytics`; the + first fatal dbt line was `Database Error in model deduped_activities`, with + ClickHouse error `Sorting key contains nullable columns, but merge tree + setting allow_nullable_key is disabled`. +- Root cause: The clean e2e first build inferred `deduped_activities.user_id` + as nullable because upstream source-record models include tombstone branches + with nullable non-key fields, while `deduped_activities` sorts by + `(user_id, activity_id)`. +- Fix / mitigation: `deduped_activities` now emits + `assumeNotNull(user_id) AS user_id` in both current and stale output branches + so clean first builds create a non-null sort-key column. +- Remaining risk: Full local e2e validation remains blocked by Docker network + address-pool exhaustion in this workspace; CI rerun is the end-to-end + validation for the compose e2e path. From ac5ee7029bf26c72c63130e7bce078e5e309dd72 Mon Sep 17 00:00:00 2001 From: Asher Cohen Date: Mon, 1 Jun 2026 19:13:48 -0700 Subject: [PATCH 3/3] Address activity dedupe review comments --- .../read_models/activity_source_records.sql | 5 ++++- .../models/read_models/deduped_activities.sql | 1 + .../read_model_microbatch.sql.test.ts | 19 ++++++------------- 3 files changed, 11 insertions(+), 14 deletions(-) diff --git a/analytics/models/read_models/activity_source_records.sql b/analytics/models/read_models/activity_source_records.sql index 7cae1513fc..32e7462709 100644 --- a/analytics/models/read_models/activity_source_records.sql +++ b/analytics/models/read_models/activity_source_records.sql @@ -37,7 +37,10 @@ device_priority_match AS ( active_device_priority.priority AS priority, row_number() OVER ( PARTITION BY active_activity.id - ORDER BY length(active_device_priority.source_name_pattern) DESC + ORDER BY + length(active_device_priority.source_name_pattern) DESC, + active_device_priority.priority ASC, + active_device_priority.source_name_pattern ASC ) AS row_number FROM active_activity INNER JOIN active_device_priority diff --git a/analytics/models/read_models/deduped_activities.sql b/analytics/models/read_models/deduped_activities.sql index 54fcce8296..8d09d8f070 100644 --- a/analytics/models/read_models/deduped_activities.sql +++ b/analytics/models/read_models/deduped_activities.sql @@ -124,6 +124,7 @@ stale_deduped_activities AS ( FROM existing_deduped_activities LEFT JOIN current_deduped_activities ON current_deduped_activities.activity_id = existing_deduped_activities.activity_id + AND current_deduped_activities.user_id = existing_deduped_activities.user_id WHERE current_deduped_activities.activity_id IS null ) {% endif %} diff --git a/analytics/models/read_models/read_model_microbatch.sql.test.ts b/analytics/models/read_models/read_model_microbatch.sql.test.ts index ebbc1fd6d3..9236f582e4 100644 --- a/analytics/models/read_models/read_model_microbatch.sql.test.ts +++ b/analytics/models/read_models/read_model_microbatch.sql.test.ts @@ -93,9 +93,6 @@ describe("production analytics read-model build", () => { expect(sql).toContain("engine='ReplacingMergeTree(refresh_version)'"); expect(sql).toContain("ref('activity_source_records')"); expect(sql).toContain("ref('activity_duplicate_groups')"); - expect(sql).not.toContain("{{ activity_dedup_graph() }}"); - expect(sql).not.toContain("connected_components AS"); - expect(sql).not.toContain("visited_activity_ids"); expect(sql).toContain("current_deduped_activities AS"); expect(sql).toContain("member_activity_ids"); expect(sql).toContain("assumeNotNull(user_id) AS user_id"); @@ -105,6 +102,9 @@ describe("production analytics read-model build", () => { expect(normalizedSql).toContain("FROM existing_deduped_activities"); expect(normalizedSql).toContain("FROM {{ this }} AS deduped FINAL"); expect(normalizedSql).toContain("WHERE deduped.is_deleted = 0"); + expect(normalizedSql).toContain( + "ON current_deduped_activities.activity_id = existing_deduped_activities.activity_id AND current_deduped_activities.user_id = existing_deduped_activities.user_id", + ); }); it("breaks activity deduplication into conceptual domain stages", () => { @@ -117,6 +117,9 @@ describe("production analytics read-model build", () => { expect(sourceRecordsSql).toContain("active_provider_priority AS"); expect(sourceRecordsSql).toContain("device_priority_match AS"); expect(sourceRecordsSql).toContain("current_source_records AS"); + expect(sourceRecordsSql).toContain("length(active_device_priority.source_name_pattern) DESC"); + expect(sourceRecordsSql).toContain("active_device_priority.priority ASC"); + expect(sourceRecordsSql).toContain("active_device_priority.source_name_pattern ASC"); expect(matchesSql).toContain("materialized='incremental'"); expect(matchesSql).toContain("ref('activity_source_records')"); @@ -130,9 +133,6 @@ describe("production analytics read-model build", () => { expect(groupsSql).toContain("duplicate_walk AS"); expect(groupsSql).toContain("current_duplicate_groups AS"); expect(groupsSql).toContain("GROUP BY activity_id"); - expect(groupsSql).not.toContain("WITH RECURSIVE"); - expect(groupsSql).not.toContain("visited_activity_ids"); - expect(groupsSql).not.toContain("activity_reachability_"); }); it("materializes deduped activity member aliases from deduped activities", () => { @@ -151,7 +151,6 @@ describe("production analytics read-model build", () => { ); expect(sql).toContain("arrayJoin(deduped_activities.member_activity_ids) AS member_activity_id"); expect(sql).toContain("stale_activity_members AS"); - expect(sql).not.toContain("bounded_activity_graph()"); expect(sql).toContain("'join_use_nulls': 1"); expect(normalizedSql).toContain("LEFT JOIN current_activity_members"); expect(normalizedSql).toContain( @@ -170,7 +169,6 @@ describe("production analytics read-model build", () => { expect(sql).toContain("lookback=3"); expect(sql).toContain("ref('deduped_sensor')"); expect(sql).toContain("ref('deduped_activities')"); - expect(sql).not.toContain("bounded_activity_graph()"); expect(sql).not.toContain("source('analytics', 'v_activity')"); expect(sql).toContain("activity_id"); }); @@ -184,7 +182,6 @@ describe("production analytics read-model build", () => { expect(sql).toContain("lookback=3"); expect(sql).toContain("source('postgres_fitness', 'metric_stream')"); expect(sql).toContain("ref('deduped_activity_members')"); - expect(sql).not.toContain("bounded_activity_graph()"); expect(sql).not.toContain("source('analytics', 'v_activity_members')"); expect(sql).toContain("channel = 'location'"); expect(sql).toContain("argMax(point, _peerdb_version) AS point"); @@ -201,13 +198,10 @@ describe("production analytics read-model build", () => { expect(matchesSql).toContain("coalesce(ended_at, started_at + INTERVAL 12 HOUR) AS ended_at"); expect(matchesSql).toContain("greatest(left_activity.started_at, right_activity.started_at)"); expect(matchesSql).toContain("least(left_activity.ended_at, right_activity.ended_at)"); - expect(matchesSql).not.toContain("INTERVAL 1 HOUR"); expect(dedupedActivitiesSql).toContain("min(ranked.started_at) AS started_at"); expect(dedupedActivitiesSql).toContain("max(coalesce(ranked.ended_at, ranked.started_at + INTERVAL 12 HOUR)) AS ended_at"); expect(dedupedActivitiesSql).toContain("max(ranked.source_synced_at) AS source_synced_at"); - expect(dedupedActivitiesSql).not.toContain("any(best.started_at) AS started_at"); - expect(dedupedActivitiesSql).not.toContain("any(best.ended_at) AS ended_at"); }); it("carries upstream source freshness through lookback microbatch intermediaries", () => { @@ -267,7 +261,6 @@ describe("production analytics read-model build", () => { expect(sql).toContain("ref('deduped_activities')"); expect(sql).toContain("ref('deduped_activity_members')"); - expect(sql).not.toContain("bounded_activity_graph()"); expect(sql).toContain("ref('activity_sensor_summary_rows')"); expect(sql).toContain("ref('activity_location_summary_rows')"); expect(sql).toContain("(user_id, activity_id) IN");