Revert "ml: implement fixed time-based training windows (#20638)" (#21045)
This reverts commit 4f68bb0b197f2ed7246ae314149afef27afaaa9e.
vkalintiris committed
Sep 24, 2025 at 20:41 UTC
e4d44b5df954c7c3651372abf032c00915a721b2
6 files changed
+47
-160
src/ml/ml.cc
+9
-25
@@ -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_training_window / dim->rd->rrdset->update_every;
60
- size_t max_n = Cfg.training_window / dim->rd->rrdset->update_every;
59
+ size_t min_n = Cfg.min_train_samples;
60
+ size_t max_n = Cfg.max_train_samples;
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 - Cfg.training_window, // Fixed time window
65
+ training_response.query_before_t - static_cast<time_t>((max_n - 1) * dim->rd->rrdset->update_every),
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.train_every));
382
+ rc = sqlite3_bind_int64(res, ++param, now_realtime_sec() - (Cfg.num_models_to_use * Cfg.max_train_samples));
383
if (unlikely(rc != SQLITE_OK))
384
goto bind_fail;
385
@@ -674,29 +674,13 @@ 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
-
677
ml_features_t features = {
680
- Cfg.diff_n, smoothing_window, Cfg.lag_n,
678
+ Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
679
worker->scratch_training_cns, training_response.total_values,
680
worker->training_cns, training_response.total_values,
681
worker->training_samples
682
};
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);
683
+ ml_features_preprocess(&features);
684
685
ml_kmeans_init(&dim->kmeans);
686
ml_kmeans_train(&dim->kmeans, &features, Cfg.max_kmeans_iters, training_response.query_after_t, training_response.query_before_t);
@@ -722,7 +706,7 @@ ml_dimension_predict(ml_dimension_t *dim, calculated_number_t value, bool exists
706
}
707
708
// Save the value and return if we don't have enough values for a sample
725
- unsigned n = Cfg.diff_n + Cfg.max_samples_to_smooth + Cfg.lag_n;
709
+ unsigned n = Cfg.diff_n + Cfg.smooth_n + Cfg.lag_n;
710
if (dim->cns.size() < n) {
711
dim->cns.push_back(value);
712
return false;
@@ -747,11 +731,11 @@ ml_dimension_predict(ml_dimension_t *dim, calculated_number_t value, bool exists
731
memcpy(dst_cns, dim->cns.data(), n * sizeof(calculated_number_t));
732
733
ml_features_t features = {
750
- Cfg.diff_n, Cfg.max_samples_to_smooth, Cfg.lag_n,
734
+ Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
735
dst_cns, n, src_cns, n,
736
dim->feature
737
};
754
- ml_features_preprocess(&features, 1.0);
738
+ ml_features_preprocess(&features);
739
740
/*
741
* Lock to predict
src/ml/ml_config.cc
+17
-115
@@ -2,102 +2,6 @@
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
-
5
/*
6
* Global configuration instance to be shared between training and
7
* prediction threads.
@@ -115,27 +19,24 @@ static T clamp(const T& Value, const T& Min, const T& Max) {
19
void ml_config_load(ml_config_t *cfg) {
20
const char *config_section_ml = CONFIG_SECTION_ML;
21
118
- // Migrate old configuration if needed
119
- ml_config_migrate();
120
-
22
int enable_anomaly_detection = inicfg_get_boolean_ondemand(&netdata_config, config_section_ml, "enabled", nd_profile.ml_enabled);
23
24
/*
25
* Read values
26
*/
27
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);
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);
30
unsigned train_every = inicfg_get_duration_seconds(&netdata_config, config_section_ml, "train every", 3 * 3600);
31
32
unsigned num_models_to_use = inicfg_get_number(&netdata_config, config_section_ml, "number of models per dimension", 18);
33
unsigned delete_models_older_than = inicfg_get_duration_seconds(&netdata_config, config_section_ml, "delete models older than", 60 * 60 * 24 * 7);
34
35
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);
37
unsigned lag_n = inicfg_get_number(&netdata_config, config_section_ml, "num samples to lag", 5);
38
39
+ double random_sampling_ratio = inicfg_get_double(&netdata_config, config_section_ml, "random sampling ratio", 1.0 / 5.0 /* default lag_n */);
40
unsigned max_kmeans_iters = inicfg_get_number(&netdata_config, config_section_ml, "maximum number of k-means iterations", 1000);
41
42
double dimension_anomaly_rate_threshold = inicfg_get_double(&netdata_config, config_section_ml, "dimension anomaly score threshold", 0.99);
@@ -163,17 +64,18 @@ void ml_config_load(ml_config_t *cfg) {
64
* Clamp
65
*/
66
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);
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);
69
train_every = clamp<unsigned>(train_every, 1 * 3600, 6 * 3600);
70
71
num_models_to_use = clamp<unsigned>(num_models_to_use, 1, 7 * 24);
72
delete_models_older_than = clamp<unsigned>(delete_models_older_than, 60 * 60 * 24 * 1, 60 * 60 * 24 * 7);
73
74
diff_n = clamp(diff_n, 0u, 1u);
174
- max_samples_to_smooth = clamp<size_t>(max_samples_to_smooth, 0, 5);
75
+ smooth_n = clamp(smooth_n, 0u, 5u);
76
lag_n = clamp(lag_n, 1u, 5u);
77
78
+ random_sampling_ratio = clamp(random_sampling_ratio, 0.2, 1.0);
79
max_kmeans_iters = clamp(max_kmeans_iters, 500u, 1000u);
80
81
dimension_anomaly_rate_threshold = clamp(dimension_anomaly_rate_threshold, 0.01, 5.00);
@@ -184,18 +86,18 @@ void ml_config_load(ml_config_t *cfg) {
86
num_worker_threads = clamp<size_t>(num_worker_threads, 4, netdata_conf_cpus());
87
flush_models_batch_size = clamp<size_t>(flush_models_batch_size, 8, 512);
88
187
- suppression_window = clamp<size_t>(suppression_window, 1, training_window);
89
+ suppression_window = clamp<size_t>(suppression_window, 1, max_train_samples);
90
suppression_threshold = clamp<size_t>(suppression_threshold, 1, suppression_window);
91
92
/*
93
* Validate
94
*/
95
194
- if (min_training_window >= training_window) {
195
- netdata_log_error("invalid min/max training window found (%ld >= %ld)", min_training_window, training_window);
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);
98
197
- min_training_window = 1 * 3600;
198
- training_window = 6 * 3600;
99
+ min_train_samples = 1 * 3600;
100
+ max_train_samples = 6 * 3600;
101
}
102
103
/*
@@ -204,18 +106,18 @@ void ml_config_load(ml_config_t *cfg) {
106
107
cfg->enable_anomaly_detection = enable_anomaly_detection;
108
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;
109
+ cfg->max_train_samples = max_train_samples;
110
+ cfg->min_train_samples = min_train_samples;
111
cfg->train_every = train_every;
112
113
cfg->num_models_to_use = num_models_to_use;
114
cfg->delete_models_older_than = delete_models_older_than;
115
116
cfg->diff_n = diff_n;
117
+ cfg->smooth_n = smooth_n;
118
cfg->lag_n = lag_n;
119
120
+ cfg->random_sampling_ratio = random_sampling_ratio;
121
cfg->max_kmeans_iters = max_kmeans_iters;
122
123
cfg->host_anomaly_rate_threshold = host_anomaly_rate_threshold;
src/ml/ml_config.h
+5
-4
@@ -8,10 +8,8 @@
8
typedef struct {
9
int enable_anomaly_detection;
10
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)
11
+ unsigned max_train_samples;
12
+ unsigned min_train_samples;
13
unsigned train_every;
14
15
unsigned num_models_to_use;
@@ -20,7 +18,10 @@ typedef struct {
18
unsigned db_engine_anomaly_rate_every;
19
20
unsigned diff_n;
21
+ unsigned smooth_n;
22
unsigned lag_n;
23
+
24
+ double random_sampling_ratio;
25
unsigned max_kmeans_iters;
26
27
double dimension_anomaly_score_threshold;
src/ml/ml_features.cc
+7
-4
@@ -41,11 +41,14 @@ 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, double sampling_ratio)
44
+static void ml_features_lag(ml_features_t *features)
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
+
52
uint32_t max_mt = std::numeric_limits<uint32_t>::max();
53
uint32_t cutoff = static_cast<double>(max_mt) * sampling_ratio;
54
@@ -55,7 +58,7 @@ static void ml_features_lag(ml_features_t *features, double sampling_ratio)
58
DSample &DS = features->preprocessed_features[sample_idx++];
59
DS.set_size(features->lag_n);
60
58
- if (Cfg.random_nums[idx % Cfg.random_nums.size()] > cutoff) {
61
+ if (Cfg.random_nums[idx] > cutoff) {
62
sample_idx--;
63
continue;
64
}
@@ -67,9 +70,9 @@ static void ml_features_lag(ml_features_t *features, double sampling_ratio)
70
features->preprocessed_features.resize(sample_idx);
71
}
72
70
-void ml_features_preprocess(ml_features_t *features, double sampling_ratio)
73
+void ml_features_preprocess(ml_features_t *features)
74
{
75
ml_features_diff(features);
76
ml_features_smooth(features);
74
- ml_features_lag(features, sampling_ratio);
77
+ ml_features_lag(features);
78
}
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, double sampling_ratio);
24
+void ml_features_preprocess(ml_features_t *features);
25
26
#endif /* ML_FEATURES_H */
src/ml/ml_public.cc
+8
-11
@@ -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, "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);
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);
146
buffer_json_member_add_uint64(wb, "train-every", Cfg.train_every);
147
148
buffer_json_member_add_uint64(wb, "diff-n", Cfg.diff_n);
149
+ buffer_json_member_add_uint64(wb, "smooth-n", Cfg.smooth_n);
150
buffer_json_member_add_uint64(wb, "lag-n", Cfg.lag_n);
151
153
- buffer_json_member_add_uint64(wb, "max-kmeans-iters", Cfg.max_kmeans_iters);
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);
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_training_vectors);
373
- for (size_t Idx = 0; Idx != Cfg.max_training_vectors; Idx++)
372
+ Cfg.random_nums.reserve(Cfg.max_train_samples);
373
+ for (size_t Idx = 0; Idx != Cfg.max_train_samples; Idx++)
374
Cfg.random_nums.push_back(Gen());
375
376
// init training thread-specific data
@@ -378,10 +378,7 @@ 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
- // 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);
381
+ size_t max_elements_needed_for_training = (size_t) Cfg.max_train_samples * (size_t) (Cfg.lag_n + 1);
382
worker->training_cns = new calculated_number_t[max_elements_needed_for_training]();
383
worker->scratch_training_cns = new calculated_number_t[max_elements_needed_for_training]();
384