|
| 1 | +-- rank_packages() applied its ranking to all of `packages` in one UPDATE (9.66M |
| 2 | +-- rows), which under REPLICA IDENTITY FULL blew up the Sequin replication slot. |
| 3 | +-- rank_packages_chunked() stages the same ranking into an UNLOGGED table, then |
| 4 | +-- applies it in committed keyset chunks so the slot advances continuously. |
| 5 | +-- rank_packages() is left in place until the new worker deploys. |
| 6 | + |
| 7 | +CREATE UNLOGGED TABLE IF NOT EXISTS staging.package_rank ( |
| 8 | + package_id bigint PRIMARY KEY, |
| 9 | + impact numeric(10, 4), |
| 10 | + is_critical bool NOT NULL, |
| 11 | + rank_in_ecosystem int NOT NULL |
| 12 | +); |
| 13 | + |
| 14 | +CREATE OR REPLACE PROCEDURE rank_packages_chunked( |
| 15 | + coverage_cutoff numeric DEFAULT 0.90, |
| 16 | + ecosystems text[] DEFAULT NULL, |
| 17 | + chunk_size int DEFAULT 25000, |
| 18 | + INOUT applied_rows int DEFAULT 0 |
| 19 | +) |
| 20 | +LANGUAGE plpgsql AS $$ |
| 21 | +DECLARE |
| 22 | + effective_ecosystems text[]; |
| 23 | + staged_count int; |
| 24 | + batch_rows int; |
| 25 | + cursor_id bigint := 0; |
| 26 | +BEGIN |
| 27 | + SET LOCAL max_parallel_workers_per_gather = 4; |
| 28 | + |
| 29 | + IF chunk_size IS NULL OR chunk_size <= 0 THEN |
| 30 | + RAISE EXCEPTION 'rank_packages_chunked: chunk_size must be a positive integer, got %', chunk_size; |
| 31 | + END IF; |
| 32 | + |
| 33 | + -- Session-level: survives the internal COMMITs below |
| 34 | + IF NOT pg_try_advisory_lock(hashtextextended('rank_packages_chunked', 0)) THEN |
| 35 | + RAISE EXCEPTION 'rank_packages_chunked: another execution is already in progress'; |
| 36 | + END IF; |
| 37 | + |
| 38 | + applied_rows := 0; |
| 39 | + |
| 40 | + IF ecosystems IS NULL THEN |
| 41 | + SELECT ARRAY_AGG(DISTINCT ecosystem) |
| 42 | + INTO effective_ecosystems |
| 43 | + FROM packages; |
| 44 | + ELSE |
| 45 | + effective_ecosystems := ecosystems; |
| 46 | + END IF; |
| 47 | + |
| 48 | + TRUNCATE staging.package_rank; |
| 49 | + |
| 50 | + -- Scoring CTE chain, unchanged from rank_packages() (V1783123201). |
| 51 | + INSERT INTO staging.package_rank (package_id, impact, is_critical, rank_in_ecosystem) |
| 52 | + WITH base AS ( |
| 53 | + SELECT |
| 54 | + id, |
| 55 | + ecosystem, |
| 56 | + COALESCE(downloads_last_30d, 0) AS downloads, |
| 57 | + COALESCE(dependent_count, 0) AS direct_dependents, |
| 58 | + COALESCE(transitive_dependent_count, 0) AS transitive_dependents, |
| 59 | + COALESCE(sonatype_popularity_score, 0) AS sonatype_popularity, |
| 60 | + SUM(COALESCE(downloads_last_30d, 0)) OVER (PARTITION BY ecosystem) AS ecosystem_total_downloads, |
| 61 | + SUM(COALESCE(dependent_count, 0)) OVER (PARTITION BY ecosystem) AS ecosystem_total_direct_dependents, |
| 62 | + SUM(COALESCE(transitive_dependent_count, 0)) OVER (PARTITION BY ecosystem) AS ecosystem_total_transitive_dependents, |
| 63 | + SUM(COALESCE(sonatype_popularity_score, 0)) OVER (PARTITION BY ecosystem) AS ecosystem_total_sonatype |
| 64 | + FROM packages |
| 65 | + WHERE ecosystem = ANY(effective_ecosystems) |
| 66 | + ), |
| 67 | + walked AS ( |
| 68 | + SELECT |
| 69 | + id, |
| 70 | + ecosystem, |
| 71 | + SUM(signal_value) OVER coverage_window / ecosystem_signal_total::numeric AS cumulative_share_inclusive, |
| 72 | + (SUM(signal_value) OVER coverage_window - signal_value) / ecosystem_signal_total::numeric AS cumulative_share_exclusive |
| 73 | + FROM base |
| 74 | + CROSS JOIN LATERAL (VALUES |
| 75 | + ('downloads', downloads, ecosystem_total_downloads), |
| 76 | + ('direct_dependents', direct_dependents, ecosystem_total_direct_dependents), |
| 77 | + ('transitive_dependents', transitive_dependents, ecosystem_total_transitive_dependents), |
| 78 | + ('sonatype_popularity', sonatype_popularity, ecosystem_total_sonatype) |
| 79 | + ) AS signal(signal_name, signal_value, ecosystem_signal_total) |
| 80 | + WHERE ecosystem_signal_total > 0 |
| 81 | + WINDOW coverage_window AS ( |
| 82 | + PARTITION BY ecosystem, signal_name |
| 83 | + ORDER BY signal_value DESC, id |
| 84 | + ROWS UNBOUNDED PRECEDING |
| 85 | + ) |
| 86 | + ), |
| 87 | + combined AS ( |
| 88 | + SELECT |
| 89 | + id, |
| 90 | + ecosystem, |
| 91 | + AVG(1.0 - cumulative_share_inclusive)::numeric(10, 4) AS new_impact, |
| 92 | + BOOL_OR(cumulative_share_exclusive < coverage_cutoff) AS new_is_critical |
| 93 | + FROM walked |
| 94 | + GROUP BY id, ecosystem |
| 95 | + ), |
| 96 | + final AS ( |
| 97 | + SELECT |
| 98 | + combined.id, |
| 99 | + combined.new_impact, |
| 100 | + combined.new_is_critical OR (spotlight.package_id IS NOT NULL) AS new_is_critical, |
| 101 | + ROW_NUMBER() OVER ( |
| 102 | + PARTITION BY combined.ecosystem |
| 103 | + ORDER BY combined.new_impact DESC NULLS LAST, combined.id |
| 104 | + ) AS new_rank_in_ecosystem |
| 105 | + FROM combined |
| 106 | + LEFT JOIN package_criticality_spotlight spotlight ON spotlight.package_id = combined.id |
| 107 | + ) |
| 108 | + SELECT id, new_impact, new_is_critical, new_rank_in_ecosystem::int |
| 109 | + FROM final; |
| 110 | + |
| 111 | + GET DIAGNOSTICS staged_count = ROW_COUNT; |
| 112 | + |
| 113 | + IF staged_count = 0 THEN |
| 114 | + RAISE EXCEPTION 'rank_packages_chunked: computed 0 rows, refusing to apply an empty ranking'; |
| 115 | + END IF; |
| 116 | + |
| 117 | + ANALYZE staging.package_rank; |
| 118 | + |
| 119 | + COMMIT; |
| 120 | + |
| 121 | + LOOP |
| 122 | + WITH batch AS ( |
| 123 | + SELECT package_id, impact, is_critical, rank_in_ecosystem |
| 124 | + FROM staging.package_rank |
| 125 | + WHERE package_id > cursor_id |
| 126 | + ORDER BY package_id |
| 127 | + LIMIT chunk_size |
| 128 | + ), |
| 129 | + updated AS ( |
| 130 | + UPDATE packages p |
| 131 | + SET impact = b.impact, |
| 132 | + is_critical = b.is_critical, |
| 133 | + rank_in_ecosystem = b.rank_in_ecosystem, |
| 134 | + last_rank_pass_at = NOW(), |
| 135 | + last_synced_at = NOW() |
| 136 | + FROM batch b |
| 137 | + WHERE p.id = b.package_id |
| 138 | + RETURNING p.id |
| 139 | + ) |
| 140 | + SELECT COUNT(*), COALESCE(MAX(b.package_id), cursor_id) |
| 141 | + INTO batch_rows, cursor_id |
| 142 | + FROM batch b; |
| 143 | + |
| 144 | + applied_rows := applied_rows + batch_rows; |
| 145 | + |
| 146 | + COMMIT; |
| 147 | + |
| 148 | + EXIT WHEN batch_rows < chunk_size; |
| 149 | + END LOOP; |
| 150 | + |
| 151 | + PERFORM pg_advisory_unlock(hashtextextended('rank_packages_chunked', 0)); |
| 152 | +END; |
| 153 | +$$; |
0 commit comments