Skip to content

Commit c4a71ff

Browse files
committed
Introduce continuous_aggs_tenant_tracking catalog
Add catalog definitions and helper functions for continuous_aggs_tenant_tracking. This new catalog table is used to persist granular trackings, so we can perform granular refresh for continuous aggregates.
1 parent 5ebc0fa commit c4a71ff

10 files changed

Lines changed: 233 additions & 2 deletions

File tree

sql/pre_install/tables.sql

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -440,6 +440,21 @@ SELECT pg_catalog.pg_extension_config_dump('_timescaledb_catalog.continuous_aggs
440440

441441
CREATE INDEX continuous_aggs_jobs_refresh_ranges_idx ON _timescaledb_catalog.continuous_aggs_jobs_refresh_ranges (materialization_id);
442442

443+
-- Per-tenant invalidation tracking, used in granular refresh of contiuous aggregates.
444+
CREATE TABLE _timescaledb_catalog.continuous_aggs_tenant_tracking (
445+
hypertable_id integer NOT NULL,
446+
tenant_id text,
447+
min_timestamp bigint,
448+
max_timestamp bigint,
449+
seqnum integer NOT NULL,
450+
-- table constraints
451+
CONSTRAINT continuous_aggs_tenant_tracking_hypertable_id_fkey FOREIGN KEY (hypertable_id) REFERENCES _timescaledb_catalog.hypertable (id) ON DELETE CASCADE
452+
);
453+
454+
SELECT pg_catalog.pg_extension_config_dump('_timescaledb_catalog.continuous_aggs_tenant_tracking', '');
455+
456+
CREATE INDEX continuous_aggs_tenant_tracking_idx ON _timescaledb_catalog.continuous_aggs_tenant_tracking (hypertable_id, seqnum);
457+
443458
/* the source of this data is the enum from the source code that lists
444459
* the algorithms. This table is NOT dumped.
445460
*/

sql/updates/latest-dev.sql

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,20 @@ $$;
5959
-- END hypertable.compressed_hypertable_id no longer used
6060
--
6161

62+
-- add continuous_aggs_tenant_tracking
63+
CREATE TABLE _timescaledb_catalog.continuous_aggs_tenant_tracking (
64+
hypertable_id integer NOT NULL,
65+
tenant_id text,
66+
min_timestamp bigint,
67+
max_timestamp bigint,
68+
seqnum integer NOT NULL,
69+
CONSTRAINT continuous_aggs_tenant_tracking_hypertable_id_fkey FOREIGN KEY (hypertable_id) REFERENCES _timescaledb_catalog.hypertable (id) ON DELETE CASCADE
70+
);
71+
72+
SELECT pg_catalog.pg_extension_config_dump('_timescaledb_catalog.continuous_aggs_tenant_tracking', '');
73+
74+
CREATE INDEX continuous_aggs_tenant_tracking_idx ON _timescaledb_catalog.continuous_aggs_tenant_tracking (hypertable_id, seqnum);
75+
6276
--
6377
-- BEGIN set compression status flag on hypertables
6478
--

sql/updates/reverse-dev.sql

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,3 +260,6 @@ DROP FUNCTION IF EXISTS _timescaledb_functions.policy_compaction_check(JSONB);
260260

261261
DROP FUNCTION IF EXISTS @extschema@.alter_job(job_id INTEGER, schedule_interval INTERVAL, max_runtime INTERVAL, max_retries INTEGER, retry_period INTERVAL, scheduled BOOL, config JSONB, next_start TIMESTAMPTZ, if_exists BOOL, check_config REGPROC, fixed_schedule BOOL, initial_start TIMESTAMPTZ, timezone TEXT, job_name TEXT, config_merge JSONB);
262262

263+
-- drop continuous_aggs_tenant_tracking
264+
ALTER EXTENSION timescaledb DROP TABLE _timescaledb_catalog.continuous_aggs_tenant_tracking;
265+
DROP TABLE _timescaledb_catalog.continuous_aggs_tenant_tracking;

src/ts_catalog/CMakeLists.txt

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ set(SOURCES
77
${CMAKE_CURRENT_SOURCE_DIR}/compression_settings.c
88
${CMAKE_CURRENT_SOURCE_DIR}/continuous_agg.c
99
${CMAKE_CURRENT_SOURCE_DIR}/continuous_aggs_jobs_refresh_ranges.c
10+
${CMAKE_CURRENT_SOURCE_DIR}/continuous_aggs_tenant_tracking.c
1011
${CMAKE_CURRENT_SOURCE_DIR}/continuous_aggs_watermark.c
1112
${CMAKE_CURRENT_SOURCE_DIR}/metadata.c
1213
${CMAKE_CURRENT_SOURCE_DIR}/tablespace.c)

src/ts_catalog/catalog.c

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,10 @@ static const TableInfoDef catalog_table_names[_MAX_CATALOG_TABLES + 1] = {
9393
.schema_name = CATALOG_SCHEMA_NAME,
9494
.table_name = CONTINUOUS_AGGS_JOBS_REFRESH_RANGES_TABLE_NAME,
9595
},
96+
[CONTINUOUS_AGGS_TENANT_TRACKING] = {
97+
.schema_name = CATALOG_SCHEMA_NAME,
98+
.table_name = CONTINUOUS_AGGS_TENANT_TRACKING_TABLE_NAME,
99+
},
96100
[COMPRESSION_SETTINGS] = {
97101
.schema_name = CATALOG_SCHEMA_NAME,
98102
.table_name = COMPRESSION_SETTINGS_TABLE_NAME,
@@ -242,6 +246,12 @@ static const TableIndexDef catalog_table_index_definitions[_MAX_CATALOG_TABLES]
242246
[CONTINUOUS_AGGS_JOBS_REFRESH_RANGES_IDX] = "continuous_aggs_jobs_refresh_ranges_idx",
243247
},
244248
},
249+
[CONTINUOUS_AGGS_TENANT_TRACKING] = {
250+
.length = _MAX_CONTINUOUS_AGGS_TENANT_TRACKING_INDEX,
251+
.names = (char *[]) {
252+
[CONTINUOUS_AGGS_TENANT_TRACKING_IDX] = "continuous_aggs_tenant_tracking_idx",
253+
},
254+
},
245255
[CONTINUOUS_AGGS_WATERMARK] = {
246256
.length = _MAX_CONTINUOUS_AGGS_WATERMARK_INDEX,
247257
.names = (char *[]) {

src/ts_catalog/catalog.h

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ typedef enum CatalogTable
4949
CONTINUOUS_AGGS_MATERIALIZATION_INVALIDATION_LOG,
5050
CONTINUOUS_AGGS_MATERIALIZATION_RANGES,
5151
CONTINUOUS_AGGS_JOBS_REFRESH_RANGES,
52+
CONTINUOUS_AGGS_TENANT_TRACKING,
5253
COMPRESSION_SETTINGS,
5354
COMPRESSION_CHUNK_SIZE,
5455
CONTINUOUS_AGGS_BUCKET_FUNCTION,
@@ -1080,6 +1081,47 @@ enum
10801081
_MAX_CONTINUOUS_AGGS_JOBS_REFRESH_RANGES_INDEX,
10811082
};
10821083

1084+
/****** CONTINUOUS_AGGS_TENANT_TRACKING definitions */
1085+
#define CONTINUOUS_AGGS_TENANT_TRACKING_TABLE_NAME "continuous_aggs_tenant_tracking"
1086+
1087+
typedef enum Anum_continuous_aggs_tenant_tracking
1088+
{
1089+
Anum_continuous_aggs_tenant_tracking_hypertable_id = 1,
1090+
Anum_continuous_aggs_tenant_tracking_tenant_id,
1091+
Anum_continuous_aggs_tenant_tracking_min_timestamp,
1092+
Anum_continuous_aggs_tenant_tracking_max_timestamp,
1093+
Anum_continuous_aggs_tenant_tracking_seqnum,
1094+
_Anum_continuous_aggs_tenant_tracking_max,
1095+
} Anum_continuous_aggs_tenant_tracking;
1096+
1097+
#define Natts_continuous_aggs_tenant_tracking (_Anum_continuous_aggs_tenant_tracking_max - 1)
1098+
1099+
typedef struct FormData_continuous_aggs_tenant_tracking
1100+
{
1101+
int32 hypertable_id;
1102+
text *tenant_id;
1103+
int64 min_timestamp;
1104+
int64 max_timestamp;
1105+
int32 seqnum;
1106+
} FormData_continuous_aggs_tenant_tracking;
1107+
1108+
typedef FormData_continuous_aggs_tenant_tracking *Form_continuous_aggs_tenant_tracking;
1109+
1110+
enum
1111+
{
1112+
CONTINUOUS_AGGS_TENANT_TRACKING_IDX = 0,
1113+
_MAX_CONTINUOUS_AGGS_TENANT_TRACKING_INDEX,
1114+
};
1115+
typedef enum Anum_continuous_aggs_tenant_tracking_idx
1116+
{
1117+
Anum_continuous_aggs_tenant_tracking_idx_hypertable_id = 1,
1118+
Anum_continuous_aggs_tenant_tracking_idx_seqnum,
1119+
_Anum_continuous_aggs_tenant_tracking_idx_max,
1120+
} Anum_continuous_aggs_tenant_tracking_idx;
1121+
1122+
#define Natts_continuous_aggs_tenant_tracking_idx \
1123+
(_Anum_continuous_aggs_tenant_tracking_idx_max - 1)
1124+
10831125
typedef enum Anum_continuous_aggs_jobs_refresh_ranges_idx
10841126
{
10851127
Anum_continuous_aggs_jobs_refresh_ranges_idx_materialization_id = 1,
Lines changed: 115 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,115 @@
1+
/*
2+
* This file and its contents are licensed under the Apache License 2.0.
3+
* Please see the included NOTICE for copyright information and
4+
* LICENSE-APACHE for a copy of the license.
5+
*/
6+
7+
#include <postgres.h>
8+
#include <access/table.h>
9+
#include <utils/rel.h>
10+
11+
#include "scan_iterator.h"
12+
#include "ts_catalog/catalog.h"
13+
#include "ts_catalog/continuous_aggs_tenant_tracking.h"
14+
15+
static void
16+
init_scan_by_hypertable_id(ScanIterator *iterator, const int32 hypertable_id)
17+
{
18+
iterator->ctx.index = catalog_get_index(ts_catalog_get(),
19+
CONTINUOUS_AGGS_TENANT_TRACKING,
20+
CONTINUOUS_AGGS_TENANT_TRACKING_IDX);
21+
22+
ts_scan_iterator_scan_key_init(iterator,
23+
Anum_continuous_aggs_tenant_tracking_idx_hypertable_id,
24+
BTEqualStrategyNumber,
25+
F_INT4EQ,
26+
Int32GetDatum(hypertable_id));
27+
}
28+
29+
/*
30+
* Insert one tenant tracking row. tenant_id is copied into a bytea holding the
31+
* exact tenant key bytes.
32+
*/
33+
TSDLLEXPORT void
34+
ts_cagg_tenant_tracking_insert(int32 hypertable_id, const char *tenant_id, int tenant_id_len,
35+
int64 min_timestamp, int64 max_timestamp, int32 seqnum)
36+
{
37+
Catalog *catalog = ts_catalog_get();
38+
Relation rel = table_open(catalog_get_table_id(catalog, CONTINUOUS_AGGS_TENANT_TRACKING),
39+
RowExclusiveLock);
40+
TupleDesc desc = RelationGetDescr(rel);
41+
Datum values[Natts_continuous_aggs_tenant_tracking];
42+
bool nulls[Natts_continuous_aggs_tenant_tracking] = { false };
43+
CatalogSecurityContext sec_ctx;
44+
bytea *tenant_bytea = (bytea *) palloc(VARHDRSZ + tenant_id_len);
45+
46+
SET_VARSIZE(tenant_bytea, VARHDRSZ + tenant_id_len);
47+
memcpy(VARDATA(tenant_bytea), tenant_id, tenant_id_len);
48+
49+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_hypertable_id)] =
50+
Int32GetDatum(hypertable_id);
51+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_tenant_id)] =
52+
PointerGetDatum(tenant_bytea);
53+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_min_timestamp)] =
54+
Int64GetDatum(min_timestamp);
55+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_max_timestamp)] =
56+
Int64GetDatum(max_timestamp);
57+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_seqnum)] =
58+
Int32GetDatum(seqnum);
59+
60+
ts_catalog_database_info_become_owner(ts_catalog_database_info_get(), &sec_ctx);
61+
ts_catalog_insert_values(rel, desc, values, nulls);
62+
ts_catalog_restore_user(&sec_ctx);
63+
64+
table_close(rel, NoLock);
65+
}
66+
67+
/*
68+
* Insert the "invalid" marker row <null, null, null, seqnum>, signalling that
69+
* tenant tracking for this seqnum is incomplete and the refresh must fall back
70+
* to the full invalidation log.
71+
*/
72+
TSDLLEXPORT void
73+
ts_cagg_tenant_tracking_insert_invalid_marker(int32 hypertable_id, int32 seqnum)
74+
{
75+
Catalog *catalog = ts_catalog_get();
76+
Relation rel = table_open(catalog_get_table_id(catalog, CONTINUOUS_AGGS_TENANT_TRACKING),
77+
RowExclusiveLock);
78+
TupleDesc desc = RelationGetDescr(rel);
79+
Datum values[Natts_continuous_aggs_tenant_tracking] = { 0 };
80+
bool nulls[Natts_continuous_aggs_tenant_tracking] = { false };
81+
CatalogSecurityContext sec_ctx;
82+
83+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_hypertable_id)] =
84+
Int32GetDatum(hypertable_id);
85+
nulls[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_tenant_id)] = true;
86+
nulls[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_min_timestamp)] = true;
87+
nulls[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_max_timestamp)] = true;
88+
values[AttrNumberGetAttrOffset(Anum_continuous_aggs_tenant_tracking_seqnum)] =
89+
Int32GetDatum(seqnum);
90+
91+
ts_catalog_database_info_become_owner(ts_catalog_database_info_get(), &sec_ctx);
92+
ts_catalog_insert_values(rel, desc, values, nulls);
93+
ts_catalog_restore_user(&sec_ctx);
94+
95+
table_close(rel, NoLock);
96+
}
97+
98+
/* Delete all tracking rows for a hypertable. */
99+
TSDLLEXPORT void
100+
ts_cagg_tenant_tracking_delete_by_hypertable_id(int32 hypertable_id)
101+
{
102+
ScanIterator iterator = ts_scan_iterator_create(CONTINUOUS_AGGS_TENANT_TRACKING,
103+
RowExclusiveLock,
104+
CurrentMemoryContext);
105+
106+
init_scan_by_hypertable_id(&iterator, hypertable_id);
107+
108+
ts_scanner_foreach(&iterator)
109+
{
110+
TupleInfo *ti = ts_scan_iterator_tuple_info(&iterator);
111+
112+
ts_catalog_delete_tid(ti->scanrel, ts_scanner_get_tuple_tid(ti));
113+
}
114+
ts_scan_iterator_close(&iterator);
115+
}
Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
/*
2+
* This file and its contents are licensed under the Apache License 2.0.
3+
* Please see the included NOTICE for copyright information and
4+
* LICENSE-APACHE for a copy of the license.
5+
*/
6+
#pragma once
7+
8+
#include <postgres.h>
9+
10+
#include "export.h"
11+
12+
/*
13+
* Catalog access for _timescaledb_catalog.continuous_aggs_tenant_tracking.
14+
*
15+
* Rows are flushed here from the shared-memory per-tenant invalidation tracker
16+
* during a continuous aggregate refresh (Transaction 2). tenant_id holds the
17+
* exact tenant key bytes; min/max_timestamp are internal time values. The
18+
* "invalid marker" row (tenant_id/min/max NULL) signals that tenant tracking
19+
* for that seqnum is incomplete, forcing a fall back to the full invalidation
20+
* log.
21+
*/
22+
23+
extern TSDLLEXPORT void ts_cagg_tenant_tracking_insert(int32 hypertable_id, const char *tenant_id,
24+
int tenant_id_len, int64 min_timestamp,
25+
int64 max_timestamp, int32 seqnum);
26+
27+
extern TSDLLEXPORT void ts_cagg_tenant_tracking_insert_invalid_marker(int32 hypertable_id,
28+
int32 seqnum);
29+
30+
extern TSDLLEXPORT void ts_cagg_tenant_tracking_delete_by_hypertable_id(int32 hypertable_id);

