Keep a count of metrics and samples collected (#17042)
* count metrics and sample points * Adjust calculation for samples when flushing a hot page * Restore deleted call to mrg_metric_set_clean_latest_time_s * Simplify code * Additional check when removing samples
Stelios Fragkakis committed
Feb 22, 2024 at 18:31 UTC
2e841054371338340137d327c8a9d1e6cb6074ac
9 files changed
+105
-7
src/database/contexts/api_v2.c
+2
@@ -1131,6 +1131,8 @@ void buffer_json_agents_v2(BUFFER *wb, struct query_timings *timings, time_t now
1131
1132
buffer_json_add_array_item_object(wb);
1133
buffer_json_member_add_uint64(wb, "tier", tier);
1134
+ buffer_json_member_add_uint64(wb, "metrics", storage_engine_metrics(eng->seb, localhost->db[tier].si));
1135
+ buffer_json_member_add_uint64(wb, "samples", storage_engine_samples(eng->seb, localhost->db[tier].si));
1136
1137
if(used || max) {
1138
buffer_json_member_add_uint64(wb, "disk_used", used);
src/database/engine/journalfile.c
+7
-1
@@ -708,8 +708,14 @@ static void journalfile_restore_extent_metadata(struct rrdengine_instance *ctx,
708
709
bool added;
710
metric = mrg_metric_add_and_acquire(main_mrg, entry, &added);
711
- if(added)
711
+ if(added) {
712
+ __atomic_add_fetch(&ctx->atomic.metrics, 1, __ATOMIC_RELAXED);
713
update_metric_time = false;
714
+ }
715
+ if (vd.update_every_s) {
716
+ uint64_t samples = (vd.end_time_s - vd.start_time_s) / vd.update_every_s;
717
+ __atomic_add_fetch(&ctx->atomic.samples, samples, __ATOMIC_RELAXED);
718
+ }
719
}
720
Word_t metric_id = mrg_metric_id(main_mrg, metric);
721
src/database/engine/metric.c
+27
-1
@@ -531,6 +531,11 @@ inline bool mrg_metric_set_hot_latest_time_s(MRG *mrg __maybe_unused, METRIC *me
531
return false;
532
}
533
534
+inline time_t mrg_metric_get_latest_clean_time_s(MRG *mrg __maybe_unused, METRIC *metric) {
535
+ time_t clean = __atomic_load_n(&metric->latest_time_s_clean, __ATOMIC_RELAXED);
536
+ return clean;
537
+}
538
+
539
inline time_t mrg_metric_get_latest_time_s(MRG *mrg __maybe_unused, METRIC *metric) {
540
time_t clean = __atomic_load_n(&metric->latest_time_s_clean, __ATOMIC_RELAXED);
541
time_t hot = __atomic_load_n(&metric->latest_time_s_hot, __ATOMIC_RELAXED);
@@ -641,9 +646,30 @@ inline void mrg_update_metric_retention_and_granularity_by_uuid(
646
metric = mrg_metric_add_and_acquire(mrg, entry, &added);
647
}
648
644
- if (likely(!added))
649
+ struct rrdengine_instance *ctx = (struct rrdengine_instance *) section;
650
+ if (likely(!added)) {
651
+ uint64_t old_samples = 0;
652
+
653
+ if (update_every_s && metric->latest_update_every_s && metric->latest_time_s_clean)
654
+ old_samples = (metric->latest_time_s_clean - metric->first_time_s) / metric->latest_update_every_s;
655
+
656
mrg_metric_expand_retention(mrg, metric, first_time_s, last_time_s, update_every_s);
657
658
+ uint64_t new_samples = 0;
659
+ if (update_every_s && metric->latest_update_every_s && metric->latest_time_s_clean)
660
+ new_samples = (metric->latest_time_s_clean - metric->first_time_s) / metric->latest_update_every_s;
661
+
662
+ __atomic_add_fetch(&ctx->atomic.samples, new_samples - old_samples, __ATOMIC_RELAXED);
663
+ }
664
+ else {
665
+ // Newly added
666
+ if (update_every_s) {
667
+ uint64_t samples = (last_time_s - first_time_s) / update_every_s;
668
+ __atomic_add_fetch(&ctx->atomic.samples, samples, __ATOMIC_RELAXED);
669
+ }
670
+ __atomic_add_fetch(&ctx->atomic.metrics, 1, __ATOMIC_RELAXED);
671
+ }
672
+
673
mrg_metric_release(mrg, metric);
674
}
675
src/database/engine/metric.h
+1
@@ -69,6 +69,7 @@ time_t mrg_metric_get_first_time_s(MRG *mrg, METRIC *metric);
69
bool mrg_metric_set_clean_latest_time_s(MRG *mrg, METRIC *metric, time_t latest_time_s);
70
bool mrg_metric_set_hot_latest_time_s(MRG *mrg, METRIC *metric, time_t latest_time_s);
71
time_t mrg_metric_get_latest_time_s(MRG *mrg, METRIC *metric);
72
+time_t mrg_metric_get_latest_clean_time_s(MRG *mrg, METRIC *metric);
73
74
bool mrg_metric_set_update_every(MRG *mrg, METRIC *metric, uint32_t update_every_s);
75
bool mrg_metric_set_update_every_s_if_zero(MRG *mrg, METRIC *metric, uint32_t update_every_s);
src/database/engine/rrdengine.c
+19
-1
@@ -1171,7 +1171,17 @@ static void update_metrics_first_time_s(struct rrdengine_instance *ctx, struct r
1171
for (size_t index = 0; index < added; ++index) {
1172
uuid_first_t_entry = &uuid_first_entry_list[index];
1173
if (likely(uuid_first_t_entry->first_time_s != LONG_MAX)) {
1174
- mrg_metric_set_first_time_s_if_bigger(main_mrg, uuid_first_t_entry->metric, uuid_first_t_entry->first_time_s);
1174
+
1175
+ time_t old_first_time_s = mrg_metric_get_first_time_s(main_mrg, uuid_first_t_entry->metric);
1176
+
1177
+ bool changed = mrg_metric_set_first_time_s_if_bigger(main_mrg, uuid_first_t_entry->metric, uuid_first_t_entry->first_time_s);
1178
+ if (changed) {
1179
+ uint32_t update_every_s = mrg_metric_get_update_every_s(main_mrg, uuid_first_t_entry->metric);
1180
+ if (update_every_s && old_first_time_s && uuid_first_t_entry->first_time_s > old_first_time_s) {
1181
+ uint64_t remove_samples = (uuid_first_t_entry->first_time_s - old_first_time_s) / update_every_s;
1182
+ __atomic_sub_fetch(&ctx->atomic.samples, remove_samples, __ATOMIC_RELAXED);
1183
+ }
1184
+ }
1185
mrg_metric_release(main_mrg, uuid_first_t_entry->metric);
1186
}
1187
else {
@@ -1180,6 +1190,14 @@ static void update_metrics_first_time_s(struct rrdengine_instance *ctx, struct r
1190
// there is no retention for this metric
1191
bool has_retention = mrg_metric_zero_disk_retention(main_mrg, uuid_first_t_entry->metric);
1192
if (!has_retention) {
1193
+ time_t first_time_s = mrg_metric_get_first_time_s(main_mrg, uuid_first_t_entry->metric);
1194
+ time_t last_time_s = mrg_metric_get_latest_time_s(main_mrg, uuid_first_t_entry->metric);
1195
+ time_t update_every_s = mrg_metric_get_update_every_s(main_mrg, uuid_first_t_entry->metric);
1196
+ if (update_every_s && first_time_s && last_time_s) {
1197
+ uint64_t remove_samples = (first_time_s - last_time_s) / update_every_s;
1198
+ __atomic_sub_fetch(&ctx->atomic.samples, remove_samples, __ATOMIC_RELAXED);
1199
+ }
1200
+
1201
bool deleted = mrg_metric_release_and_delete(main_mrg, uuid_first_t_entry->metric);
1202
if(deleted)
1203
deleted_metrics++;
src/database/engine/rrdengine.h
+2
@@ -387,6 +387,8 @@ struct rrdengine_instance {
387
unsigned extents_currently_being_flushed; // non-zero until we commit data to disk (both datafile and journal file)
388
389
time_t first_time_s;
390
+ uint64_t metrics;
391
+ uint64_t samples;
392
} atomic;
393
394
struct {
src/database/engine/rrdengineapi.c
+26
-1
@@ -145,7 +145,10 @@ static METRIC *rrdeng_metric_create(STORAGE_INSTANCE *si, uuid_t *uuid) {
145
.latest_update_every_s = 0,
146
};
147
148
- METRIC *metric = mrg_metric_add_and_acquire(main_mrg, entry, NULL);
148
+ bool added;
149
+ METRIC *metric = mrg_metric_add_and_acquire(main_mrg, entry, &added);
150
+ if (added)
151
+ __atomic_add_fetch(&ctx->atomic.metrics, 1, __ATOMIC_RELAXED);
152
return metric;
153
}
154
@@ -307,6 +310,16 @@ void rrdeng_store_metric_flush_current_page(STORAGE_COLLECT_HANDLE *sch) {
310
else {
311
check_completed_page_consistency(handle);
312
mrg_metric_set_clean_latest_time_s(main_mrg, handle->metric, pgc_page_end_time_s(handle->pgc_page));
313
+
314
+ struct rrdengine_instance *ctx = mrg_metric_ctx(handle->metric);
315
+ time_t start_time_s = pgc_page_start_time_s(handle->pgc_page);
316
+ time_t end_time_s = pgc_page_end_time_s(handle->pgc_page);
317
+ uint32_t update_every_s = mrg_metric_get_update_every_s(main_mrg, handle->metric);
318
+ if (end_time_s && start_time_s && end_time_s > start_time_s && update_every_s) {
319
+ uint64_t add_samples = (end_time_s - start_time_s) / update_every_s;
320
+ __atomic_add_fetch(&ctx->atomic.samples, add_samples, __ATOMIC_RELAXED);
321
+ }
322
+
323
pgc_page_hot_to_dirty_and_release(main_cache, handle->pgc_page);
324
}
325
@@ -967,6 +980,16 @@ uint64_t rrdeng_disk_space_used(STORAGE_INSTANCE *si) {
980
return __atomic_load_n(&ctx->atomic.current_disk_space, __ATOMIC_RELAXED);
981
}
982
983
+uint64_t rrdeng_metrics(STORAGE_INSTANCE *si) {
984
+ struct rrdengine_instance *ctx = (struct rrdengine_instance *)si;
985
+ return __atomic_load_n(&ctx->atomic.metrics, __ATOMIC_RELAXED);
986
+}
987
+
988
+uint64_t rrdeng_samples(STORAGE_INSTANCE *si) {
989
+ struct rrdengine_instance *ctx = (struct rrdengine_instance *)si;
990
+ return __atomic_load_n(&ctx->atomic.samples, __ATOMIC_RELAXED);
991
+}
992
+
993
time_t rrdeng_global_first_time_s(STORAGE_INSTANCE *si) {
994
struct rrdengine_instance *ctx = (struct rrdengine_instance *)si;
995
@@ -1154,6 +1177,8 @@ int rrdeng_init(struct rrdengine_instance **ctxp, const char *dbfiles_path,
1177
1178
rw_spinlock_init(&ctx->njfv2idx.spinlock);
1179
ctx->atomic.first_time_s = LONG_MAX;
1180
+ ctx->atomic.metrics = 0;
1181
+ ctx->atomic.samples = 0;
1182
1183
if (rrdeng_dbengine_spawn(ctx) && !init_rrd_files(ctx)) {
1184
// success - we run this ctx too
src/database/engine/rrdengineapi.h
-3
@@ -223,7 +223,4 @@ RRDENG_SIZE_STATS rrdeng_size_statistics(struct rrdengine_instance *ctx);
223
size_t rrdeng_collectors_running(struct rrdengine_instance *ctx);
224
bool rrdeng_is_legacy(STORAGE_INSTANCE *si);
225
226
-uint64_t rrdeng_disk_space_max(STORAGE_INSTANCE *si);
227
-uint64_t rrdeng_disk_space_used(STORAGE_INSTANCE *si);
228
-
226
#endif /* NETDATA_RRDENGINEAPI_H */
src/database/rrd.h
+21
@@ -446,6 +446,27 @@ static inline uint64_t storage_engine_disk_space_used(STORAGE_ENGINE_BACKEND seb
446
return 0;
447
}
448
449
+uint64_t rrdeng_metrics(STORAGE_INSTANCE *si);
450
+static inline uint64_t storage_engine_metrics(STORAGE_ENGINE_BACKEND seb __maybe_unused, STORAGE_INSTANCE *si __maybe_unused) {
451
+#ifdef ENABLE_DBENGINE
452
+ if(likely(seb == STORAGE_ENGINE_BACKEND_DBENGINE))
453
+ return rrdeng_metrics(si);
454
+#endif
455
+
456
+ // TODO - calculate the total host disk space for memory mode save and map
457
+ return 0;
458
+}
459
+
460
+uint64_t rrdeng_samples(STORAGE_INSTANCE *si);
461
+static inline uint64_t storage_engine_samples(STORAGE_ENGINE_BACKEND seb __maybe_unused, STORAGE_INSTANCE *si __maybe_unused) {
462
+#ifdef ENABLE_DBENGINE
463
+ if(likely(seb == STORAGE_ENGINE_BACKEND_DBENGINE))
464
+ return rrdeng_samples(si);
465
+#endif
466
+ return 0;
467
+}
468
+
469
+
470
time_t rrdeng_global_first_time_s(STORAGE_INSTANCE *si);
471
static inline time_t storage_engine_global_first_time_s(STORAGE_ENGINE_BACKEND seb __maybe_unused, STORAGE_INSTANCE *si __maybe_unused) {
472
#ifdef ENABLE_DBENGINE