@cryptotaxi247 / netdata-1 / commits / 556bdad9b

Revert ML changes. (#14908)

vkalintiris committed Apr 14, 2023 at 10:49 UTC 556bdad9be687a917901fa77e0d2ffeb6d0b4a47
11 files changed +431 -774
daemon/global_statistics.c
+27 -1
@@ -827,7 +827,33 @@ static void global_statistics_charts(void) {
827 rrdset_done(st_points_stored);
828 }
829
830 - ml_update_global_statistics_charts(gs.ml_models_consulted);
830 + {
831 + static RRDSET *st = NULL;
832 + static RRDDIM *rd = NULL;
833 +
834 + if (unlikely(!st)) {
835 + st = rrdset_create_localhost(
836 + "netdata" // type
837 + , "ml_models_consulted" // id
838 + , NULL // name
839 + , NETDATA_ML_CHART_FAMILY // family
840 + , NULL // context
841 + , "KMeans models used for prediction" // title
842 + , "models" // units
843 + , NETDATA_ML_PLUGIN // plugin
844 + , NETDATA_ML_MODULE_DETECTION // module
845 + , NETDATA_ML_CHART_PRIO_MACHINE_LEARNING_STATUS // priority
846 + , localhost->rrd_update_every // update_every
847 + , RRDSET_TYPE_AREA // chart_type
848 + );
849 +
850 + rd = rrddim_add(st, "num_models_consulted", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
851 + }
852 +
853 + rrddim_set_by_pointer(st, rd, (collected_number) gs.ml_models_consulted);
854 +
855 + rrdset_done(st);
856 + }
857 }
858
859 // ----------------------------------------------------------------------------
daemon/main.c
+9 -9
@@ -148,6 +148,10 @@ static void service_to_buffer(BUFFER *wb, SERVICE_TYPE service) {
148 buffer_strcat(wb, "MAINTENANCE ");
149 if(service & SERVICE_COLLECTORS)
150 buffer_strcat(wb, "COLLECTORS ");
151 + if(service & SERVICE_ML_TRAINING)
152 + buffer_strcat(wb, "ML_TRAINING ");
153 + if(service & SERVICE_ML_PREDICTION)
154 + buffer_strcat(wb, "ML_PREDICTION ");
155 if(service & SERVICE_REPLICATION)
156 buffer_strcat(wb, "REPLICATION ");
157 if(service & ABILITY_DATA_QUERIES)
@@ -336,11 +340,6 @@ void netdata_cleanup_and_exit(int ret) {
340 }
341 #endif
342
339 - delta_shutdown_time("disable ML detection and training threads");
340 -
341 - ml_stop_threads();
342 - ml_fini();
343 -
343 delta_shutdown_time("disable maintenance, new queries, new web requests, new streaming connections and aclk");
344
345 service_signal_exit(
@@ -352,11 +351,12 @@ void netdata_cleanup_and_exit(int ret) {
351 | SERVICE_ACLKSYNC
352 );
353
355 - delta_shutdown_time("stop replication, exporters, health and web servers threads");
354 + delta_shutdown_time("stop replication, exporters, ML training, health and web servers threads");
355
356 timeout = !service_wait_exit(
357 SERVICE_REPLICATION
358 | SERVICE_EXPORTERS
359 + | SERVICE_ML_TRAINING
360 | SERVICE_HEALTH
361 | SERVICE_WEB_SERVER
362 , 3 * USEC_PER_SEC);
@@ -368,10 +368,11 @@ void netdata_cleanup_and_exit(int ret) {
368 | SERVICE_STREAMING
369 , 3 * USEC_PER_SEC);
370
371 - delta_shutdown_time("stop context thread");
371 + delta_shutdown_time("stop ML prediction and context threads");
372
373 timeout = !service_wait_exit(
374 - SERVICE_CONTEXT
374 + SERVICE_ML_PREDICTION
375 + | SERVICE_CONTEXT
376 , 3 * USEC_PER_SEC);
377
378 delta_shutdown_time("stop maintenance thread");
@@ -2084,7 +2085,6 @@ int main(int argc, char **argv) {
2085 }
2086 else debug(D_SYSTEM, "Not starting thread %s.", st->name);
2087 }
2087 - ml_start_threads();
2088
2089 // ------------------------------------------------------------------------
2090 // Initialize netdata agent command serving from cli and signals
daemon/main.h
+11 -9
@@ -33,15 +33,17 @@ typedef enum {
33 ABILITY_STREAMING_CONNECTIONS = (1 << 2),
34 SERVICE_MAINTENANCE = (1 << 3),
35 SERVICE_COLLECTORS = (1 << 4),
36 - SERVICE_REPLICATION = (1 << 5),
37 - SERVICE_WEB_SERVER = (1 << 6),
38 - SERVICE_ACLK = (1 << 7),
39 - SERVICE_HEALTH = (1 << 8),
40 - SERVICE_STREAMING = (1 << 9),
41 - SERVICE_CONTEXT = (1 << 10),
42 - SERVICE_ANALYTICS = (1 << 11),
43 - SERVICE_EXPORTERS = (1 << 12),
44 - SERVICE_ACLKSYNC = (1 << 13)
36 + SERVICE_ML_TRAINING = (1 << 5),
37 + SERVICE_ML_PREDICTION = (1 << 6),
38 + SERVICE_REPLICATION = (1 << 7),
39 + SERVICE_WEB_SERVER = (1 << 8),
40 + SERVICE_ACLK = (1 << 9),
41 + SERVICE_HEALTH = (1 << 10),
42 + SERVICE_STREAMING = (1 << 11),
43 + SERVICE_CONTEXT = (1 << 12),
44 + SERVICE_ANALYTICS = (1 << 13),
45 + SERVICE_EXPORTERS = (1 << 14),
46 + SERVICE_ACLKSYNC = (1 << 15)
47 } SERVICE_TYPE;
48
49 typedef enum {
database/rrdhost.c
+3
@@ -524,6 +524,7 @@ int is_legacy = 1;
524 rrdhost_load_rrdcontext_data(host);
525 // rrdhost_flag_set(host, RRDHOST_FLAG_METADATA_INFO | RRDHOST_FLAG_METADATA_UPDATE);
526 ml_host_new(host);
527 + ml_host_start_training_thread(host);
528 } else
529 rrdhost_flag_set(host, RRDHOST_FLAG_PENDING_CONTEXT_LOAD | RRDHOST_FLAG_ARCHIVED | RRDHOST_FLAG_ORPHAN);
530
@@ -640,6 +641,7 @@ static void rrdhost_update(RRDHOST *host
641 host->rrdpush_replication_step = rrdpush_replication_step;
642
643 ml_host_new(host);
644 + ml_host_start_training_thread(host);
645
646 rrdhost_load_rrdcontext_data(host);
647 info("Host %s is not in archived mode anymore", rrdhost_hostname(host));
@@ -1141,6 +1143,7 @@ void rrdhost_free___while_having_rrd_wrlock(RRDHOST *host, bool force) {
1143 rrdcalctemplate_index_destroy(host);
1144
1145 // cleanup ML resources
1146 + ml_host_stop_training_thread(host);
1147 ml_host_delete(host);
1148
1149 freez(host->exporting_flags);
ml/Config.cc
+1 -11
@@ -34,7 +34,7 @@ void ml_config_load(ml_config_t *cfg) {
34 unsigned smooth_n = config_get_number(config_section_ml, "num samples to smooth", 3);
35 unsigned lag_n = config_get_number(config_section_ml, "num samples to lag", 5);
36
37 - double random_sampling_ratio = config_get_float(config_section_ml, "random sampling ratio", 1.0 / 5.0 /* default lag_n */);
37 + double random_sampling_ratio = config_get_float(config_section_ml, "random sampling ratio", 1.0 / lag_n);
38 unsigned max_kmeans_iters = config_get_number(config_section_ml, "maximum number of k-means iterations", 1000);
39
40 double dimension_anomaly_rate_threshold = config_get_float(config_section_ml, "dimension anomaly score threshold", 0.99);
@@ -43,10 +43,6 @@ void ml_config_load(ml_config_t *cfg) {
43 std::string anomaly_detection_grouping_method = config_get(config_section_ml, "anomaly detection grouping method", "average");
44 time_t anomaly_detection_query_duration = config_get_number(config_section_ml, "anomaly detection grouping duration", 5 * 60);
45
46 - size_t num_training_threads = config_get_number(config_section_ml, "num training threads", 4);
47 -
48 - bool enable_statistics_charts = config_get_boolean(config_section_ml, "enable statistics charts", false);
49 -
46 /*
47 * Clamp
48 */
@@ -68,8 +64,6 @@ void ml_config_load(ml_config_t *cfg) {
64 host_anomaly_rate_threshold = clamp(host_anomaly_rate_threshold, 0.1, 10.0);
65 anomaly_detection_query_duration = clamp<time_t>(anomaly_detection_query_duration, 60, 15 * 60);
66
71 - num_training_threads = clamp<size_t>(num_training_threads, 1, 128);
72 -
67 /*
68 * Validate
69 */
@@ -115,8 +109,4 @@ void ml_config_load(ml_config_t *cfg) {
109 cfg->sp_charts_to_skip = simple_pattern_create(cfg->charts_to_skip.c_str(), NULL, SIMPLE_PATTERN_EXACT, true);
110
111 cfg->stream_anomaly_detection_charts = config_get_boolean(config_section_ml, "stream anomaly detection charts", true);
118 -
119 - cfg->num_training_threads = num_training_threads;
120 -
121 - cfg->enable_statistics_charts = enable_statistics_charts;
112 }
ml/ad_charts.cc
+69 -98
@@ -6,7 +6,7 @@ void ml_update_dimensions_chart(ml_host_t *host, const ml_machine_learning_stats
6 /*
7 * Machine learning status
8 */
9 - if (Cfg.enable_statistics_charts) {
9 + {
10 if (!host->machine_learning_status_rs) {
11 char id_buf[1024];
12 char name_buf[1024];
@@ -48,7 +48,7 @@ void ml_update_dimensions_chart(ml_host_t *host, const ml_machine_learning_stats
48 /*
49 * Metric type
50 */
51 - if (Cfg.enable_statistics_charts) {
51 + {
52 if (!host->metric_type_rs) {
53 char id_buf[1024];
54 char name_buf[1024];
@@ -90,7 +90,7 @@ void ml_update_dimensions_chart(ml_host_t *host, const ml_machine_learning_stats
90 /*
91 * Training status
92 */
93 - if (Cfg.enable_statistics_charts) {
93 + {
94 if (!host->training_status_rs) {
95 char id_buf[1024];
96 char name_buf[1024];
@@ -179,6 +179,7 @@ void ml_update_dimensions_chart(ml_host_t *host, const ml_machine_learning_stats
179
180 rrdset_done(host->dimensions_rs);
181 }
182 +
183 }
184
185 void ml_update_host_and_detection_rate_charts(ml_host_t *host, collected_number AnomalyRate) {
@@ -300,20 +301,20 @@ void ml_update_host_and_detection_rate_charts(ml_host_t *host, collected_number
301 }
302 }
303
303 -void ml_update_training_statistics_chart(ml_training_thread_t *training_thread, const ml_training_stats_t &ts) {
304 +void ml_update_training_statistics_chart(ml_host_t *host, const ml_training_stats_t &ts) {
305 /*
306 * queue stats
307 */
308 {
308 - if (!training_thread->queue_stats_rs) {
309 + if (!host->queue_stats_rs) {
310 char id_buf[1024];
311 char name_buf[1024];
312
312 - snprintfz(id_buf, 1024, "training_queue_%zu_stats", training_thread->id);
313 - snprintfz(name_buf, 1024, "training_queue_%zu_stats", training_thread->id);
313 + snprintfz(id_buf, 1024, "queue_stats_on_%s", localhost->machine_guid);
314 + snprintfz(name_buf, 1024, "queue_stats_on_%s", rrdhost_hostname(localhost));
315
315 - training_thread->queue_stats_rs = rrdset_create(
316 - localhost,
316 + host->queue_stats_rs = rrdset_create(
317 + host->rh,
318 "netdata", // type
319 id_buf, // id
320 name_buf, // name
@@ -327,35 +328,35 @@ void ml_update_training_statistics_chart(ml_training_thread_t *training_thread,
328 localhost->rrd_update_every, // update_every
329 RRDSET_TYPE_LINE// chart_type
330 );
330 - rrdset_flag_set(training_thread->queue_stats_rs, RRDSET_FLAG_ANOMALY_DETECTION);
331 + rrdset_flag_set(host->queue_stats_rs, RRDSET_FLAG_ANOMALY_DETECTION);
332
332 - training_thread->queue_stats_queue_size_rd =
333 - rrddim_add(training_thread->queue_stats_rs, "queue_size", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
334 - training_thread->queue_stats_popped_items_rd =
335 - rrddim_add(training_thread->queue_stats_rs, "popped_items", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
333 + host->queue_stats_queue_size_rd =
334 + rrddim_add(host->queue_stats_rs, "queue_size", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
335 + host->queue_stats_popped_items_rd =
336 + rrddim_add(host->queue_stats_rs, "popped_items", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
337 }
338
338 - rrddim_set_by_pointer(training_thread->queue_stats_rs,
339 - training_thread->queue_stats_queue_size_rd, ts.queue_size);
340 - rrddim_set_by_pointer(training_thread->queue_stats_rs,
341 - training_thread->queue_stats_popped_items_rd, ts.num_popped_items);
339 + rrddim_set_by_pointer(host->queue_stats_rs,
340 + host->queue_stats_queue_size_rd, ts.queue_size);
341 + rrddim_set_by_pointer(host->queue_stats_rs,
342 + host->queue_stats_popped_items_rd, ts.num_popped_items);
343
343 - rrdset_done(training_thread->queue_stats_rs);
344 + rrdset_done(host->queue_stats_rs);
345 }
346
347 /*
348 * training stats
349 */
350 {
350 - if (!training_thread->training_time_stats_rs) {
351 + if (!host->training_time_stats_rs) {
352 char id_buf[1024];
353 char name_buf[1024];
354
354 - snprintfz(id_buf, 1024, "training_queue_%zu_time_stats", training_thread->id);
355 - snprintfz(name_buf, 1024, "training_queue_%zu_time_stats", training_thread->id);
355 + snprintfz(id_buf, 1024, "training_time_stats_on_%s", localhost->machine_guid);
356 + snprintfz(name_buf, 1024, "training_time_stats_on_%s", rrdhost_hostname(localhost));
357
357 - training_thread->training_time_stats_rs = rrdset_create(
358 - localhost,
358 + host->training_time_stats_rs = rrdset_create(
359 + host->rh,
360 "netdata", // type
361 id_buf, // id
362 name_buf, // name
@@ -369,39 +370,39 @@ void ml_update_training_statistics_chart(ml_training_thread_t *training_thread,
370 localhost->rrd_update_every, // update_every
371 RRDSET_TYPE_LINE// chart_type
372 );
372 - rrdset_flag_set(training_thread->training_time_stats_rs, RRDSET_FLAG_ANOMALY_DETECTION);
373 -
374 - training_thread->training_time_stats_allotted_rd =
375 - rrddim_add(training_thread->training_time_stats_rs, "allotted", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
376 - training_thread->training_time_stats_consumed_rd =
377 - rrddim_add(training_thread->training_time_stats_rs, "consumed", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
378 - training_thread->training_time_stats_remaining_rd =
379 - rrddim_add(training_thread->training_time_stats_rs, "remaining", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
373 + rrdset_flag_set(host->training_time_stats_rs, RRDSET_FLAG_ANOMALY_DETECTION);
374 +
375 + host->training_time_stats_allotted_rd =
376 + rrddim_add(host->training_time_stats_rs, "allotted", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
377 + host->training_time_stats_consumed_rd =
378 + rrddim_add(host->training_time_stats_rs, "consumed", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
379 + host->training_time_stats_remaining_rd =
380 + rrddim_add(host->training_time_stats_rs, "remaining", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
381 }
382
382 - rrddim_set_by_pointer(training_thread->training_time_stats_rs,
383 - training_thread->training_time_stats_allotted_rd, ts.allotted_ut);
384 - rrddim_set_by_pointer(training_thread->training_time_stats_rs,
385 - training_thread->training_time_stats_consumed_rd, ts.consumed_ut);
386 - rrddim_set_by_pointer(training_thread->training_time_stats_rs,
387 - training_thread->training_time_stats_remaining_rd, ts.remaining_ut);
383 + rrddim_set_by_pointer(host->training_time_stats_rs,
384 + host->training_time_stats_allotted_rd, ts.allotted_ut);
385 + rrddim_set_by_pointer(host->training_time_stats_rs,
386 + host->training_time_stats_consumed_rd, ts.consumed_ut);
387 + rrddim_set_by_pointer(host->training_time_stats_rs,
388 + host->training_time_stats_remaining_rd, ts.remaining_ut);
389
389 - rrdset_done(training_thread->training_time_stats_rs);
390 + rrdset_done(host->training_time_stats_rs);
391 }
392
393 /*
394 * training result stats
395 */
396 {
396 - if (!training_thread->training_results_rs) {
397 + if (!host->training_results_rs) {
398 char id_buf[1024];
399 char name_buf[1024];
400
400 - snprintfz(id_buf, 1024, "training_queue_%zu_results", training_thread->id);
401 - snprintfz(name_buf, 1024, "training_queue_%zu_results", training_thread->id);
401 + snprintfz(id_buf, 1024, "training_results_on_%s", localhost->machine_guid);
402 + snprintfz(name_buf, 1024, "training_results_on_%s", rrdhost_hostname(localhost));
403
403 - training_thread->training_results_rs = rrdset_create(
404 - localhost,
404 + host->training_results_rs = rrdset_create(
405 + host->rh,
406 "netdata", // type
407 id_buf, // id
408 name_buf, // name
@@ -415,61 +416,31 @@ void ml_update_training_statistics_chart(ml_training_thread_t *training_thread,
416 localhost->rrd_update_every, // update_every
417 RRDSET_TYPE_LINE// chart_type
418 );
418 - rrdset_flag_set(training_thread->training_results_rs, RRDSET_FLAG_ANOMALY_DETECTION);
419 -
420 - training_thread->training_results_ok_rd =
421 - rrddim_add(training_thread->training_results_rs, "ok", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
422 - training_thread->training_results_invalid_query_time_range_rd =
423 - rrddim_add(training_thread->training_results_rs, "invalid-queries", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
424 - training_thread->training_results_not_enough_collected_values_rd =
425 - rrddim_add(training_thread->training_results_rs, "not-enough-values", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
426 - training_thread->training_results_null_acquired_dimension_rd =
427 - rrddim_add(training_thread->training_results_rs, "null-acquired-dimensions", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
428 - training_thread->training_results_chart_under_replication_rd =
429 - rrddim_add(training_thread->training_results_rs, "chart-under-replication", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
419 + rrdset_flag_set(host->training_results_rs, RRDSET_FLAG_ANOMALY_DETECTION);
420 +
421 + host->training_results_ok_rd =
422 + rrddim_add(host->training_results_rs, "ok", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
423 + host->training_results_invalid_query_time_range_rd =
424 + rrddim_add(host->training_results_rs, "invalid-queries", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
425 + host->training_results_not_enough_collected_values_rd =
426 + rrddim_add(host->training_results_rs, "not-enough-values", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
427 + host->training_results_null_acquired_dimension_rd =
428 + rrddim_add(host->training_results_rs, "null-acquired-dimensions", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
429 + host->training_results_chart_under_replication_rd =
430 + rrddim_add(host->training_results_rs, "chart-under-replication", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
431 }
432
432 - rrddim_set_by_pointer(training_thread->training_results_rs,
433 - training_thread->training_results_ok_rd, ts.training_result_ok);
434 - rrddim_set_by_pointer(training_thread->training_results_rs,
435 - training_thread->training_results_invalid_query_time_range_rd, ts.training_result_invalid_query_time_range);
436 - rrddim_set_by_pointer(training_thread->training_results_rs,
437 - training_thread->training_results_not_enough_collected_values_rd, ts.training_result_not_enough_collected_values);
438 - rrddim_set_by_pointer(training_thread->training_results_rs,
439 - training_thread->training_results_null_acquired_dimension_rd, ts.training_result_null_acquired_dimension);
440 - rrddim_set_by_pointer(training_thread->training_results_rs,
441 - training_thread->training_results_chart_under_replication_rd, ts.training_result_chart_under_replication);
442 -
443 - rrdset_done(training_thread->training_results_rs);
444 - }
445 -}
446 -
447 -void ml_update_global_statistics_charts(uint64_t models_consulted) {
448 - if (Cfg.enable_statistics_charts) {
449 - static RRDSET *st = NULL;
450 - static RRDDIM *rd = NULL;
451 -
452 - if (unlikely(!st)) {
453 - st = rrdset_create_localhost(
454 - "netdata" // type
455 - , "ml_models_consulted" // id
456 - , NULL // name
457 - , NETDATA_ML_CHART_FAMILY // family
458 - , NULL // context
459 - , "KMeans models used for prediction" // title
460 - , "models" // units
461 - , NETDATA_ML_PLUGIN // plugin
462 - , NETDATA_ML_MODULE_DETECTION // module
463 - , NETDATA_ML_CHART_PRIO_MACHINE_LEARNING_STATUS // priority
464 - , localhost->rrd_update_every // update_every
465 - , RRDSET_TYPE_AREA // chart_type
466 - );
467 -
468 - rd = rrddim_add(st, "num_models_consulted", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
469 - }
470 -
471 - rrddim_set_by_pointer(st, rd, (collected_number) models_consulted);
472 -
473 - rrdset_done(st);
433 + rrddim_set_by_pointer(host->training_results_rs,
434 + host->training_results_ok_rd, ts.training_result_ok);
435 + rrddim_set_by_pointer(host->training_results_rs,
436 + host->training_results_invalid_query_time_range_rd, ts.training_result_invalid_query_time_range);
437 + rrddim_set_by_pointer(host->training_results_rs,
438 + host->training_results_not_enough_collected_values_rd, ts.training_result_not_enough_collected_values);
439 + rrddim_set_by_pointer(host->training_results_rs,
440 + host->training_results_null_acquired_dimension_rd, ts.training_result_null_acquired_dimension);
441 + rrddim_set_by_pointer(host->training_results_rs,
442 + host->training_results_chart_under_replication_rd, ts.training_result_chart_under_replication);
443 +
444 + rrdset_done(host->training_results_rs);
445 }
446 }
ml/ad_charts.h
+1 -1
@@ -9,6 +9,6 @@ void ml_update_dimensions_chart(ml_host_t *host, const ml_machine_learning_stats
9
10 void ml_update_host_and_detection_rate_charts(ml_host_t *host, collected_number anomaly_rate);
11
12 -void ml_update_training_statistics_chart(ml_training_thread_t *training_thread, const ml_training_stats_t &ts);
12 +void ml_update_training_statistics_chart(ml_host_t *host, const ml_training_stats_t &ts);
13
14 #endif /* ML_ADCHARTS_H */
ml/ml-dummy.c
-10
@@ -19,12 +19,6 @@ bool ml_streaming_enabled() {
19
20 void ml_init(void) {}
21
22 -void ml_fini(void) {}
23 -
24 -void ml_start_threads(void) {}
25 -
26 -void ml_stop_threads(void) {}
27 -
22 void ml_host_new(RRDHOST *rh) {
23 UNUSED(rh);
24 }
@@ -92,8 +86,4 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
86 return false;
87 }
88
95 -void ml_update_global_statistics_charts(uint64_t models_consulted) {
96 - UNUSED(models_consulted);
97 -}
98 -
89 #endif
ml/ml-private.h
+13 -29
@@ -33,15 +33,14 @@ typedef struct {
33 /*
34 * KMeans
35 */
36 -
36 typedef struct {
37 + size_t num_clusters;
38 + size_t max_iterations;
39 +
40 std::vector<DSample> cluster_centers;
41
42 calculated_number_t min_dist;
43 calculated_number_t max_dist;
42 -
43 - uint32_t after;
44 - uint32_t before;
44 } ml_kmeans_t;
45
46 typedef struct machine_learning_stats_t {
@@ -124,7 +123,6 @@ enum ml_training_result {
123
124 typedef struct {
125 // Chart/dimension we want to train
127 - STRING *host_id;
126 STRING *chart_id;
127 STRING *dimension_id;
128
@@ -170,7 +168,6 @@ typedef struct {
168 /*
169 * Queue
170 */
173 -
171 typedef struct {
172 std::queue<ml_training_request_t> internal;
173 netdata_mutex_t mutex;
@@ -178,6 +175,7 @@ typedef struct {
175 std::atomic<bool> exit;
176 } ml_queue_t;
177
178 +
179 typedef struct {
180 RRDDIM *rd;
181
@@ -209,13 +207,20 @@ typedef struct {
207 RRDHOST *rh;
208
209 ml_machine_learning_stats_t mls;
210 + ml_training_stats_t ts;
211
212 calculated_number_t host_anomaly_rate;
213
215 - netdata_mutex_t mutex;
214 + std::atomic<bool> threads_running;
215 + std::atomic<bool> threads_cancelled;
216 + std::atomic<bool> threads_joined;
217
218 ml_queue_t *training_queue;
219
220 + netdata_mutex_t mutex;
221 +
222 + netdata_thread_t training_thread;
223 +
224 /*
225 * bookkeeping for anomaly detection charts
226 */
@@ -244,19 +249,6 @@ typedef struct {
249 RRDSET *detector_events_rs;
250 RRDDIM *detector_events_above_threshold_rd;
251 RRDDIM *detector_events_new_anomaly_event_rd;
247 -} ml_host_t;
248 -
249 -typedef struct {
250 - size_t id;
251 - netdata_thread_t nd_thread;
252 - netdata_mutex_t nd_mutex;
253 -
254 - ml_queue_t *training_queue;
255 - ml_training_stats_t training_stats;
256 -
257 - calculated_number_t *training_cns;
258 - calculated_number_t *scratch_training_cns;
259 - std::vector<DSample> training_samples;
252
253 RRDSET *queue_stats_rs;
254 RRDDIM *queue_stats_queue_size_rd;
@@ -273,7 +265,7 @@ typedef struct {
265 RRDDIM *training_results_not_enough_collected_values_rd;
266 RRDDIM *training_results_null_acquired_dimension_rd;
267 RRDDIM *training_results_chart_under_replication_rd;
276 -} ml_training_thread_t;
268 +} ml_host_t;
269
270 typedef struct {
271 bool enable_anomaly_detection;
@@ -310,14 +302,6 @@ typedef struct {
302 std::vector<uint32_t> random_nums;
303
304 netdata_thread_t detection_thread;
313 - std::atomic<bool> detection_stop;
314 -
315 - size_t num_training_threads;
316 -
317 - std::vector<ml_training_thread_t> training_threads;
318 - std::atomic<bool> training_stop;
319 -
320 - bool enable_statistics_charts;
305 } ml_config_t;
306
307 void ml_config_load(ml_config_t *cfg);
ml/ml.cc
+293 -599
@@ -7,18 +7,15 @@
7 #include <random>
8
9 #include "ad_charts.h"
10 -#include "database/sqlite/sqlite3.h"
10
12 -#define WORKER_TRAIN_QUEUE_POP 0
13 -#define WORKER_TRAIN_ACQUIRE_DIMENSION 1
14 -#define WORKER_TRAIN_QUERY 2
15 -#define WORKER_TRAIN_KMEANS 3
16 -#define WORKER_TRAIN_UPDATE_MODELS 4
17 -#define WORKER_TRAIN_RELEASE_DIMENSION 5
18 -#define WORKER_TRAIN_UPDATE_HOST 6
19 -#define WORKER_TRAIN_LOAD_MODELS 7
11 +typedef struct {
12 + calculated_number_t *training_cns;
13 + calculated_number_t *scratch_training_cns;
14 +
15 + std::vector<DSample> training_samples;
16 +} ml_tls_data_t;
17
21 -static sqlite3 *db = NULL;
18 +static thread_local ml_tls_data_t tls_data;
19
20 /*
21 * Functions to convert enums to strings
@@ -176,26 +173,26 @@ ml_features_preprocess(ml_features_t *features)
173 */
174
175 static void
179 -ml_kmeans_init(ml_kmeans_t *kmeans)
176 +ml_kmeans_init(ml_kmeans_t *kmeans, size_t num_clusters, size_t max_iterations)
177 {
181 - kmeans->cluster_centers.reserve(2);
178 + kmeans->num_clusters = num_clusters;
179 + kmeans->max_iterations = max_iterations;
180 +
181 + kmeans->cluster_centers.reserve(kmeans->num_clusters);
182 kmeans->min_dist = std::numeric_limits<calculated_number_t>::max();
183 kmeans->max_dist = std::numeric_limits<calculated_number_t>::min();
184 }
185
186 static void
187 -ml_kmeans_train(ml_kmeans_t *kmeans, const ml_features_t *features, time_t after, time_t before)
187 +ml_kmeans_train(ml_kmeans_t *kmeans, const ml_features_t *features)
188 {
189 - kmeans->after = (uint32_t) after;
190 - kmeans->before = (uint32_t) before;
191 -
189 kmeans->min_dist = std::numeric_limits<calculated_number_t>::max();
190 kmeans->max_dist = std::numeric_limits<calculated_number_t>::min();
191
192 kmeans->cluster_centers.clear();
193
197 - dlib::pick_initial_centers(2, kmeans->cluster_centers, features->preprocessed_features);
198 - dlib::find_clusters_using_kmeans(features->preprocessed_features, kmeans->cluster_centers, Cfg.max_kmeans_iters);
194 + dlib::pick_initial_centers(kmeans->num_clusters, kmeans->cluster_centers, features->preprocessed_features);
195 + dlib::find_clusters_using_kmeans(features->preprocessed_features, kmeans->cluster_centers, kmeans->max_iterations);
196
197 for (const auto &preprocessed_feature : features->preprocessed_features) {
198 calculated_number_t mean_dist = 0.0;
@@ -204,7 +201,7 @@ ml_kmeans_train(ml_kmeans_t *kmeans, const ml_features_t *features, time_t after
201 mean_dist += dlib::length(cluster_center - preprocessed_feature);
202 }
203
207 - mean_dist /= kmeans->cluster_centers.size();
204 + mean_dist /= kmeans->num_clusters;
205
206 if (mean_dist < kmeans->min_dist)
207 kmeans->min_dist = mean_dist;
@@ -221,7 +218,7 @@ ml_kmeans_anomaly_score(const ml_kmeans_t *kmeans, const DSample &DS)
218 for (const auto &CC: kmeans->cluster_centers)
219 mean_dist += dlib::length(CC - DS);
220
224 - mean_dist /= kmeans->cluster_centers.size();
221 + mean_dist /= kmeans->num_clusters;
222
223 if (kmeans->max_dist == kmeans->min_dist)
224 return 0.0;
@@ -267,14 +264,7 @@ ml_queue_pop(ml_queue_t *q)
264 {
265 netdata_mutex_lock(&q->mutex);
266
270 - ml_training_request_t req = {
271 - NULL, // host_id
272 - NULL, // chart id
273 - NULL, // dimension id
274 - 0, // current time
275 - 0, // first entry
276 - 0 // last entry
277 - };
267 + ml_training_request_t req = { NULL, NULL, 0, 0, 0 };
268
269 while (q->internal.empty()) {
270 pthread_cond_wait(&q->cond_var, &q->mutex);
@@ -317,7 +307,7 @@ ml_queue_signal(ml_queue_t *q)
307 */
308
309 static std::pair<calculated_number_t *, ml_training_response_t>
320 -ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
310 +ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t &training_request)
311 {
312 ml_training_response_t training_response = {};
313
@@ -358,7 +348,7 @@ ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimens
348 STORAGE_PRIORITY_BEST_EFFORT);
349
350 size_t idx = 0;
361 - memset(training_thread->training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
351 + memset(tls_data.training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
352 calculated_number_t last_value = std::numeric_limits<calculated_number_t>::quiet_NaN();
353
354 while (!storage_engine_query_is_finished(&handle)) {
@@ -375,11 +365,11 @@ ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimens
365 training_response.db_after_t = timestamp;
366 training_response.db_before_t = timestamp;
367
378 - training_thread->training_cns[idx] = value;
379 - last_value = training_thread->training_cns[idx];
368 + tls_data.training_cns[idx] = value;
369 + last_value = tls_data.training_cns[idx];
370 training_response.collected_values++;
371 } else
382 - training_thread->training_cns[idx] = last_value;
372 + tls_data.training_cns[idx] = last_value;
373
374 idx++;
375 }
@@ -394,270 +384,20 @@ ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimens
384 }
385
386 // Find first non-NaN value.
397 - for (idx = 0; std::isnan(training_thread->training_cns[idx]); idx++, training_response.total_values--) { }
387 + for (idx = 0; std::isnan(tls_data.training_cns[idx]); idx++, training_response.total_values--) { }
388
389 // Overwrite NaN values.
390 if (idx != 0)
401 - memmove(training_thread->training_cns, &training_thread->training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
391 + memmove(tls_data.training_cns, &tls_data.training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
392
393 training_response.result = TRAINING_RESULT_OK;
404 - return { training_thread->training_cns, training_response };
405 -}
406 -
407 -const char *db_models_create_table =
408 - "CREATE TABLE IF NOT EXISTS models("
409 - " dim_id BLOB, dim_str TEXT, after INT, before INT,"
410 - " min_dist REAL, max_dist REAL,"
411 - " c00 REAL, c01 REAL, c02 REAL, c03 REAL, c04 REAL, c05 REAL,"
412 - " c10 REAL, c11 REAL, c12 REAL, c13 REAL, c14 REAL, c15 REAL,"
413 - " PRIMARY KEY(dim_id, after)"
414 - ");";
415 -
416 -const char *db_models_add_model =
417 - "INSERT OR REPLACE INTO models("
418 - " dim_id, dim_str, after, before,"
419 - " min_dist, max_dist,"
420 - " c00, c01, c02, c03, c04, c05,"
421 - " c10, c11, c12, c13, c14, c15)"
422 - "VALUES("
423 - " @dim_id, @dim_str, @after, @before,"
424 - " @min_dist, @max_dist,"
425 - " @c00, @c01, @c02, @c03, @c04, @c05,"
426 - " @c10, @c11, @c12, @c13, @c14, @c15);";
427 -
428 -const char *db_models_load =
429 - "SELECT * FROM models "
430 - "WHERE dim_id == @dim_id AND after >= @after ORDER BY before ASC;";
431 -
432 -const char *db_models_delete =
433 - "DELETE FROM models "
434 - "WHERE dim_id = @dim_id AND before < @before;";
435 -
436 -static int
437 -ml_dimension_add_model(ml_dimension_t *dim)
438 -{
439 - static __thread sqlite3_stmt *res = NULL;
440 - int param = 0;
441 - int rc = 0;
442 -
443 - if (unlikely(!db)) {
444 - error_report("Database has not been initialized");
445 - return 1;
446 - }
447 -
448 - if (unlikely(!res)) {
449 - rc = prepare_statement(db, db_models_add_model, &res);
450 - if (unlikely(rc != SQLITE_OK)) {
451 - error_report("Failed to prepare statement to store model, rc = %d", rc);
452 - return 1;
453 - }
454 - }
455 -
456 - rc = sqlite3_bind_blob(res, ++param, &dim->rd->metric_uuid, sizeof(dim->rd->metric_uuid), SQLITE_STATIC);
457 - if (unlikely(rc != SQLITE_OK))
458 - goto bind_fail;
459 -
460 - char id[1024];
461 - snprintfz(id, 1024 - 1, "%s.%s", rrdset_id(dim->rd->rrdset), rrddim_id(dim->rd));
462 - rc = sqlite3_bind_text(res, ++param, id, -1, SQLITE_STATIC);
463 - if (unlikely(rc != SQLITE_OK))
464 - goto bind_fail;
465 -
466 - rc = sqlite3_bind_int(res, ++param, (int) dim->kmeans.after);
467 - if (unlikely(rc != SQLITE_OK))
468 - goto bind_fail;
469 -
470 - rc = sqlite3_bind_int(res, ++param, (int) dim->kmeans.before);
471 - if (unlikely(rc != SQLITE_OK))
472 - goto bind_fail;
473 -
474 - rc = sqlite3_bind_double(res, ++param, dim->kmeans.min_dist);
475 - if (unlikely(rc != SQLITE_OK))
476 - goto bind_fail;
477 -
478 - rc = sqlite3_bind_double(res, ++param, dim->kmeans.max_dist);
479 - if (unlikely(rc != SQLITE_OK))
480 - goto bind_fail;
481 -
482 - if (dim->kmeans.cluster_centers.size() != 2)
483 - fatal("Expected 2 cluster centers, got %zu", dim->kmeans.cluster_centers.size());
484 -
485 - for (const DSample &ds : dim->kmeans.cluster_centers) {
486 - if (ds.size() != 6)
487 - fatal("Expected dsample with 6 dimensions, got %ld", ds.size());
488 -
489 - for (long idx = 0; idx != ds.size(); idx++) {
490 - calculated_number_t cn = ds(idx);
491 - int rc = sqlite3_bind_double(res, ++param, cn);
492 - if (unlikely(rc != SQLITE_OK))
493 - goto bind_fail;
494 - }
495 - }
496 -
497 - rc = execute_insert(res);
498 - if (unlikely(rc != SQLITE_DONE))
499 - error_report("Failed to store model, rc = %d", rc);
500 -
501 - rc = sqlite3_reset(res);
502 - if (unlikely(rc != SQLITE_OK))
503 - error_report("Failed to reset statement when storing model, rc = %d", rc);
504 -
505 - return 0;
506 -
507 -bind_fail:
508 - error_report("Failed to bind parameter %d to store model, rc = %d", param, rc);
509 - rc = sqlite3_reset(res);
510 - if (unlikely(rc != SQLITE_OK))
511 - error_report("Failed to reset statement to store model, rc = %d", rc);
512 - return 1;
513 -}
514 -
515 -static int
516 -ml_dimension_delete_models(ml_dimension_t *dim)
517 -{
518 - static __thread sqlite3_stmt *res = NULL;
519 - int rc = 0;
520 - int param = 0;
521 -
522 - if (unlikely(!db)) {
523 - error_report("Database has not been initialized");
524 - return 1;
525 - }
526 -
527 - if (unlikely(!res)) {
528 - rc = prepare_statement(db, db_models_delete, &res);
529 - if (unlikely(rc != SQLITE_OK)) {
530 - error_report("Failed to prepare statement to delete models, rc = %d", rc);
531 - return 1;
532 - }
533 - }
534 -
535 - rc = sqlite3_bind_blob(res, ++param, &dim->rd->metric_uuid, sizeof(dim->rd->metric_uuid), SQLITE_STATIC);
536 - if (unlikely(rc != SQLITE_OK))
537 - goto bind_fail;
538 -
539 - rc = sqlite3_bind_int(res, ++param, (int) dim->kmeans.before - (Cfg.num_models_to_use * Cfg.train_every));
540 - if (unlikely(rc != SQLITE_OK))
541 - goto bind_fail;
542 -
543 - rc = execute_insert(res);
544 - if (unlikely(rc != SQLITE_DONE))
545 - error_report("Failed to delete models, rc = %d", rc);
546 -
547 - rc = sqlite3_reset(res);
548 - if (unlikely(rc != SQLITE_OK))
549 - error_report("Failed to reset statement when deleting models, rc = %d", rc);
550 -
551 - return 0;
552 -
553 -bind_fail:
554 - error_report("Failed to bind parameter %d to delete models, rc = %d", param, rc);
555 - rc = sqlite3_reset(res);
556 - if (unlikely(rc != SQLITE_OK))
557 - error_report("Failed to reset statement to delete models, rc = %d", rc);
558 - return 1;
559 -}
560 -
561 -static int
562 -ml_dimension_load_models(ml_dimension_t *dim) {
563 - std::vector<ml_kmeans_t> V;
564 -
565 - static __thread sqlite3_stmt *res = NULL;
566 - int rc = 0;
567 - int param = 0;
568 -
569 - if (unlikely(!db)) {
570 - error_report("Database has not been initialized");
571 - return 1;
572 - }
573 -
574 - if (unlikely(!res)) {
575 - rc = prepare_statement(db, db_models_load, &res);
576 - if (unlikely(rc != SQLITE_OK)) {
577 - error_report("Failed to prepare statement to load models, rc = %d", rc);
578 - return 1;
579 - }
580 - }
581 -
582 - rc = sqlite3_bind_blob(res, ++param, &dim->rd->metric_uuid, sizeof(dim->rd->metric_uuid), SQLITE_STATIC);
583 - if (unlikely(rc != SQLITE_OK))
584 - goto bind_fail;
585 -
586 - rc = sqlite3_bind_int(res, ++param, now_realtime_usec() - (Cfg.num_models_to_use * Cfg.max_train_samples));
587 - if (unlikely(rc != SQLITE_OK))
588 - goto bind_fail;
589 -
590 - dim->km_contexts.reserve(Cfg.num_models_to_use);
591 - while ((rc = sqlite3_step_monitored(res)) == SQLITE_ROW) {
592 - ml_kmeans_t km;
593 -
594 - km.after = sqlite3_column_int(res, 2);
595 - km.before = sqlite3_column_int(res, 3);
596 -
597 - km.min_dist = sqlite3_column_int(res, 4);
598 - km.max_dist = sqlite3_column_int(res, 5);
599 -
600 - km.cluster_centers.resize(2);
601 -
602 - km.cluster_centers[0].set_size(Cfg.lag_n + 1);
603 - km.cluster_centers[0](0) = sqlite3_column_double(res, 6);
604 - km.cluster_centers[0](1) = sqlite3_column_double(res, 7);
605 - km.cluster_centers[0](2) = sqlite3_column_double(res, 8);
606 - km.cluster_centers[0](3) = sqlite3_column_double(res, 9);
607 - km.cluster_centers[0](4) = sqlite3_column_double(res, 10);
608 - km.cluster_centers[0](5) = sqlite3_column_double(res, 11);
609 -
610 - km.cluster_centers[1].set_size(Cfg.lag_n + 1);
611 - km.cluster_centers[1](0) = sqlite3_column_double(res, 12);
612 - km.cluster_centers[1](1) = sqlite3_column_double(res, 13);
613 - km.cluster_centers[1](2) = sqlite3_column_double(res, 14);
614 - km.cluster_centers[1](3) = sqlite3_column_double(res, 15);
615 - km.cluster_centers[1](4) = sqlite3_column_double(res, 16);
616 - km.cluster_centers[1](5) = sqlite3_column_double(res, 17);
617 -
618 - dim->km_contexts.push_back(km);
619 - }
620 -
621 - if (unlikely(rc != SQLITE_DONE))
622 - error_report("Failed to load models, rc = %d", rc);
623 -
624 - rc = sqlite3_reset(res);
625 - if (unlikely(rc != SQLITE_OK))
626 - error_report("Failed to reset statement when loading models, rc = %d", rc);
627 -
628 - return 0;
629 -
630 -bind_fail:
631 - error_report("Failed to bind parameter %d to load models, rc = %d", param, rc);
632 - rc = sqlite3_reset(res);
633 - if (unlikely(rc != SQLITE_OK))
634 - error_report("Failed to reset statement to load models, rc = %d", rc);
635 - return 1;
636 -}
637 -
638 -static int
639 -ml_dimension_update_models(ml_dimension_t *dim)
640 -{
641 - int rc;
642 -
643 - if (dim->km_contexts.empty()) {
644 - rc = ml_dimension_load_models(dim);
645 - if (rc)
646 - return rc;
647 - }
648 -
649 - rc = ml_dimension_add_model(dim);
650 - if (rc)
651 - return rc;
652 -
653 - return ml_dimension_delete_models(dim);
394 + return { tls_data.training_cns, training_response };
395 }
396
397 static enum ml_training_result
657 -ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
398 +ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &training_request)
399 {
659 - worker_is_busy(WORKER_TRAIN_QUERY);
660 - auto P = ml_dimension_calculated_numbers(training_thread, dim, training_request);
400 + auto P = ml_dimension_calculated_numbers(dim, training_request);
401 ml_training_response_t training_response = P.second;
402
403 if (training_response.result != TRAINING_RESULT_OK) {
@@ -686,56 +426,31 @@ ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *
426 }
427
428 // compute kmeans
689 - worker_is_busy(WORKER_TRAIN_KMEANS);
429 {
691 - memcpy(training_thread->scratch_training_cns, training_thread->training_cns,
430 + memcpy(tls_data.scratch_training_cns, tls_data.training_cns,
431 training_response.total_values * sizeof(calculated_number_t));
432
433 ml_features_t features = {
434 Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
696 - training_thread->scratch_training_cns, training_response.total_values,
697 - training_thread->training_cns, training_response.total_values,
698 - training_thread->training_samples
435 + tls_data.scratch_training_cns, training_response.total_values,
436 + tls_data.training_cns, training_response.total_values,
437 + tls_data.training_samples
438 };
439 ml_features_preprocess(&features);
440
702 - ml_kmeans_init(&dim->kmeans);
703 - ml_kmeans_train(&dim->kmeans, &features, training_response.query_after_t, training_response.query_before_t);
441 + ml_kmeans_init(&dim->kmeans, 2, 1000);
442 + ml_kmeans_train(&dim->kmeans, &features);
443 }
444
706 - // update models
445 + // update kmeans models
446 {
447 netdata_mutex_lock(&dim->mutex);
448
710 - worker_is_busy(WORKER_TRAIN_LOAD_MODELS);
711 -
712 - int rc = ml_dimension_update_models(dim);
713 - if (rc) {
714 - error("Failed to update models for %s [%u, %u]", rrddim_id(dim->rd), dim->kmeans.after, dim->kmeans.before);
715 - }
716 -
717 - worker_is_busy(WORKER_TRAIN_UPDATE_MODELS);
718 -
449 if (dim->km_contexts.size() < Cfg.num_models_to_use) {
450 dim->km_contexts.push_back(std::move(dim->kmeans));
451 } else {
722 - bool can_drop_middle_km = false;
723 -
724 - if (Cfg.num_models_to_use > 2) {
725 - const ml_kmeans_t *old_km = &dim->km_contexts[dim->km_contexts.size() - 1];
726 - const ml_kmeans_t *middle_km = &dim->km_contexts[dim->km_contexts.size() - 2];
727 - const ml_kmeans_t *new_km = &dim->kmeans;
728 -
729 - can_drop_middle_km = (middle_km->after < old_km->before) &&
730 - (middle_km->before > new_km->after);
731 - }
732 -
733 - if (can_drop_middle_km) {
734 - dim->km_contexts.back() = dim->kmeans;
735 - } else {
736 - std::rotate(std::begin(dim->km_contexts), std::begin(dim->km_contexts) + 1, std::end(dim->km_contexts));
737 - dim->km_contexts[dim->km_contexts.size() - 1] = std::move(dim->kmeans);
738 - }
452 + std::rotate(std::begin(dim->km_contexts), std::begin(dim->km_contexts) + 1, std::end(dim->km_contexts));
453 + dim->km_contexts[dim->km_contexts.size() - 1] = std::move(dim->kmeans);
454 }
455
456 dim->mt = METRIC_TYPE_CONSTANT;
@@ -779,16 +494,11 @@ ml_dimension_schedule_for_training(ml_dimension_t *dim, time_t curr_time)
494 }
495
496 if (schedule_for_training) {
497 + ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
498 ml_training_request_t req = {
783 - string_dup(dim->rd->rrdset->rrdhost->hostname),
784 - string_dup(dim->rd->rrdset->id),
785 - string_dup(dim->rd->id),
786 - curr_time,
787 - rrddim_first_entry_s(dim->rd),
788 - rrddim_last_entry_s(dim->rd),
499 + string_dup(dim->rd->rrdset->id), string_dup(dim->rd->id),
500 + curr_time, rrddim_first_entry_s(dim->rd), rrddim_last_entry_s(dim->rd),
501 };
790 -
791 - ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
502 ml_queue_push(host->training_queue, req);
503 }
504 }
@@ -964,6 +674,7 @@ ml_host_detect_once(ml_host_t *host)
674
675 host->mls = {};
676 ml_machine_learning_stats_t mls_copy = {};
677 + ml_training_stats_t ts_copy = {};
678
679 {
680 netdata_mutex_lock(&host->mutex);
@@ -1007,14 +718,54 @@ ml_host_detect_once(ml_host_t *host)
718
719 mls_copy = host->mls;
720
721 + /*
722 + * training stats
723 + */
724 + ts_copy = host->ts;
725 +
726 + host->ts.queue_size = 0;
727 + host->ts.num_popped_items = 0;
728 +
729 + host->ts.allotted_ut = 0;
730 + host->ts.consumed_ut = 0;
731 + host->ts.remaining_ut = 0;
732 +
733 + host->ts.training_result_ok = 0;
734 + host->ts.training_result_invalid_query_time_range = 0;
735 + host->ts.training_result_not_enough_collected_values = 0;
736 + host->ts.training_result_null_acquired_dimension = 0;
737 + host->ts.training_result_chart_under_replication = 0;
738 +
739 netdata_mutex_unlock(&host->mutex);
740 }
741
742 + // Calc the avg values
743 + if (ts_copy.num_popped_items) {
744 + ts_copy.queue_size /= ts_copy.num_popped_items;
745 + ts_copy.allotted_ut /= ts_copy.num_popped_items;
746 + ts_copy.consumed_ut /= ts_copy.num_popped_items;
747 + ts_copy.remaining_ut /= ts_copy.num_popped_items;
748 +
749 + ts_copy.training_result_ok /= ts_copy.num_popped_items;
750 + ts_copy.training_result_invalid_query_time_range /= ts_copy.num_popped_items;
751 + ts_copy.training_result_not_enough_collected_values /= ts_copy.num_popped_items;
752 + ts_copy.training_result_null_acquired_dimension /= ts_copy.num_popped_items;
753 + ts_copy.training_result_chart_under_replication /= ts_copy.num_popped_items;
754 + } else {
755 + ts_copy.queue_size = 0;
756 + ts_copy.allotted_ut = 0;
757 + ts_copy.consumed_ut = 0;
758 + ts_copy.remaining_ut = 0;
759 + }
760 +
761 worker_is_busy(WORKER_JOB_DETECTION_DIM_CHART);
762 ml_update_dimensions_chart(host, mls_copy);
763
764 worker_is_busy(WORKER_JOB_DETECTION_HOST_CHART);
765 ml_update_host_and_detection_rate_charts(host, host->host_anomaly_rate * 10000.0);
766 +
767 + worker_is_busy(WORKER_JOB_DETECTION_STATS);
768 + ml_update_training_statistics_chart(host, ts_copy);
769 }
770
771 typedef struct {
@@ -1023,21 +774,18 @@ typedef struct {
774 } ml_acquired_dimension_t;
775
776 static ml_acquired_dimension_t
1026 -ml_acquired_dimension_get(STRING *host_id, STRING *chart_id, STRING *dimension_id)
777 +ml_acquired_dimension_get(RRDHOST *rh, STRING *chart_id, STRING *dimension_id)
778 {
779 RRDDIM_ACQUIRED *acq_rd = NULL;
780 ml_dimension_t *dim = NULL;
781
1031 - RRDHOST *rh = rrdhost_find_by_hostname(string2str(host_id));
1032 - if (rh) {
1033 - RRDSET *rs = rrdset_find(rh, string2str(chart_id));
1034 - if (rs) {
1035 - acq_rd = rrddim_find_and_acquire(rs, string2str(dimension_id));
1036 - if (acq_rd) {
1037 - RRDDIM *rd = rrddim_acquired_to_rrddim(acq_rd);
1038 - if (rd)
1039 - dim = (ml_dimension_t *) rd->ml_dimension;
1040 - }
782 + RRDSET *rs = rrdset_find(rh, string2str(chart_id));
783 + if (rs) {
784 + acq_rd = rrddim_find_and_acquire(rs, string2str(dimension_id));
785 + if (acq_rd) {
786 + RRDDIM *rd = rrddim_acquired_to_rrddim(acq_rd);
787 + if (rd)
788 + dim = (ml_dimension_t *) rd->ml_dimension;
789 }
790 }
791
@@ -1058,12 +806,110 @@ ml_acquired_dimension_release(ml_acquired_dimension_t acq_dim)
806 }
807
808 static enum ml_training_result
1061 -ml_acquired_dimension_train(ml_training_thread_t *training_thread, ml_acquired_dimension_t acq_dim, const ml_training_request_t &tr)
809 +ml_acquired_dimension_train(ml_acquired_dimension_t acq_dim, const ml_training_request_t &TR)
810 {
811 if (!acq_dim.dim)
812 return TRAINING_RESULT_NULL_ACQUIRED_DIMENSION;
813
1066 - return ml_dimension_train_model(training_thread, acq_dim.dim, tr);
814 + return ml_dimension_train_model(acq_dim.dim, TR);
815 +}
816 +
817 +#define WORKER_JOB_TRAINING_FIND 0
818 +#define WORKER_JOB_TRAINING_TRAIN 1
819 +#define WORKER_JOB_TRAINING_STATS 2
820 +
821 +static void
822 +ml_host_train(ml_host_t *host)
823 +{
824 + worker_register("MLTRAIN");
825 + worker_register_job_name(WORKER_JOB_TRAINING_FIND, "find");
826 + worker_register_job_name(WORKER_JOB_TRAINING_TRAIN, "train");
827 + worker_register_job_name(WORKER_JOB_TRAINING_STATS, "stats");
828 +
829 + service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, (force_quit_t ) ml_host_cancel_training_thread, host->rh, true);
830 +
831 + while (service_running(SERVICE_ML_TRAINING)) {
832 + ml_training_request_t training_req = ml_queue_pop(host->training_queue);
833 + size_t queue_size = ml_queue_size(host->training_queue) + 1;
834 +
835 + if (host->threads_cancelled) {
836 + info("Stopping training thread for host %s because it was cancelled", rrdhost_hostname(host->rh));
837 + break;
838 + }
839 +
840 + usec_t allotted_ut = (Cfg.train_every * host->rh->rrd_update_every * USEC_PER_SEC) / queue_size;
841 + if (allotted_ut > USEC_PER_SEC)
842 + allotted_ut = USEC_PER_SEC;
843 +
844 + usec_t start_ut = now_monotonic_usec();
845 + enum ml_training_result training_res;
846 + {
847 + worker_is_busy(WORKER_JOB_TRAINING_FIND);
848 + ml_acquired_dimension_t acq_dim = ml_acquired_dimension_get(host->rh, training_req.chart_id, training_req.dimension_id);
849 +
850 + worker_is_busy(WORKER_JOB_TRAINING_TRAIN);
851 + training_res = ml_acquired_dimension_train(acq_dim, training_req);
852 +
853 + string_freez(training_req.chart_id);
854 + string_freez(training_req.dimension_id);
855 +
856 + ml_acquired_dimension_release(acq_dim);
857 + }
858 + usec_t consumed_ut = now_monotonic_usec() - start_ut;
859 +
860 + worker_is_busy(WORKER_JOB_TRAINING_STATS);
861 +
862 + usec_t remaining_ut = 0;
863 + if (consumed_ut < allotted_ut)
864 + remaining_ut = allotted_ut - consumed_ut;
865 +
866 + {
867 + netdata_mutex_lock(&host->mutex);
868 +
869 + host->ts.queue_size += queue_size;
870 + host->ts.num_popped_items += 1;
871 +
872 + host->ts.allotted_ut += allotted_ut;
873 + host->ts.consumed_ut += consumed_ut;
874 + host->ts.remaining_ut += remaining_ut;
875 +
876 + switch (training_res) {
877 + case TRAINING_RESULT_OK:
878 + host->ts.training_result_ok += 1;
879 + break;
880 + case TRAINING_RESULT_INVALID_QUERY_TIME_RANGE:
881 + host->ts.training_result_invalid_query_time_range += 1;
882 + break;
883 + case TRAINING_RESULT_NOT_ENOUGH_COLLECTED_VALUES:
884 + host->ts.training_result_not_enough_collected_values += 1;
885 + break;
886 + case TRAINING_RESULT_NULL_ACQUIRED_DIMENSION:
887 + host->ts.training_result_null_acquired_dimension += 1;
888 + break;
889 + case TRAINING_RESULT_CHART_UNDER_REPLICATION:
890 + host->ts.training_result_chart_under_replication += 1;
891 + break;
892 + }
893 +
894 + netdata_mutex_unlock(&host->mutex);
895 + }
896 +
897 + worker_is_idle();
898 + std::this_thread::sleep_for(std::chrono::microseconds{remaining_ut});
899 + worker_is_busy(0);
900 + }
901 +}
902 +
903 +static void *
904 +train_main(void *arg)
905 +{
906 + size_t max_elements_needed_for_training = Cfg.max_train_samples * (Cfg.lag_n + 1);
907 + tls_data.training_cns = new calculated_number_t[max_elements_needed_for_training]();
908 + tls_data.scratch_training_cns = new calculated_number_t[max_elements_needed_for_training]();
909 +
910 + ml_host_t *host = (ml_host_t *) arg;
911 + ml_host_train(host);
912 + return NULL;
913 }
914
915 static void *
@@ -1077,55 +923,25 @@ ml_detect_main(void *arg)
923 worker_register_job_name(WORKER_JOB_DETECTION_HOST_CHART, "host chart");
924 worker_register_job_name(WORKER_JOB_DETECTION_STATS, "training stats");
925
926 + service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, NULL, NULL, true);
927 +
928 heartbeat_t hb;
929 heartbeat_init(&hb);
930
1083 - while (!Cfg.detection_stop) {
931 + while (service_running((SERVICE_TYPE)(SERVICE_ML_PREDICTION | SERVICE_COLLECTORS))) {
932 worker_is_idle();
933 heartbeat_next(&hb, USEC_PER_SEC);
934
1087 - RRDHOST *rh;
1088 - rrd_rdlock();
1089 - rrdhost_foreach_read(rh) {
935 + void *rhp;
936 + dfe_start_reentrant(rrdhost_root_index, rhp) {
937 + RRDHOST *rh = (RRDHOST *) rhp;
938 +
939 if (!rh->ml_host)
940 continue;
941
942 ml_host_detect_once((ml_host_t *) rh->ml_host);
943 }
1095 - rrd_unlock();
1096 -
1097 - if (Cfg.enable_statistics_charts) {
1098 - // collect and update training thread stats
1099 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1100 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1101 -
1102 - netdata_mutex_lock(&training_thread->nd_mutex);
1103 - ml_training_stats_t training_stats = training_thread->training_stats;
1104 - training_thread->training_stats = {};
1105 - netdata_mutex_unlock(&training_thread->nd_mutex);
1106 -
1107 - // calc the avg values
1108 - if (training_stats.num_popped_items) {
1109 - training_stats.queue_size /= training_stats.num_popped_items;
1110 - training_stats.allotted_ut /= training_stats.num_popped_items;
1111 - training_stats.consumed_ut /= training_stats.num_popped_items;
1112 - training_stats.remaining_ut /= training_stats.num_popped_items;
1113 - } else {
1114 - training_stats.queue_size = 0;
1115 - training_stats.allotted_ut = 0;
1116 - training_stats.consumed_ut = 0;
1117 - training_stats.remaining_ut = 0;
1118 -
1119 - training_stats.training_result_ok = 0;
1120 - training_stats.training_result_invalid_query_time_range = 0;
1121 - training_stats.training_result_not_enough_collected_values = 0;
1122 - training_stats.training_result_null_acquired_dimension = 0;
1123 - training_stats.training_result_chart_under_replication = 0;
1124 - }
1125 -
1126 - ml_update_training_statistics_chart(training_thread, training_stats);
1127 - }
1128 - }
944 + dfe_done(rhp);
945 }
946
947 return NULL;
@@ -1159,6 +975,31 @@ bool ml_streaming_enabled()
975 return Cfg.stream_anomaly_detection_charts;
976 }
977
978 +void ml_init()
979 +{
980 + // Read config values
981 + ml_config_load(&Cfg);
982 +
983 + if (!Cfg.enable_anomaly_detection)
984 + return;
985 +
986 + // Generate random numbers to efficiently sample the features we need
987 + // for KMeans clustering.
988 + std::random_device RD;
989 + std::mt19937 Gen(RD());
990 +
991 + Cfg.random_nums.reserve(Cfg.max_train_samples);
992 + for (size_t Idx = 0; Idx != Cfg.max_train_samples; Idx++)
993 + Cfg.random_nums.push_back(Gen());
994 +
995 +
996 + // start detection & training threads
997 + char tag[NETDATA_THREAD_TAG_MAX + 1];
998 +
999 + snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "PREDICT");
1000 + netdata_thread_create(&Cfg.detection_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_detect_main, NULL);
1001 +}
1002 +
1003 void ml_host_new(RRDHOST *rh)
1004 {
1005 if (!ml_enabled(rh))
@@ -1168,12 +1009,14 @@ void ml_host_new(RRDHOST *rh)
1009
1010 host->rh = rh;
1011 host->mls = ml_machine_learning_stats_t();
1171 - //host->ts = ml_training_stats_t();
1172 -
1173 - static std::atomic<size_t> times_called(0);
1174 - host->training_queue = Cfg.training_threads[times_called++ % Cfg.num_training_threads].training_queue;
1012 + host->ts = ml_training_stats_t();
1013
1014 host->host_anomaly_rate = 0.0;
1015 + host->threads_running = false;
1016 + host->threads_cancelled = false;
1017 + host->threads_joined = false;
1018 +
1019 + host->training_queue = ml_queue_init();
1020
1021 netdata_mutex_init(&host->mutex);
1022
@@ -1187,6 +1030,7 @@ void ml_host_delete(RRDHOST *rh)
1030 return;
1031
1032 netdata_mutex_destroy(&host->mutex);
1033 + ml_queue_destroy(host->training_queue);
1034
1035 delete host;
1036 rh->ml_host = NULL;
@@ -1253,6 +1097,69 @@ void ml_host_get_models(RRDHOST *rh, BUFFER *wb)
1097 error("Fetching KMeans models is not supported yet");
1098 }
1099
1100 +void ml_host_start_training_thread(RRDHOST *rh)
1101 +{
1102 + if (!rh || !rh->ml_host)
1103 + return;
1104 +
1105 + ml_host_t *host = (ml_host_t *) rh->ml_host;
1106 +
1107 + if (host->threads_running) {
1108 + error("Anomaly detections threads for host %s are already-up and running.", rrdhost_hostname(host->rh));
1109 + return;
1110 + }
1111 +
1112 + host->threads_running = true;
1113 + host->threads_cancelled = false;
1114 + host->threads_joined = false;
1115 +
1116 + char tag[NETDATA_THREAD_TAG_MAX + 1];
1117 +
1118 + snprintfz(tag, NETDATA_THREAD_TAG_MAX, "MLTR[%s]", rrdhost_hostname(host->rh));
1119 + netdata_thread_create(&host->training_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, train_main, static_cast<void *>(host));
1120 +}
1121 +
1122 +void ml_host_cancel_training_thread(RRDHOST *rh)
1123 +{
1124 + if (!rh || !rh->ml_host)
1125 + return;
1126 +
1127 + ml_host_t *host = (ml_host_t *) rh->ml_host;
1128 +
1129 + if (!host->threads_running) {
1130 + error("Anomaly detections threads for host %s have already been stopped.", rrdhost_hostname(host->rh));
1131 + return;
1132 + }
1133 +
1134 + if (!host->threads_cancelled) {
1135 + host->threads_cancelled = true;
1136 +
1137 + // Signal the training queue to stop popping-items
1138 + ml_queue_signal(host->training_queue);
1139 + netdata_thread_cancel(host->training_thread);
1140 + }
1141 +}
1142 +
1143 +void ml_host_stop_training_thread(RRDHOST *rh)
1144 +{
1145 + if (!rh || !rh->ml_host)
1146 + return;
1147 +
1148 + ml_host_cancel_training_thread(rh);
1149 +
1150 + ml_host_t *host = (ml_host_t *) rh->ml_host;
1151 +
1152 + if (!host->threads_joined) {
1153 + host->threads_joined = true;
1154 + host->threads_running = false;
1155 +
1156 + delete[] tls_data.training_cns;
1157 + delete[] tls_data.scratch_training_cns;
1158 +
1159 + netdata_thread_join(host->training_thread, NULL);
1160 + }
1161 +}
1162 +
1163 void ml_chart_new(RRDSET *rs)
1164 {
1165 ml_host_t *host = (ml_host_t *) rs->rrdhost->ml_host;
@@ -1318,7 +1225,7 @@ void ml_dimension_new(RRDDIM *rd)
1225
1226 dim->last_training_time = 0;
1227
1321 - ml_kmeans_init(&dim->kmeans);
1228 + ml_kmeans_init(&dim->kmeans, 2, 1000);
1229
1230 if (simple_pattern_matches(Cfg.sp_charts_to_skip, rrdset_name(rd->rrdset)))
1231 dim->mls = MACHINE_LEARNING_STATUS_DISABLED_DUE_TO_EXCLUDED_CHART;
@@ -1357,216 +1264,3 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
1264
1265 return is_anomalous;
1266 }
1360 -
1361 -static void *ml_train_main(void *arg) {
1362 - ml_training_thread_t *training_thread = (ml_training_thread_t *) arg;
1363 -
1364 - char worker_name[1024];
1365 - snprintfz(worker_name, 1024, "training_thread_%zu", training_thread->id);
1366 - worker_register("MLTRAIN");
1367 -
1368 - worker_register_job_name(WORKER_TRAIN_QUEUE_POP, "pop queue");
1369 - worker_register_job_name(WORKER_TRAIN_ACQUIRE_DIMENSION, "acquire");
1370 - worker_register_job_name(WORKER_TRAIN_QUERY, "query");
1371 - worker_register_job_name(WORKER_TRAIN_KMEANS, "kmeans");
1372 - worker_register_job_name(WORKER_TRAIN_UPDATE_MODELS, "update models");
1373 - worker_register_job_name(WORKER_TRAIN_LOAD_MODELS, "load models");
1374 - worker_register_job_name(WORKER_TRAIN_RELEASE_DIMENSION, "release");
1375 - worker_register_job_name(WORKER_TRAIN_UPDATE_HOST, "update host");
1376 -
1377 - while (!Cfg.training_stop) {
1378 - worker_is_busy(WORKER_TRAIN_QUEUE_POP);
1379 -
1380 - ml_training_request_t training_req = ml_queue_pop(training_thread->training_queue);
1381 -
1382 - // we know this thread has been cancelled, when the queue starts
1383 - // returning "null" requests without blocking on queue's pop().
1384 - if (training_req.host_id == NULL)
1385 - break;
1386 -
1387 - size_t queue_size = ml_queue_size(training_thread->training_queue) + 1;
1388 -
1389 - usec_t allotted_ut = (Cfg.train_every * USEC_PER_SEC) / queue_size;
1390 - if (allotted_ut > USEC_PER_SEC)
1391 - allotted_ut = USEC_PER_SEC;
1392 -
1393 - usec_t start_ut = now_monotonic_usec();
1394 -
1395 - enum ml_training_result training_res;
1396 - {
1397 - worker_is_busy(WORKER_TRAIN_ACQUIRE_DIMENSION);
1398 - ml_acquired_dimension_t acq_dim = ml_acquired_dimension_get(
1399 - training_req.host_id,
1400 - training_req.chart_id,
1401 - training_req.dimension_id);
1402 -
1403 - training_res = ml_acquired_dimension_train(training_thread, acq_dim, training_req);
1404 -
1405 - string_freez(training_req.host_id);
1406 - string_freez(training_req.chart_id);
1407 - string_freez(training_req.dimension_id);
1408 -
1409 - worker_is_busy(WORKER_TRAIN_RELEASE_DIMENSION);
1410 - ml_acquired_dimension_release(acq_dim);
1411 - }
1412 -
1413 - usec_t consumed_ut = now_monotonic_usec() - start_ut;
1414 -
1415 - usec_t remaining_ut = 0;
1416 - if (consumed_ut < allotted_ut)
1417 - remaining_ut = allotted_ut - consumed_ut;
1418 -
1419 - if (Cfg.enable_statistics_charts) {
1420 - worker_is_busy(WORKER_TRAIN_UPDATE_HOST);
1421 -
1422 - netdata_mutex_lock(&training_thread->nd_mutex);
1423 -
1424 - training_thread->training_stats.queue_size += queue_size;
1425 - training_thread->training_stats.num_popped_items += 1;
1426 -
1427 - training_thread->training_stats.allotted_ut += allotted_ut;
1428 - training_thread->training_stats.consumed_ut += consumed_ut;
1429 - training_thread->training_stats.remaining_ut += remaining_ut;
1430 -
1431 - switch (training_res) {
1432 - case TRAINING_RESULT_OK:
1433 - training_thread->training_stats.training_result_ok += 1;
1434 - break;
1435 - case TRAINING_RESULT_INVALID_QUERY_TIME_RANGE:
1436 - training_thread->training_stats.training_result_invalid_query_time_range += 1;
1437 - break;
1438 - case TRAINING_RESULT_NOT_ENOUGH_COLLECTED_VALUES:
1439 - training_thread->training_stats.training_result_not_enough_collected_values += 1;
1440 - break;
1441 - case TRAINING_RESULT_NULL_ACQUIRED_DIMENSION:
1442 - training_thread->training_stats.training_result_null_acquired_dimension += 1;
1443 - break;
1444 - case TRAINING_RESULT_CHART_UNDER_REPLICATION:
1445 - training_thread->training_stats.training_result_chart_under_replication += 1;
1446 - break;
1447 - }
1448 -
1449 - netdata_mutex_unlock(&training_thread->nd_mutex);
1450 - }
1451 -
1452 - worker_is_idle();
1453 - std::this_thread::sleep_for(std::chrono::microseconds{remaining_ut});
1454 - }
1455 -
1456 - return NULL;
1457 -}
1458 -
1459 -void ml_init()
1460 -{
1461 - // Read config values
1462 - ml_config_load(&Cfg);
1463 -
1464 - if (!Cfg.enable_anomaly_detection)
1465 - return;
1466 -
1467 - // Generate random numbers to efficiently sample the features we need
1468 - // for KMeans clustering.
1469 - std::random_device RD;
1470 - std::mt19937 Gen(RD());
1471 -
1472 - Cfg.random_nums.reserve(Cfg.max_train_samples);
1473 - for (size_t Idx = 0; Idx != Cfg.max_train_samples; Idx++)
1474 - Cfg.random_nums.push_back(Gen());
1475 -
1476 - // init training thread-specific data
1477 - Cfg.training_threads.resize(Cfg.num_training_threads);
1478 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1479 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1480 -
1481 - size_t max_elements_needed_for_training = Cfg.max_train_samples * (Cfg.lag_n + 1);
1482 - training_thread->training_cns = new calculated_number_t[max_elements_needed_for_training]();
1483 - training_thread->scratch_training_cns = new calculated_number_t[max_elements_needed_for_training]();
1484 -
1485 - training_thread->id = idx;
1486 - training_thread->training_queue = ml_queue_init();
1487 - netdata_mutex_init(&training_thread->nd_mutex);
1488 - }
1489 -
1490 - // open sqlite db
1491 - char path[FILENAME_MAX];
1492 - snprintfz(path, FILENAME_MAX - 1, "%s/%s", netdata_configured_cache_dir, "ml.db");
1493 - int rc = sqlite3_open(path, &db);
1494 - if (rc != SQLITE_OK) {
1495 - error_report("Failed to initialize database at %s, due to \"%s\"", path, sqlite3_errstr(rc));
1496 - sqlite3_close(db);
1497 - db = NULL;
1498 - }
1499 -
1500 - if (db) {
1501 - char *err = NULL;
1502 - int rc = sqlite3_exec(db, db_models_create_table, NULL, NULL, &err);
1503 - if (rc != SQLITE_OK) {
1504 - error_report("Failed to create models table (%s, %s)", sqlite3_errstr(rc), err ? err : "");
1505 - sqlite3_close(db);
1506 - db = NULL;
1507 - }
1508 - }
1509 -}
1510 -
1511 -void ml_fini() {
1512 - int rc = sqlite3_close_v2(db);
1513 - if (unlikely(rc != SQLITE_OK))
1514 - error_report("Error %d while closing the SQLite database, %s", rc, sqlite3_errstr(rc));
1515 -}
1516 -
1517 -void ml_start_threads() {
1518 - // start detection & training threads
1519 - Cfg.detection_stop = false;
1520 - Cfg.training_stop = false;
1521 -
1522 - char tag[NETDATA_THREAD_TAG_MAX + 1];
1523 -
1524 - snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "PREDICT");
1525 - netdata_thread_create(&Cfg.detection_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_detect_main, NULL);
1526 -
1527 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1528 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1529 - snprintfz(tag, NETDATA_THREAD_TAG_MAX, "TRAIN[%zu]", training_thread->id);
1530 - netdata_thread_create(&training_thread->nd_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_train_main, training_thread);
1531 - }
1532 -}
1533 -
1534 -void ml_stop_threads()
1535 -{
1536 - Cfg.detection_stop = true;
1537 - Cfg.training_stop = true;
1538 -
1539 - netdata_thread_cancel(Cfg.detection_thread);
1540 - netdata_thread_join(Cfg.detection_thread, NULL);
1541 -
1542 - // signal the training queue of each thread
1543 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1544 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1545 -
1546 - ml_queue_signal(training_thread->training_queue);
1547 - }
1548 -
1549 - // cancel training threads
1550 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1551 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1552 -
1553 - netdata_thread_cancel(training_thread->nd_thread);
1554 - }
1555 -
1556 - // join training threads
1557 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1558 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1559 -
1560 - netdata_thread_join(training_thread->nd_thread, NULL);
1561 - }
1562 -
1563 - // clear training thread data
1564 - for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1565 - ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1566 -
1567 - delete[] training_thread->training_cns;
1568 - delete[] training_thread->scratch_training_cns;
1569 - ml_queue_destroy(training_thread->training_queue);
1570 - netdata_mutex_destroy(&training_thread->nd_mutex);
1571 - }
1572 -}
ml/ml.h
+4 -7
@@ -13,12 +13,7 @@ extern "C" {
13 bool ml_capable();
14 bool ml_enabled(RRDHOST *rh);
15 bool ml_streaming_enabled();
16 -
16 void ml_init(void);
18 -void ml_fini(void);
19 -
20 -void ml_start_threads(void);
21 -void ml_stop_threads(void);
17
18 void ml_host_new(RRDHOST *rh);
19 void ml_host_delete(RRDHOST *rh);
@@ -27,6 +22,10 @@ void ml_host_get_info(RRDHOST *RH, BUFFER *wb);
22 void ml_host_get_detection_info(RRDHOST *RH, BUFFER *wb);
23 void ml_host_get_models(RRDHOST *RH, BUFFER *wb);
24
25 +void ml_host_start_training_thread(RRDHOST *rh);
26 +void ml_host_cancel_training_thread(RRDHOST *rh);
27 +void ml_host_stop_training_thread(RRDHOST *rh);
28 +
29 void ml_chart_new(RRDSET *rs);
30 void ml_chart_delete(RRDSET *rs);
31 bool ml_chart_update_begin(RRDSET *rs);
@@ -36,8 +35,6 @@ void ml_dimension_new(RRDDIM *rd);
35 void ml_dimension_delete(RRDDIM *rd);
36 bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool exists);
37
39 -void ml_update_global_statistics_charts(uint64_t models_consulted);
40 -
38 #ifdef __cplusplus
39 };
40 #endif