· 8 years ago · Dec 26, 2017, 02:12 PM
1CREATE TABLE IF NOT EXISTS rollup_generations (
2 function_name text primary key,
3 current_generation int default 0,
4 processed_generation int default -1
5);
6
7CREATE OR REPLACE FUNCTION current_generation(rollup text)
8RETURNS int LANGUAGE plpgsql STABLE
9AS $function$
10DECLARE
11 generation int;
12BEGIN
13 /* get the current generation of the rollup */
14 SELECT current_generation INTO generation
15 FROM rollup_generations WHERE function_name = rollup;
16
17 /* take a share lock on the current generation to prevent concurrent rollup */
18 PERFORM pg_try_advisory_xact_lock_shared(888888, generation);
19
20 RETURN generation;
21END;
22$function$;
23
24SELECT run_command_on_workers($$
25CREATE OR REPLACE FUNCTION current_generation(rollup text)
26RETURNS int LANGUAGE sql IMMUTABLE
27AS $function$
28 SELECT -1;
29$function$;
30$$);
31
32/* bump new writes to the next generation, always run in a separate transaction */
33CREATE OR REPLACE FUNCTION next_generation(rollup text)
34RETURNS int LANGUAGE sql
35AS $function$
36 UPDATE rollup_generations SET current_generation = current_generation + 1
37 WHERE function_name = rollup RETURNING current_generation;
38$function$;
39
40/* bump the genereation, always run */
41CREATE OR REPLACE FUNCTION do_generation_rollup(rollup text, generation int)
42RETURNS void LANGUAGE plpgsql
43AS $function$
44BEGIN
45 /* block until all transactions for the generation have finished */
46 PERFORM pg_try_advisory_xact_lock(888888, generation);
47
48 /* call the rollup function */
49 EXECUTE format('SELECT %I(%s)', rollup, generation);
50END;
51$function$;
52
53/* bump the genereation, always run */
54CREATE OR REPLACE FUNCTION do_rollup(rollup text)
55RETURNS int LANGUAGE plpgsql
56AS $function$
57DECLARE
58 last_rolled_up int;
59BEGIN
60 PERFORM do_generation_rollup(rollup, g)
61 FROM rollup_generations, generate_series(processed_generation + 1, current_generation - 1) g FOR UPDATE;
62
63 UPDATE rollup_generations
64 SET processed_generation = current_generation - 1
65 WHERE function_name = rollup
66 RETURNING processed_generation
67 INTO last_rolled_up;
68
69 RETURN last_rolled_up;
70END
71$function$;