test/expected/drop_rename_hypertable.out

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -195,6 +195,7 @@ SELECT schema, name FROM test.relation WHERE schema IN ('public', '_timescaledb_
195195
_timescaledb_catalog | continuous_aggs_jobs_refresh_ranges
196196
_timescaledb_catalog | continuous_aggs_materialization_invalidation_log
197197
_timescaledb_catalog | continuous_aggs_materialization_ranges
198+
_timescaledb_catalog | continuous_aggs_tenant_tracking
198199
_timescaledb_catalog | continuous_aggs_watermark
199200
_timescaledb_catalog | dimension
200201
_timescaledb_catalog | dimension_slice

tsl/test/expected/explain_wal.out

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,9 @@ INSERT INTO metrics SELECT '2025-01-01', 'd1', i::float FROM generate_series(0,1
2626
--- QUERY PLAN ---
2727
Custom Scan (ModifyHypertable) (actual rows=0.00 loops=1)
2828
Direct Compress: true
29-
WAL: records=3 bytes=487
29+
WAL: records=2 bytes=401
3030
-> Insert on metrics (actual rows=0.00 loops=1)
31-
WAL: records=3 bytes=487
31+
WAL: records=2 bytes=401
3232
-> Function Scan on generate_series i (actual rows=11.00 loops=1)
3333

3434
EXPLAIN (ANALYZE, BUFFERS OFF, COSTS OFF, SUMMARY OFF, TIMING OFF, WAL ON)

0 commit comments

Comments
 (0)