ml: implement fixed time-based training windows (#20638)
Costa Tsaousis committed
Sep 22, 2025 at 15:01 UTC
4f68bb0b197f2ed7246ae314149afef27afaaa9e
6 files changed
+160
-47
src/ml/ml.cc
+25
-9
@@ -56,13 +56,13 @@ ml_dimension_calculated_numbers(ml_worker_t *worker, ml_dimension_t *dim)
56
training_response.first_entry_on_response = rrddim_first_entry_s_of_tier(dim->rd, 0);
57
training_response.last_entry_on_response = rrddim_last_entry_s_of_tier(dim->rd, 0);
58
59
- size_t min_n = Cfg.min_train_samples;
60
- size_t max_n = Cfg.max_train_samples;
59
+ size_t min_n = Cfg.min_training_window / dim->rd->rrdset->update_every;
60
+ size_t max_n = Cfg.training_window / dim->rd->rrdset->update_every;
61
62
// Figure out what our time window should be.
63
training_response.query_before_t = training_response.last_entry_on_response;
64
training_response.query_after_t = std::max(
65
- training_response.query_before_t - static_cast<time_t>((max_n - 1) * dim->rd->rrdset->update_every),
65
+ training_response.query_before_t - Cfg.training_window, // Fixed time window
66
training_response.first_entry_on_response
67
);
68
@@ -379,7 +379,7 @@ int ml_dimension_load_models(RRDDIM *rd, sqlite3_stmt **active_stmt) {
379
if (unlikely(rc != SQLITE_OK))
380
goto bind_fail;
381
382
- rc = sqlite3_bind_int64(res, ++param, now_realtime_sec() - (Cfg.num_models_to_use * Cfg.max_train_samples));
382
+ rc = sqlite3_bind_int64(res, ++param, now_realtime_sec() - (Cfg.num_models_to_use * Cfg.train_every));
383
if (unlikely(rc != SQLITE_OK))
384
goto bind_fail;
385
@@ -674,13 +674,29 @@ ml_dimension_train_model(ml_worker_t *worker, ml_dimension_t *dim)
674
memcpy(worker->scratch_training_cns, worker->training_cns,
675
training_response.total_values * sizeof(calculated_number_t));
676
677
+ size_t smoothing_window = (dim->rd->rrdset->update_every > nd_profile.update_every) ? 1 : Cfg.max_samples_to_smooth;
678
+
679
ml_features_t features = {
678
- Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
680
+ Cfg.diff_n, smoothing_window, Cfg.lag_n,
681
worker->scratch_training_cns, training_response.total_values,
682
worker->training_cns, training_response.total_values,
683
worker->training_samples
684
};
683
- ml_features_preprocess(&features);
685
+
686
+ // Calculate dynamic sampling ratio based on expected output size
687
+ // After diff and smooth, we'll have approximately this many vectors
688
+ size_t expected_vectors = training_response.total_values;
689
+ if (Cfg.diff_n > 0) expected_vectors--;
690
+ if (smoothing_window > 1) expected_vectors = expected_vectors - smoothing_window + 1;
691
+ expected_vectors = expected_vectors - Cfg.lag_n;
692
+
693
+ double sampling_ratio = 1.0;
694
+ if (expected_vectors > Cfg.max_training_vectors) {
695
+ sampling_ratio = (double)Cfg.max_training_vectors / expected_vectors;
696
+ }
697
+
698
+ // Apply sampling during lag feature extraction
699
+ ml_features_preprocess(&features, sampling_ratio);
700
701
ml_kmeans_init(&dim->kmeans);
702
ml_kmeans_train(&dim->kmeans, &features, Cfg.max_kmeans_iters, training_response.query_after_t, training_response.query_before_t);
@@ -706,7 +722,7 @@ ml_dimension_predict(ml_dimension_t *dim, calculated_number_t value, bool exists
722
}
723
724
// Save the value and return if we don't have enough values for a sample
709
- unsigned n = Cfg.diff_n + Cfg.smooth_n + Cfg.lag_n;
725
+ unsigned n = Cfg.diff_n + Cfg.max_samples_to_smooth + Cfg.lag_n;
726
if (dim->cns.size() < n) {
727
dim->cns.push_back(value);
728
return false;
@@ -731,11 +747,11 @@ ml_dimension_predict(ml_dimension_t *dim, calculated_number_t value, bool exists
747
memcpy(dst_cns, dim->cns.data(), n * sizeof(calculated_number_t));
748
749
ml_features_t features = {
734
- Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
750
+ Cfg.diff_n, Cfg.max_samples_to_smooth, Cfg.lag_n,
751
dst_cns, n, src_cns, n,
752
dim->feature
753
};
738
- ml_features_preprocess(&features);
754
+ ml_features_preprocess(&features, 1.0);
755
756
/*
757
* Lock to predict
src/ml/ml_config.cc
+115
-17
@@ -2,6 +2,102 @@
2
3
#include "ml_config.h"
4
5
+static void ml_config_migrate() {
6
+ const char *config_section_ml = CONFIG_SECTION_ML;
7
+
8
+ // Check if migration is needed by looking for old keys
9
+ bool has_old_keys = false;
10
+ if (inicfg_exists(&netdata_config, config_section_ml, "maximum num samples to train") ||
11
+ inicfg_exists(&netdata_config, config_section_ml, "minimum num samples to train") ||
12
+ inicfg_exists(&netdata_config, config_section_ml, "num samples to diff") ||
13
+ inicfg_exists(&netdata_config, config_section_ml, "num samples to smooth") ||
14
+ inicfg_exists(&netdata_config, config_section_ml, "num samples to lag") ||
15
+ inicfg_exists(&netdata_config, config_section_ml, "random sampling ratio")) {
16
+ has_old_keys = true;
17
+ }
18
+
19
+ // Check if new keys already exist (user manually migrated)
20
+ bool has_new_keys = false;
21
+ if (inicfg_exists(&netdata_config, config_section_ml, "training window") ||
22
+ inicfg_exists(&netdata_config, config_section_ml, "max training vectors")) {
23
+ has_new_keys = true;
24
+ }
25
+
26
+ // Only migrate if we have old keys but no new keys
27
+ if (!has_old_keys || has_new_keys) {
28
+ return;
29
+ }
30
+
31
+ // Get the user's "high resolution" setting
32
+ // This is what their configuration was designed for
33
+ time_t global_update_every = nd_profile.update_every;
34
+
35
+ // Read all old configuration values with defaults
36
+ // Users may have changed only some values, so we need proper defaults
37
+ unsigned old_max_train_samples = inicfg_get_number(&netdata_config, config_section_ml,
38
+ "maximum num samples to train", 21600);
39
+ unsigned old_min_train_samples = inicfg_get_number(&netdata_config, config_section_ml,
40
+ "minimum num samples to train", 900);
41
+ unsigned old_train_every = inicfg_get_duration_seconds(&netdata_config, config_section_ml,
42
+ "train every", 10800);
43
+ unsigned old_diff_n = inicfg_get_number(&netdata_config, config_section_ml,
44
+ "num samples to diff", 1);
45
+ unsigned old_smooth_n = inicfg_get_number(&netdata_config, config_section_ml,
46
+ "num samples to smooth", 3);
47
+ unsigned old_lag_n = inicfg_get_number(&netdata_config, config_section_ml,
48
+ "num samples to lag", 5);
49
+ double old_sampling_ratio = inicfg_get_double(&netdata_config, config_section_ml,
50
+ "random sampling ratio", 0.2);
51
+
52
+ // Calculate time-based equivalents
53
+ // These preserve the exact behavior the user had configured
54
+ time_t training_window = old_max_train_samples * global_update_every;
55
+ time_t min_training_window = old_min_train_samples * global_update_every;
56
+
57
+ // Calculate target training vectors based on old pipeline
58
+ // Account for data reduction from diff, smooth, and sampling
59
+ size_t effective_samples = old_max_train_samples;
60
+ if (old_diff_n > 0) effective_samples--; // Lose one sample to differencing
61
+ size_t max_training_vectors = (size_t)(effective_samples * old_sampling_ratio);
62
+
63
+ // Write new configuration values
64
+ char window_str[32];
65
+ snprintf(window_str, sizeof(window_str), "%ldh", training_window / 3600);
66
+ inicfg_set(&netdata_config, config_section_ml, "training window", window_str);
67
+
68
+ snprintf(window_str, sizeof(window_str), "%ldm", min_training_window / 60);
69
+ inicfg_set(&netdata_config, config_section_ml, "min training window", window_str);
70
+
71
+ inicfg_set_number(&netdata_config, config_section_ml, "max training vectors", max_training_vectors);
72
+ inicfg_set_number(&netdata_config, config_section_ml, "max samples to smooth", old_smooth_n);
73
+
74
+ // Migrate unchanged values
75
+ inicfg_set_duration_seconds(&netdata_config, config_section_ml, "train every", old_train_every);
76
+ inicfg_set_number(&netdata_config, config_section_ml, "num samples to diff", old_diff_n);
77
+ inicfg_set_number(&netdata_config, config_section_ml, "num samples to lag", old_lag_n);
78
+
79
+ // Mark old keys as migrated by moving them to avoid showing in netdata.conf
80
+ // This uses Netdata's config migration pattern
81
+ inicfg_move(&netdata_config, config_section_ml, "maximum num samples to train",
82
+ config_section_ml, "obsolete maximum num samples to train");
83
+ inicfg_move(&netdata_config, config_section_ml, "minimum num samples to train",
84
+ config_section_ml, "obsolete minimum num samples to train");
85
+ inicfg_move(&netdata_config, config_section_ml, "num samples to smooth",
86
+ config_section_ml, "obsolete num samples to smooth");
87
+ inicfg_move(&netdata_config, config_section_ml, "random sampling ratio",
88
+ config_section_ml, "obsolete random sampling ratio");
89
+
90
+ // Log the migration
91
+ nd_log(NDLS_DAEMON, NDLP_NOTICE,
92
+ "ML configuration migrated from sample-based to time-based:");
93
+ nd_log(NDLS_DAEMON, NDLP_NOTICE,
94
+ " Training window: %ld seconds (%ld hours) - was %u samples at %ld second intervals",
95
+ training_window, training_window / 3600, old_max_train_samples, global_update_every);
96
+ nd_log(NDLS_DAEMON, NDLP_NOTICE,
97
+ " Target training vectors: %zu - calculated from smoothing and sampling",
98
+ max_training_vectors);
99
+}
100
+
101
/*
102
* Global configuration instance to be shared between training and
103
* prediction threads.
@@ -19,24 +115,27 @@ static T clamp(const T& Value, const T& Min, const T& Max) {
115
void ml_config_load(ml_config_t *cfg) {
116
const char *config_section_ml = CONFIG_SECTION_ML;
117
118
+ // Migrate old configuration if needed
119
+ ml_config_migrate();
120
+
121
int enable_anomaly_detection = inicfg_get_boolean_ondemand(&netdata_config, config_section_ml, "enabled", nd_profile.ml_enabled);
122
123
/*
124
* Read values
125
*/
126
28
- unsigned max_train_samples = inicfg_get_number(&netdata_config, config_section_ml, "maximum num samples to train", 6 * 3600);
29
- unsigned min_train_samples = inicfg_get_number(&netdata_config, config_section_ml, "minimum num samples to train", 1 * 900);
127
+ time_t training_window = inicfg_get_duration_seconds(&netdata_config, config_section_ml, "training window", 6 * 3600);
128
+ time_t min_training_window = inicfg_get_duration_seconds(&netdata_config, config_section_ml, "min training window", 15 * 60);
129
+ size_t max_training_vectors = inicfg_get_number(&netdata_config, config_section_ml, "max training vectors", 1440);
130
+ size_t max_samples_to_smooth = inicfg_get_number(&netdata_config, config_section_ml, "max samples to smooth", 3);
131
unsigned train_every = inicfg_get_duration_seconds(&netdata_config, config_section_ml, "train every", 3 * 3600);
132
133
unsigned num_models_to_use = inicfg_get_number(&netdata_config, config_section_ml, "number of models per dimension", 18);
134
unsigned delete_models_older_than = inicfg_get_duration_seconds(&netdata_config, config_section_ml, "delete models older than", 60 * 60 * 24 * 7);
135
136
unsigned diff_n = inicfg_get_number(&netdata_config, config_section_ml, "num samples to diff", 1);
36
- unsigned smooth_n = inicfg_get_number(&netdata_config, config_section_ml, "num samples to smooth", 3);
137
unsigned lag_n = inicfg_get_number(&netdata_config, config_section_ml, "num samples to lag", 5);
138
39
- double random_sampling_ratio = inicfg_get_double(&netdata_config, config_section_ml, "random sampling ratio", 1.0 / 5.0 /* default lag_n */);
139
unsigned max_kmeans_iters = inicfg_get_number(&netdata_config, config_section_ml, "maximum number of k-means iterations", 1000);
140
141
double dimension_anomaly_rate_threshold = inicfg_get_double(&netdata_config, config_section_ml, "dimension anomaly score threshold", 0.99);
@@ -64,18 +163,17 @@ void ml_config_load(ml_config_t *cfg) {
163
* Clamp
164
*/
165
67
- max_train_samples = clamp<unsigned>(max_train_samples, 1 * 3600, 24 * 3600);
68
- min_train_samples = clamp<unsigned>(min_train_samples, 1 * 900, 6 * 3600);
166
+ training_window = clamp<time_t>(training_window, 1 * 3600, 24 * 3600);
167
+ min_training_window = clamp<time_t>(min_training_window, 1 * 900, 6 * 3600);
168
train_every = clamp<unsigned>(train_every, 1 * 3600, 6 * 3600);
169
170
num_models_to_use = clamp<unsigned>(num_models_to_use, 1, 7 * 24);
171
delete_models_older_than = clamp<unsigned>(delete_models_older_than, 60 * 60 * 24 * 1, 60 * 60 * 24 * 7);
172
173
diff_n = clamp(diff_n, 0u, 1u);
75
- smooth_n = clamp(smooth_n, 0u, 5u);
174
+ max_samples_to_smooth = clamp<size_t>(max_samples_to_smooth, 0, 5);
175
lag_n = clamp(lag_n, 1u, 5u);
176
78
- random_sampling_ratio = clamp(random_sampling_ratio, 0.2, 1.0);
177
max_kmeans_iters = clamp(max_kmeans_iters, 500u, 1000u);
178
179
dimension_anomaly_rate_threshold = clamp(dimension_anomaly_rate_threshold, 0.01, 5.00);
@@ -86,18 +184,18 @@ void ml_config_load(ml_config_t *cfg) {
184
num_worker_threads = clamp<size_t>(num_worker_threads, 4, netdata_conf_cpus());
185
flush_models_batch_size = clamp<size_t>(flush_models_batch_size, 8, 512);
186
89
- suppression_window = clamp<size_t>(suppression_window, 1, max_train_samples);
187
+ suppression_window = clamp<size_t>(suppression_window, 1, training_window);
188
suppression_threshold = clamp<size_t>(suppression_threshold, 1, suppression_window);
189
190
/*
191
* Validate
192
*/
193
96
- if (min_train_samples >= max_train_samples) {
97
- netdata_log_error("invalid min/max train samples found (%u >= %u)", min_train_samples, max_train_samples);
194
+ if (min_training_window >= training_window) {
195
+ netdata_log_error("invalid min/max training window found (%ld >= %ld)", min_training_window, training_window);
196
99
- min_train_samples = 1 * 3600;
100
- max_train_samples = 6 * 3600;
197
+ min_training_window = 1 * 3600;
198
+ training_window = 6 * 3600;
199
}
200
201
/*
@@ -106,18 +204,18 @@ void ml_config_load(ml_config_t *cfg) {
204
205
cfg->enable_anomaly_detection = enable_anomaly_detection;
206
109
- cfg->max_train_samples = max_train_samples;
110
- cfg->min_train_samples = min_train_samples;
207
+ cfg->training_window = training_window;
208
+ cfg->min_training_window = min_training_window;
209
+ cfg->max_training_vectors = max_training_vectors;
210
+ cfg->max_samples_to_smooth = max_samples_to_smooth;
211
cfg->train_every = train_every;
212
213
cfg->num_models_to_use = num_models_to_use;
214
cfg->delete_models_older_than = delete_models_older_than;
215
216
cfg->diff_n = diff_n;
117
- cfg->smooth_n = smooth_n;
217
cfg->lag_n = lag_n;
218
120
- cfg->random_sampling_ratio = random_sampling_ratio;
219
cfg->max_kmeans_iters = max_kmeans_iters;
220
221
cfg->host_anomaly_rate_threshold = host_anomaly_rate_threshold;
src/ml/ml_config.h
+4
-5
@@ -8,8 +8,10 @@
8
typedef struct {
9
int enable_anomaly_detection;
10
11
- unsigned max_train_samples;
12
- unsigned min_train_samples;
11
+ time_t training_window; // Training window in seconds
12
+ time_t min_training_window; // Minimum training window in seconds
13
+ size_t max_training_vectors; // Target number of vectors for training
14
+ size_t max_samples_to_smooth; // Maximum smoothing window (adaptive)
15
unsigned train_every;
16
17
unsigned num_models_to_use;
@@ -18,10 +20,7 @@ typedef struct {
20
unsigned db_engine_anomaly_rate_every;
21
22
unsigned diff_n;
21
- unsigned smooth_n;
23
unsigned lag_n;
23
-
24
- double random_sampling_ratio;
24
unsigned max_kmeans_iters;
25
26
double dimension_anomaly_score_threshold;
src/ml/ml_features.cc
+4
-7
@@ -41,14 +41,11 @@ static void ml_features_smooth(ml_features_t *features)
41
features->src[(features->src_n - 1) - idx] = 0.0;
42
}
43
44
-static void ml_features_lag(ml_features_t *features)
44
+static void ml_features_lag(ml_features_t *features, double sampling_ratio)
45
{
46
size_t n = features->src_n - features->diff_n - features->smooth_n + 1 - features->lag_n;
47
features->preprocessed_features.resize(n);
48
49
- unsigned target_num_samples = Cfg.max_train_samples * Cfg.random_sampling_ratio;
50
- double sampling_ratio = std::min(static_cast<double>(target_num_samples) / n, 1.0);
51
-
49
uint32_t max_mt = std::numeric_limits<uint32_t>::max();
50
uint32_t cutoff = static_cast<double>(max_mt) * sampling_ratio;
51
@@ -58,7 +55,7 @@ static void ml_features_lag(ml_features_t *features)
55
DSample &DS = features->preprocessed_features[sample_idx++];
56
DS.set_size(features->lag_n);
57
61
- if (Cfg.random_nums[idx] > cutoff) {
58
+ if (Cfg.random_nums[idx % Cfg.random_nums.size()] > cutoff) {
59
sample_idx--;
60
continue;
61
}
@@ -70,9 +67,9 @@ static void ml_features_lag(ml_features_t *features)
67
features->preprocessed_features.resize(sample_idx);
68
}
69
73
-void ml_features_preprocess(ml_features_t *features)
70
+void ml_features_preprocess(ml_features_t *features, double sampling_ratio)
71
{
72
ml_features_diff(features);
73
ml_features_smooth(features);
77
- ml_features_lag(features);
74
+ ml_features_lag(features, sampling_ratio);
75
}
src/ml/ml_features.h
+1
-1
@@ -21,6 +21,6 @@ typedef struct {
21
std::vector<DSample> &preprocessed_features;
22
} ml_features_t;
23
24
-void ml_features_preprocess(ml_features_t *features);
24
+void ml_features_preprocess(ml_features_t *features, double sampling_ratio);
25
26
#endif /* ML_FEATURES_H */
src/ml/ml_public.cc
+11
-8
@@ -141,16 +141,16 @@ void ml_host_get_info(RRDHOST *rh, BUFFER *wb)
141
142
buffer_json_member_add_boolean(wb, "enabled", Cfg.enable_anomaly_detection);
143
144
- buffer_json_member_add_uint64(wb, "min-train-samples", Cfg.min_train_samples);
145
- buffer_json_member_add_uint64(wb, "max-train-samples", Cfg.max_train_samples);
144
+ buffer_json_member_add_uint64(wb, "training-window", Cfg.training_window);
145
+ buffer_json_member_add_uint64(wb, "min-training-window", Cfg.min_training_window);
146
+ buffer_json_member_add_uint64(wb, "max-training-vectors", Cfg.max_training_vectors);
147
+ buffer_json_member_add_uint64(wb, "max-samples-to-smooth", Cfg.max_samples_to_smooth);
148
buffer_json_member_add_uint64(wb, "train-every", Cfg.train_every);
149
150
buffer_json_member_add_uint64(wb, "diff-n", Cfg.diff_n);
149
- buffer_json_member_add_uint64(wb, "smooth-n", Cfg.smooth_n);
151
buffer_json_member_add_uint64(wb, "lag-n", Cfg.lag_n);
152
152
- buffer_json_member_add_double(wb, "random-sampling-ratio", Cfg.random_sampling_ratio);
153
- buffer_json_member_add_uint64(wb, "max-kmeans-iters", Cfg.random_sampling_ratio);
153
+ buffer_json_member_add_uint64(wb, "max-kmeans-iters", Cfg.max_kmeans_iters);
154
155
buffer_json_member_add_double(wb, "dimension-anomaly-score-threshold", Cfg.dimension_anomaly_score_threshold);
156
@@ -369,8 +369,8 @@ void ml_init()
369
std::random_device RD;
370
std::mt19937 Gen(RD());
371
372
- Cfg.random_nums.reserve(Cfg.max_train_samples);
373
- for (size_t Idx = 0; Idx != Cfg.max_train_samples; Idx++)
372
+ Cfg.random_nums.reserve(Cfg.max_training_vectors);
373
+ for (size_t Idx = 0; Idx != Cfg.max_training_vectors; Idx++)
374
Cfg.random_nums.push_back(Gen());
375
376
// init training thread-specific data
@@ -378,7 +378,10 @@ void ml_init()
378
for (size_t idx = 0; idx != Cfg.num_worker_threads; idx++) {
379
ml_worker_t *worker = &Cfg.workers[idx];
380
381
- size_t max_elements_needed_for_training = (size_t) Cfg.max_train_samples * (size_t) (Cfg.lag_n + 1);
381
+ // Calculate max elements needed based on the highest frequency metrics
382
+ // For 1-second metrics: training_window samples
383
+ // We allocate for worst case (1-second update frequency)
384
+ size_t max_elements_needed_for_training = (size_t) Cfg.training_window * (size_t) (Cfg.lag_n + 1);
385
worker->training_cns = new calculated_number_t[max_elements_needed_for_training]();
386
worker->scratch_training_cns = new calculated_number_t[max_elements_needed_for_training]();
387