|
117 | 117 | (.execute stmt |
118 | 118 | "CREATE TABLE IF NOT EXISTS version (version INTEGER)") |
119 | 119 | (.execute stmt |
120 | | - "INSERT INTO version VALUES (4)") |
| 120 | + "INSERT INTO version VALUES (5)") |
121 | 121 | (.execute stmt |
122 | 122 | "CREATE TYPE IF NOT EXISTS agent_type_t AS ENUM ('feed', 'bot', 'browser')") |
123 | 123 | (.execute stmt |
|
139 | 139 | (.execute stmt |
140 | 140 | "CREATE TYPE IF NOT EXISTS dim_t AS ENUM ('feed', 'bot', 'browser', 'path', 'query', 'ref_domain')") |
141 | 141 | (.execute stmt |
142 | | - ;; Per-day aggregates of stats for all days <= rollup_state.last_date. |
| 142 | + ;; Per-day aggregates of stats for all completed days. |
143 | 143 | ;; For dims feed/bot/browser, value is agent (possibly NULL). |
144 | 144 | ;; For dims path/query/ref_domain, value is that column. |
145 | 145 | ;; Only the top 20 (rollup-depth) values per (date, dim) are kept; |
|
150 | 150 | dim dim_t, |
151 | 151 | value VARCHAR, |
152 | 152 | cnt BIGINT |
153 | | - )") |
154 | | - (.execute stmt |
155 | | - "CREATE TABLE IF NOT EXISTS rollup_state (last_date DATE)") |
156 | | - (.execute stmt |
157 | | - "INSERT INTO rollup_state VALUES (DATE '1970-01-01')"))) |
| 153 | + )"))) |
158 | 154 |
|
159 | 155 | (defn db-version ^long [^DuckDBConnection conn] |
160 | 156 | (or |
|
232 | 228 | (log "Migrating" db-path "to version 4") |
233 | 229 | (with-open [conn (connect db-path) |
234 | 230 | stmt (.createStatement conn)] |
235 | | - (.execute stmt "BEGIN TRANSACTION") |
| 231 | + (.execute stmt "DROP TABLE IF EXISTS daily_counts") |
236 | 232 | (.execute stmt |
237 | | - (str |
238 | | - "CREATE TABLE rollup_daily AS |
239 | | - WITH ranked AS ( |
240 | | - SELECT date, dim, value, cnt, |
241 | | - ROW_NUMBER() OVER (PARTITION BY date, dim ORDER BY cnt DESC, value) AS rn |
242 | | - FROM daily_counts |
243 | | - WHERE value IS NOT NULL |
244 | | - ) |
245 | | - SELECT date, dim, value, cnt FROM ranked WHERE rn <= " rollup-depth " |
246 | | - UNION ALL |
247 | | - SELECT date, dim, '" rollup-sentinel "', SUM(cnt)::BIGINT |
248 | | - FROM ranked WHERE rn > " rollup-depth " GROUP BY date, dim |
249 | | - UNION ALL |
250 | | - SELECT date, dim, value, cnt FROM daily_counts WHERE value IS NULL")) |
251 | | - (.execute stmt "DROP TABLE daily_counts") |
252 | | - (.execute stmt "UPDATE version SET version = 4") |
253 | | - (.execute stmt "COMMIT")) |
| 233 | + "CREATE TABLE IF NOT EXISTS rollup_daily ( |
| 234 | + date DATE, |
| 235 | + dim dim_t, |
| 236 | + value VARCHAR, |
| 237 | + cnt BIGINT |
| 238 | + )") |
| 239 | + (.execute stmt "UPDATE version SET version = 4")) |
254 | 240 | (log "Migration to version 4 complete")) |
255 | 241 |
|
| 242 | +(defn migrate-4->5! [^String db-path] |
| 243 | + (log "Migrating" db-path "to version 5") |
| 244 | + (with-open [conn (connect db-path) |
| 245 | + stmt (.createStatement conn)] |
| 246 | + (.execute stmt "DROP TABLE IF EXISTS rollup_state") |
| 247 | + (.execute stmt "UPDATE version SET version = 5")) |
| 248 | + (log "Migration to version 5 complete")) |
| 249 | + |
256 | 250 | (defn- rollup-day! |
257 | | - "Aggregates one day of stats into rollup_daily and advances rollup_state to it" |
| 251 | + "Aggregates one day of stats into rollup_daily" |
258 | 252 | [^DuckDBConnection conn ^LocalDate date] |
259 | 253 | (log-verbose "Rolling up" (str date)) |
260 | 254 | (.setAutoCommit conn false) |
|
306 | 300 | FROM ranked WHERE rn > " rollup-depth " GROUP BY date"))] |
307 | 301 | (.setObject stmt 1 date) |
308 | 302 | (.execute stmt))) |
309 | | - (with-open [stmt (.prepareStatement conn "UPDATE rollup_state SET last_date = ?")] |
310 | | - (.setObject stmt 1 date) |
311 | | - (.execute stmt)) |
312 | 303 | (.commit conn) |
313 | 304 | (catch Throwable t |
314 | 305 | (.rollback conn) |
|
317 | 308 | (.setAutoCommit conn true)))) |
318 | 309 |
|
319 | 310 | (def ^:private *rollup-dates |
320 | | - "db-path -> in-memory mirror of rollup_state.last_date, so worker ticks |
321 | | - with a current watermark don't open the database at all. Read from the |
322 | | - database once after startup, then maintained in memory" |
| 311 | + "db-path -> last rolluped up day" |
323 | 312 | (atom {})) |
324 | 313 |
|
325 | 314 | (defn rollup! |
326 | | - "Rolls up every completed UTC day after rollup_state.last_date into rollup_daily, |
327 | | - day by day, then advances the watermark to yesterday" |
| 315 | + "Rolls up every completed UTC day after the watermark into rollup_daily, |
| 316 | + day by day. The persistent watermark is MAX(date) in rollup_daily" |
328 | 317 | [db-path] |
329 | 318 | (let [yesterday (.minusDays (LocalDate/now UTC) 1) |
330 | 319 | last-date (get @*rollup-dates db-path)] |
331 | 320 | (when (or (nil? last-date) (LocalDate/.isBefore last-date yesterday)) |
332 | 321 | (with-open [conn (acquire-conn db-path)] |
333 | 322 | (let [last-date (or last-date |
334 | 323 | (with-open [stmt (.createStatement conn) |
335 | | - rs (.executeQuery stmt "SELECT last_date FROM rollup_state")] |
| 324 | + rs (.executeQuery stmt "SELECT COALESCE(MAX(date), DATE '1970-01-01') FROM rollup_daily")] |
336 | 325 | (when (.next rs) |
337 | 326 | ^LocalDate (.getObject rs 1))))] |
338 | | - (when last-date |
339 | | - (when (LocalDate/.isBefore last-date yesterday) |
340 | | - (let [dates (with-open [stmt (doto (.prepareStatement conn |
341 | | - "SELECT DISTINCT date FROM stats WHERE date > ? AND date <= ? ORDER BY date") |
342 | | - (.setObject 1 last-date) |
343 | | - (.setObject 2 yesterday)) |
344 | | - rs (.executeQuery stmt)] |
345 | | - (loop [acc []] |
346 | | - (if (.next rs) |
347 | | - (recur (conj acc (.getObject rs 1))) |
348 | | - acc)))] |
349 | | - (doseq [date dates] |
350 | | - (rollup-day! conn date)) |
351 | | - ;; advance watermark over trailing empty days too |
352 | | - (with-open [stmt (.prepareStatement conn "UPDATE rollup_state SET last_date = ?")] |
353 | | - (.setObject stmt 1 yesterday) |
354 | | - (.execute stmt)) |
355 | | - (log-verbose "Rolled up" (count dates) "day(s), watermark at" (str yesterday)))) |
356 | | - (swap! *rollup-dates assoc db-path yesterday))))))) |
| 327 | + (when (LocalDate/.isBefore last-date yesterday) |
| 328 | + (let [dates (with-open [stmt (doto (.prepareStatement conn |
| 329 | + "SELECT DISTINCT date FROM stats WHERE date > ? AND date <= ? ORDER BY date") |
| 330 | + (.setObject 1 last-date) |
| 331 | + (.setObject 2 yesterday)) |
| 332 | + rs (.executeQuery stmt)] |
| 333 | + (loop [acc []] |
| 334 | + (if (.next rs) |
| 335 | + (recur (conj acc (.getObject rs 1))) |
| 336 | + acc)))] |
| 337 | + (doseq [date dates] |
| 338 | + (rollup-day! conn date)) |
| 339 | + (log-verbose "Rolled up" (count dates) "day(s), watermark at" (str yesterday)))) |
| 340 | + (swap! *rollup-dates assoc db-path yesterday)))))) |
357 | 341 |
|
358 | 342 | (defn rebuild-rollup! |
359 | 343 | "Recompute rollup_daily from scratch" |
360 | 344 | [db-path] |
361 | 345 | (with-conn [conn db-path] |
362 | 346 | (with-open [stmt (.createStatement conn)] |
363 | | - (.execute stmt "DELETE FROM rollup_daily") |
364 | | - (.execute stmt "UPDATE rollup_state SET last_date = DATE '1970-01-01'"))) |
| 347 | + (.execute stmt "DELETE FROM rollup_daily"))) |
365 | 348 | (swap! *rollup-dates dissoc db-path) |
366 | 349 | (rollup! db-path)) |
367 | 350 |
|
|
376 | 359 | (when (<= v 2) |
377 | 360 | (migrate-2->3! db-path)) |
378 | 361 | (when (<= v 3) |
379 | | - (migrate-3->4! db-path))))) |
| 362 | + (migrate-3->4! db-path)) |
| 363 | + (when (<= v 4) |
| 364 | + (migrate-4->5! db-path))))) |
380 | 365 |
|
381 | 366 | (def ^:private *worker-pool |
382 | 367 | (atom nil)) |
|
0 commit comments