SET datestyle = 'ISO'; SET max_parallel_workers_per_gather = 0; CREATE SERVER pwagg_svr FOREIGN DATA WRAPPER clickhouse_fdw OPTIONS(dbname 'pwagg_test', driver 'http'); CREATE USER MAPPING FOR CURRENT_USER SERVER pwagg_svr; -- ClickHouse holds cold data CREATE SERVER pwagg_admin FOREIGN DATA WRAPPER clickhouse_fdw; CREATE USER MAPPING FOR CURRENT_USER SERVER pwagg_admin; CALL clickhouse_perform('pwagg_admin', 'DROP DATABASE IF EXISTS pwagg_test'); CALL clickhouse_perform('pwagg_admin', 'CREATE DATABASE pwagg_test'); CALL clickhouse_perform('pwagg_admin', 'CREATE TABLE pwagg_test.events (id Int32, ts Date, val Int32, amt Float64) ENGINE = MergeTree ORDER BY ts'); CALL clickhouse_perform('pwagg_admin', $$INSERT INTO pwagg_test.events VALUES (1,'2023-01-15',10,10),(2,'2023-02-10',20,20),(3,'2023-03-20',30,30),(4,'2023-04-05',40,40)$$); -- Partitioned table whose cold 2023 range lives on ClickHouse as a foreign -- partition while the hot 2024 range stays local: the layout a consumer builds -- when offloading old partitions (offload itself is left to the consumer) CREATE TABLE events (id int, ts date, val int, amt float8) PARTITION BY RANGE (ts); CREATE FOREIGN TABLE events_cold PARTITION OF events FOR VALUES FROM ('2023-01-01') TO ('2024-01-01') SERVER pwagg_svr OPTIONS (table_name 'events'); CREATE TABLE events_hot PARTITION OF events FOR VALUES FROM ('2024-01-01') TO ('2025-01-01'); INSERT INTO events_hot VALUES (100,'2024-01-10',5,5), (101,'2024-02-15',15,15); SET enable_partitionwise_aggregate = on; -- Decomposable aggregates push the cold partial straight to ClickHouse EXPLAIN (VERBOSE, COSTS OFF) SELECT count(*), sum(val), min(ts), max(ts) FROM events; QUERY PLAN ------------------------------------------------------------------------------------------------------------------------- Finalize Aggregate Output: count(*), sum(events.val), min(events.ts), max(events.ts) -> Append -> Foreign Scan Output: (PARTIAL count(*)), (PARTIAL sum(events.val)), (PARTIAL min(events.ts)), (PARTIAL max(events.ts)) Relations: Aggregate on (events_cold events) Remote SQL: SELECT count(*), sum(val), min(ts), max(ts) FROM pwagg_test.events -> Partial Aggregate Output: PARTIAL count(*), PARTIAL sum(events_1.val), PARTIAL min(events_1.ts), PARTIAL max(events_1.ts) -> Seq Scan on public.events_hot events_1 Output: events_1.val, events_1.ts (11 rows) SELECT count(*), sum(val), min(ts), max(ts) FROM events; count | sum | min | max -------+-----+------------+------------ 6 | 120 | 2023-01-15 | 2024-02-15 (1 row) -- avg(int) pushes its transition state as int8[2] {count, sum} EXPLAIN (VERBOSE, COSTS OFF) SELECT avg(val) FROM events; QUERY PLAN -------------------------------------------------------------------------------------------------- Finalize Aggregate Output: avg(events.val) -> Append -> Foreign Scan Output: (PARTIAL avg(events.val)) Relations: Aggregate on (events_cold events) Remote SQL: SELECT [toInt64(count(val)), toInt64(sum(val))] FROM pwagg_test.events -> Partial Aggregate Output: PARTIAL avg(events_1.val) -> Seq Scan on public.events_hot events_1 Output: events_1.val (11 rows) SELECT avg(val) FROM events; avg --------------------- 20.0000000000000000 (1 row) -- avg/var/stddev over float push float8[3] {N, sum, sum of squared deviations} EXPLAIN (VERBOSE, COSTS OFF) SELECT var_samp(amt) FROM events; QUERY PLAN ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- Finalize Aggregate Output: var_samp(events.amt) -> Append -> Foreign Scan Output: (PARTIAL var_samp(events.amt)) Relations: Aggregate on (events_cold events) Remote SQL: SELECT [toFloat64(count(amt)), sum(toFloat64(amt)), if(count(amt) > 0, sum(pow(toFloat64(amt), 2)) - pow(sum(toFloat64(amt)), 2) / count(amt), 0)] FROM pwagg_test.events -> Partial Aggregate Output: PARTIAL var_samp(events_1.amt) -> Seq Scan on public.events_hot events_1 Output: events_1.amt (11 rows) SELECT avg(amt), var_pop(amt), var_samp(amt), stddev_pop(amt), stddev_samp(amt) FROM events; avg | var_pop | var_samp | stddev_pop | stddev_samp -----+--------------------+----------+--------------------+-------------------- 20 | 141.66666666666666 | 170 | 11.902380714238083 | 13.038404810405298 (1 row) -- var_samp(int) keeps an INTERNAL numeric state, so it falls back to fetching -- the cold rows and aggregating locally EXPLAIN (VERBOSE, COSTS OFF) SELECT var_samp(val) FROM events; QUERY PLAN ------------------------------------------------------------------- Finalize Aggregate Output: var_samp(events.val) -> Append -> Partial Aggregate Output: PARTIAL var_samp(events.val) -> Foreign Scan on public.events_cold events Output: events.val Remote SQL: SELECT val FROM pwagg_test.events -> Partial Aggregate Output: PARTIAL var_samp(events_1.val) -> Seq Scan on public.events_hot events_1 Output: events_1.val (12 rows) SELECT var_samp(val) FROM events; var_samp ---------------------- 170.0000000000000000 (1 row) -- FILTER pushes too, as ClickHouse -If on each transition-state component EXPLAIN (VERBOSE, COSTS OFF) SELECT avg(val) FILTER (WHERE val > 15) FROM events; QUERY PLAN ------------------------------------------------------------------------------------------------------------------------------------------ Finalize Aggregate Output: avg(events.val) FILTER (WHERE (events.val > 15)) -> Append -> Foreign Scan Output: (PARTIAL avg(events.val) FILTER (WHERE (events.val > 15))) Relations: Aggregate on (events_cold events) Remote SQL: SELECT [toInt64(countIf(val, ((val > 15)) > 0)), toInt64(sumIf(val, ((val > 15)) > 0))] FROM pwagg_test.events -> Partial Aggregate Output: PARTIAL avg(events_1.val) FILTER (WHERE (events_1.val > 15)) -> Seq Scan on public.events_hot events_1 Output: events_1.val (11 rows) SELECT avg(val) FILTER (WHERE val > 15), avg(amt) FILTER (WHERE amt > 15), var_samp(amt) FILTER (WHERE amt > 15) FROM events; avg | avg | var_samp ---------------------+-----+---------- 30.0000000000000000 | 30 | 100 (1 row) RESET enable_partitionwise_aggregate; DROP TABLE events; DROP USER MAPPING FOR CURRENT_USER SERVER pwagg_svr; CALL clickhouse_perform('pwagg_admin', 'DROP DATABASE pwagg_test'); DROP SERVER pwagg_svr CASCADE;