@cryptotaxi247 / netdata-1 / commits / ba2a5e785

Skip ML initialization when it's been disabled in netdata.conf (#14920)

* Apply ML changes again. The ML changes in 003df5f2 wheere reverted with 556bdad9 because we were partially initializing ML even when it was explicitly disabled in netdata.conf, causing the agent to crash on startup. * Do not start/stop ML threads when ML is disabled. * Restore default config settings.

vkalintiris committed Apr 20, 2023 at 11:24 UTC ba2a5e7857435cdcd217d5b7d4f78bd51ae6a6d1
11 files changed +783 -431
daemon/global_statistics.c
+1 -27
@@ -827,33 +827,7 @@ static void global_statistics_charts(void) {
827 rrdset_done(st_points_stored);
828 }
829
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 - }
830 + ml_update_global_statistics_charts(gs.ml_models_consulted);
831 }
832
833 // ----------------------------------------------------------------------------
daemon/main.c
+9 -9
@@ -148,10 +148,6 @@ 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 ");
151 if(service & SERVICE_REPLICATION)
152 buffer_strcat(wb, "REPLICATION ");
153 if(service & ABILITY_DATA_QUERIES)
@@ -340,6 +336,11 @@ void netdata_cleanup_and_exit(int ret) {
336 }
337 #endif
338
339 + delta_shutdown_time("disable ML detection and training threads");
340 +
341 + ml_stop_threads();
342 + ml_fini();
343 +
344 delta_shutdown_time("disable maintenance, new queries, new web requests, new streaming connections and aclk");
345
346 service_signal_exit(
@@ -351,12 +352,11 @@ void netdata_cleanup_and_exit(int ret) {
352 | SERVICE_ACLKSYNC
353 );
354
354 - delta_shutdown_time("stop replication, exporters, ML training, health and web servers threads");
355 + delta_shutdown_time("stop replication, exporters, health and web servers threads");
356
357 timeout = !service_wait_exit(
358 SERVICE_REPLICATION
359 | SERVICE_EXPORTERS
359 - | SERVICE_ML_TRAINING
360 | SERVICE_HEALTH
361 | SERVICE_WEB_SERVER
362 , 3 * USEC_PER_SEC);
@@ -368,11 +368,10 @@ void netdata_cleanup_and_exit(int ret) {
368 | SERVICE_STREAMING
369 , 3 * USEC_PER_SEC);
370
371 - delta_shutdown_time("stop ML prediction and context threads");
371 + delta_shutdown_time("stop context thread");
372
373 timeout = !service_wait_exit(
374 - SERVICE_ML_PREDICTION
375 - | SERVICE_CONTEXT
374 + SERVICE_CONTEXT
375 , 3 * USEC_PER_SEC);
376
377 delta_shutdown_time("stop maintenance thread");
@@ -2085,6 +2084,7 @@ int main(int argc, char **argv) {
2084 }
2085 else debug(D_SYSTEM, "Not starting thread %s.", st->name);
2086 }
2087 + ml_start_threads();
2088
2089 // ------------------------------------------------------------------------
2090 // Initialize netdata agent command serving from cli and signals
daemon/main.h
+9 -11
@@ -33,17 +33,15 @@ typedef enum {
33 ABILITY_STREAMING_CONNECTIONS = (1 << 2),
34 SERVICE_MAINTENANCE = (1 << 3),
35 SERVICE_COLLECTORS = (1 << 4),
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)
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)
45 } SERVICE_TYPE;
46
47 typedef enum {
database/rrdhost.c
-3
@@ -524,7 +524,6 @@ 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);
527 } else
528 rrdhost_flag_set(host, RRDHOST_FLAG_PENDING_CONTEXT_LOAD | RRDHOST_FLAG_ARCHIVED | RRDHOST_FLAG_ORPHAN);
529
@@ -641,7 +640,6 @@ static void rrdhost_update(RRDHOST *host
640 host->rrdpush_replication_step = rrdpush_replication_step;
641
642 ml_host_new(host);
644 - ml_host_start_training_thread(host);
643
644 rrdhost_load_rrdcontext_data(host);
645 info("Host %s is not in archived mode anymore", rrdhost_hostname(host));
@@ -1143,7 +1141,6 @@ void rrdhost_free___while_having_rrd_wrlock(RRDHOST *host, bool force) {
1141 rrdcalctemplate_index_destroy(host);
1142
1143 // cleanup ML resources
1146 - ml_host_stop_training_thread(host);
1144 ml_host_delete(host);
1145
1146 freez(host->exporting_flags);
ml/Config.cc
+11 -1
@@ -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 / lag_n);
37 + double random_sampling_ratio = config_get_float(config_section_ml, "random sampling ratio", 1.0 / 5.0 /* default 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,6 +43,10 @@ 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 +
50 /*
51 * Clamp
52 */
@@ -64,6 +68,8 @@ void ml_config_load(ml_config_t *cfg) {
68 host_anomaly_rate_threshold = clamp(host_anomaly_rate_threshold, 0.1, 10.0);
69 anomaly_detection_query_duration = clamp<time_t>(anomaly_detection_query_duration, 60, 15 * 60);
70
71 + num_training_threads = clamp<size_t>(num_training_threads, 1, 128);
72 +
73 /*
74 * Validate
75 */
@@ -109,4 +115,8 @@ void ml_config_load(ml_config_t *cfg) {
115 cfg->sp_charts_to_skip = simple_pattern_create(cfg->charts_to_skip.c_str(), NULL, SIMPLE_PATTERN_EXACT, true);
116
117 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;
122 }
ml/ad_charts.cc
+98 -69
@@ -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 - {
9 + if (Cfg.enable_statistics_charts) {
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 - {
51 + if (Cfg.enable_statistics_charts) {
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 - {
93 + if (Cfg.enable_statistics_charts) {
94 if (!host->training_status_rs) {
95 char id_buf[1024];
96 char name_buf[1024];
@@ -179,7 +179,6 @@ void ml_update_dimensions_chart(ml_host_t *host, const ml_machine_learning_stats
179
180 rrdset_done(host->dimensions_rs);
181 }
182 -
182 }
183
184 void ml_update_host_and_detection_rate_charts(ml_host_t *host, collected_number AnomalyRate) {
@@ -301,20 +300,20 @@ void ml_update_host_and_detection_rate_charts(ml_host_t *host, collected_number
300 }
301 }
302
304 -void ml_update_training_statistics_chart(ml_host_t *host, const ml_training_stats_t &ts) {
303 +void ml_update_training_statistics_chart(ml_training_thread_t *training_thread, const ml_training_stats_t &ts) {
304 /*
305 * queue stats
306 */
307 {
309 - if (!host->queue_stats_rs) {
308 + if (!training_thread->queue_stats_rs) {
309 char id_buf[1024];
310 char name_buf[1024];
311
313 - snprintfz(id_buf, 1024, "queue_stats_on_%s", localhost->machine_guid);
314 - snprintfz(name_buf, 1024, "queue_stats_on_%s", rrdhost_hostname(localhost));
312 + snprintfz(id_buf, 1024, "training_queue_%zu_stats", training_thread->id);
313 + snprintfz(name_buf, 1024, "training_queue_%zu_stats", training_thread->id);
314
316 - host->queue_stats_rs = rrdset_create(
317 - host->rh,
315 + training_thread->queue_stats_rs = rrdset_create(
316 + localhost,
317 "netdata", // type
318 id_buf, // id
319 name_buf, // name
@@ -328,35 +327,35 @@ void ml_update_training_statistics_chart(ml_host_t *host, const ml_training_stat
327 localhost->rrd_update_every, // update_every
328 RRDSET_TYPE_LINE// chart_type
329 );
331 - rrdset_flag_set(host->queue_stats_rs, RRDSET_FLAG_ANOMALY_DETECTION);
330 + rrdset_flag_set(training_thread->queue_stats_rs, RRDSET_FLAG_ANOMALY_DETECTION);
331
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);
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);
336 }
337
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);
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);
342
344 - rrdset_done(host->queue_stats_rs);
343 + rrdset_done(training_thread->queue_stats_rs);
344 }
345
346 /*
347 * training stats
348 */
349 {
351 - if (!host->training_time_stats_rs) {
350 + if (!training_thread->training_time_stats_rs) {
351 char id_buf[1024];
352 char name_buf[1024];
353
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));
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);
356
358 - host->training_time_stats_rs = rrdset_create(
359 - host->rh,
357 + training_thread->training_time_stats_rs = rrdset_create(
358 + localhost,
359 "netdata", // type
360 id_buf, // id
361 name_buf, // name
@@ -370,39 +369,39 @@ void ml_update_training_statistics_chart(ml_host_t *host, const ml_training_stat
369 localhost->rrd_update_every, // update_every
370 RRDSET_TYPE_LINE// chart_type
371 );
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);
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);
380 }
381
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);
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);
388
390 - rrdset_done(host->training_time_stats_rs);
389 + rrdset_done(training_thread->training_time_stats_rs);
390 }
391
392 /*
393 * training result stats
394 */
395 {
397 - if (!host->training_results_rs) {
396 + if (!training_thread->training_results_rs) {
397 char id_buf[1024];
398 char name_buf[1024];
399
401 - snprintfz(id_buf, 1024, "training_results_on_%s", localhost->machine_guid);
402 - snprintfz(name_buf, 1024, "training_results_on_%s", rrdhost_hostname(localhost));
400 + snprintfz(id_buf, 1024, "training_queue_%zu_results", training_thread->id);
401 + snprintfz(name_buf, 1024, "training_queue_%zu_results", training_thread->id);
402
404 - host->training_results_rs = rrdset_create(
405 - host->rh,
403 + training_thread->training_results_rs = rrdset_create(
404 + localhost,
405 "netdata", // type
406 id_buf, // id
407 name_buf, // name
@@ -416,31 +415,61 @@ void ml_update_training_statistics_chart(ml_host_t *host, const ml_training_stat
415 localhost->rrd_update_every, // update_every
416 RRDSET_TYPE_LINE// chart_type
417 );
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);
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);
430 }
431
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);
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);
474 }
475 }
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_host_t *host, const ml_training_stats_t &ts);
12 +void ml_update_training_statistics_chart(ml_training_thread_t *training_thread, const ml_training_stats_t &ts);
13
14 #endif /* ML_ADCHARTS_H */
ml/ml-dummy.c
+10
@@ -19,6 +19,12 @@ 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 +
28 void ml_host_new(RRDHOST *rh) {
29 UNUSED(rh);
30 }
@@ -86,4 +92,8 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
92 return false;
93 }
94
95 +void ml_update_global_statistics_charts(uint64_t models_consulted) {
96 + UNUSED(models_consulted);
97 +}
98 +
99 #endif
ml/ml-private.h
+29 -13
@@ -33,14 +33,15 @@ typedef struct {
33 /*
34 * KMeans
35 */
36 -typedef struct {
37 - size_t num_clusters;
38 - size_t max_iterations;
36
37 +typedef struct {
38 std::vector<DSample> cluster_centers;
39
40 calculated_number_t min_dist;
41 calculated_number_t max_dist;
42 +
43 + uint32_t after;
44 + uint32_t before;
45 } ml_kmeans_t;
46
47 typedef struct machine_learning_stats_t {
@@ -123,6 +124,7 @@ enum ml_training_result {
124
125 typedef struct {
126 // Chart/dimension we want to train
127 + STRING *host_id;
128 STRING *chart_id;
129 STRING *dimension_id;
130
@@ -168,6 +170,7 @@ typedef struct {
170 /*
171 * Queue
172 */
173 +
174 typedef struct {
175 std::queue<ml_training_request_t> internal;
176 netdata_mutex_t mutex;
@@ -175,7 +178,6 @@ typedef struct {
178 std::atomic<bool> exit;
179 } ml_queue_t;
180
178 -
181 typedef struct {
182 RRDDIM *rd;
183
@@ -207,19 +209,12 @@ typedef struct {
209 RRDHOST *rh;
210
211 ml_machine_learning_stats_t mls;
210 - ml_training_stats_t ts;
212
213 calculated_number_t host_anomaly_rate;
214
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 -
215 netdata_mutex_t mutex;
216
222 - netdata_thread_t training_thread;
217 + ml_queue_t *training_queue;
218
219 /*
220 * bookkeeping for anomaly detection charts
@@ -249,6 +244,19 @@ typedef struct {
244 RRDSET *detector_events_rs;
245 RRDDIM *detector_events_above_threshold_rd;
246 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;
260
261 RRDSET *queue_stats_rs;
262 RRDDIM *queue_stats_queue_size_rd;
@@ -265,7 +273,7 @@ typedef struct {
273 RRDDIM *training_results_not_enough_collected_values_rd;
274 RRDDIM *training_results_null_acquired_dimension_rd;
275 RRDDIM *training_results_chart_under_replication_rd;
268 -} ml_host_t;
276 +} ml_training_thread_t;
277
278 typedef struct {
279 bool enable_anomaly_detection;
@@ -302,6 +310,14 @@ typedef struct {
310 std::vector<uint32_t> random_nums;
311
312 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;
321 } ml_config_t;
322
323 void ml_config_load(ml_config_t *cfg);
ml/ml.cc
+608 -293
@@ -7,15 +7,18 @@
7 #include <random>
8
9 #include "ad_charts.h"
10 +#include "database/sqlite/sqlite3.h"
11
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;
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
20
18 -static thread_local ml_tls_data_t tls_data;
21 +static sqlite3 *db = NULL;
22
23 /*
24 * Functions to convert enums to strings
@@ -173,26 +176,26 @@ ml_features_preprocess(ml_features_t *features)
176 */
177
178 static void
176 -ml_kmeans_init(ml_kmeans_t *kmeans, size_t num_clusters, size_t max_iterations)
179 +ml_kmeans_init(ml_kmeans_t *kmeans)
180 {
178 - kmeans->num_clusters = num_clusters;
179 - kmeans->max_iterations = max_iterations;
180 -
181 - kmeans->cluster_centers.reserve(kmeans->num_clusters);
181 + kmeans->cluster_centers.reserve(2);
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)
187 +ml_kmeans_train(ml_kmeans_t *kmeans, const ml_features_t *features, time_t after, time_t before)
188 {
189 + kmeans->after = (uint32_t) after;
190 + kmeans->before = (uint32_t) before;
191 +
192 kmeans->min_dist = std::numeric_limits<calculated_number_t>::max();
193 kmeans->max_dist = std::numeric_limits<calculated_number_t>::min();
194
195 kmeans->cluster_centers.clear();
196
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);
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);
199
200 for (const auto &preprocessed_feature : features->preprocessed_features) {
201 calculated_number_t mean_dist = 0.0;
@@ -201,7 +204,7 @@ ml_kmeans_train(ml_kmeans_t *kmeans, const ml_features_t *features)
204 mean_dist += dlib::length(cluster_center - preprocessed_feature);
205 }
206
204 - mean_dist /= kmeans->num_clusters;
207 + mean_dist /= kmeans->cluster_centers.size();
208
209 if (mean_dist < kmeans->min_dist)
210 kmeans->min_dist = mean_dist;
@@ -218,7 +221,7 @@ ml_kmeans_anomaly_score(const ml_kmeans_t *kmeans, const DSample &DS)
221 for (const auto &CC: kmeans->cluster_centers)
222 mean_dist += dlib::length(CC - DS);
223
221 - mean_dist /= kmeans->num_clusters;
224 + mean_dist /= kmeans->cluster_centers.size();
225
226 if (kmeans->max_dist == kmeans->min_dist)
227 return 0.0;
@@ -264,7 +267,14 @@ ml_queue_pop(ml_queue_t *q)
267 {
268 netdata_mutex_lock(&q->mutex);
269
267 - ml_training_request_t req = { NULL, NULL, 0, 0, 0 };
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 + };
278
279 while (q->internal.empty()) {
280 pthread_cond_wait(&q->cond_var, &q->mutex);
@@ -307,7 +317,7 @@ ml_queue_signal(ml_queue_t *q)
317 */
318
319 static std::pair<calculated_number_t *, ml_training_response_t>
310 -ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t &training_request)
320 +ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
321 {
322 ml_training_response_t training_response = {};
323
@@ -348,7 +358,7 @@ ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t
358 STORAGE_PRIORITY_BEST_EFFORT);
359
360 size_t idx = 0;
351 - memset(tls_data.training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
361 + memset(training_thread->training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
362 calculated_number_t last_value = std::numeric_limits<calculated_number_t>::quiet_NaN();
363
364 while (!storage_engine_query_is_finished(&handle)) {
@@ -365,11 +375,11 @@ ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t
375 training_response.db_after_t = timestamp;
376 training_response.db_before_t = timestamp;
377
368 - tls_data.training_cns[idx] = value;
369 - last_value = tls_data.training_cns[idx];
378 + training_thread->training_cns[idx] = value;
379 + last_value = training_thread->training_cns[idx];
380 training_response.collected_values++;
381 } else
372 - tls_data.training_cns[idx] = last_value;
382 + training_thread->training_cns[idx] = last_value;
383
384 idx++;
385 }
@@ -384,20 +394,270 @@ ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t
394 }
395
396 // Find first non-NaN value.
387 - for (idx = 0; std::isnan(tls_data.training_cns[idx]); idx++, training_response.total_values--) { }
397 + for (idx = 0; std::isnan(training_thread->training_cns[idx]); idx++, training_response.total_values--) { }
398
399 // Overwrite NaN values.
400 if (idx != 0)
391 - memmove(tls_data.training_cns, &tls_data.training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
401 + memmove(training_thread->training_cns, &training_thread->training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
402
403 training_response.result = TRAINING_RESULT_OK;
394 - return { tls_data.training_cns, training_response };
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);
654 }
655
656 static enum ml_training_result
398 -ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &training_request)
657 +ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
658 {
400 - auto P = ml_dimension_calculated_numbers(dim, training_request);
659 + worker_is_busy(WORKER_TRAIN_QUERY);
660 + auto P = ml_dimension_calculated_numbers(training_thread, dim, training_request);
661 ml_training_response_t training_response = P.second;
662
663 if (training_response.result != TRAINING_RESULT_OK) {
@@ -426,31 +686,56 @@ ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &train
686 }
687
688 // compute kmeans
689 + worker_is_busy(WORKER_TRAIN_KMEANS);
690 {
430 - memcpy(tls_data.scratch_training_cns, tls_data.training_cns,
691 + memcpy(training_thread->scratch_training_cns, training_thread->training_cns,
692 training_response.total_values * sizeof(calculated_number_t));
693
694 ml_features_t features = {
695 Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
435 - tls_data.scratch_training_cns, training_response.total_values,
436 - tls_data.training_cns, training_response.total_values,
437 - tls_data.training_samples
696 + training_thread->scratch_training_cns, training_response.total_values,
697 + training_thread->training_cns, training_response.total_values,
698 + training_thread->training_samples
699 };
700 ml_features_preprocess(&features);
701
441 - ml_kmeans_init(&dim->kmeans, 2, 1000);
442 - ml_kmeans_train(&dim->kmeans, &features);
702 + ml_kmeans_init(&dim->kmeans);
703 + ml_kmeans_train(&dim->kmeans, &features, training_response.query_after_t, training_response.query_before_t);
704 }
705
445 - // update kmeans models
706 + // update models
707 {
708 netdata_mutex_lock(&dim->mutex);
709
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 +
719 if (dim->km_contexts.size() < Cfg.num_models_to_use) {
720 dim->km_contexts.push_back(std::move(dim->kmeans));
721 } else {
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);
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 + }
739 }
740
741 dim->mt = METRIC_TYPE_CONSTANT;
@@ -494,11 +779,16 @@ ml_dimension_schedule_for_training(ml_dimension_t *dim, time_t curr_time)
779 }
780
781 if (schedule_for_training) {
497 - ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
782 ml_training_request_t req = {
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),
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),
789 };
790 +
791 + ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
792 ml_queue_push(host->training_queue, req);
793 }
794 }
@@ -674,7 +964,6 @@ ml_host_detect_once(ml_host_t *host)
964
965 host->mls = {};
966 ml_machine_learning_stats_t mls_copy = {};
677 - ml_training_stats_t ts_copy = {};
967
968 {
969 netdata_mutex_lock(&host->mutex);
@@ -718,54 +1007,14 @@ ml_host_detect_once(ml_host_t *host)
1007
1008 mls_copy = host->mls;
1009
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 -
1010 netdata_mutex_unlock(&host->mutex);
1011 }
1012
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 -
1013 worker_is_busy(WORKER_JOB_DETECTION_DIM_CHART);
1014 ml_update_dimensions_chart(host, mls_copy);
1015
1016 worker_is_busy(WORKER_JOB_DETECTION_HOST_CHART);
1017 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);
1018 }
1019
1020 typedef struct {
@@ -774,18 +1023,21 @@ typedef struct {
1023 } ml_acquired_dimension_t;
1024
1025 static ml_acquired_dimension_t
777 -ml_acquired_dimension_get(RRDHOST *rh, STRING *chart_id, STRING *dimension_id)
1026 +ml_acquired_dimension_get(STRING *host_id, STRING *chart_id, STRING *dimension_id)
1027 {
1028 RRDDIM_ACQUIRED *acq_rd = NULL;
1029 ml_dimension_t *dim = NULL;
1030
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;
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 + }
1041 }
1042 }
1043
@@ -806,110 +1058,12 @@ ml_acquired_dimension_release(ml_acquired_dimension_t acq_dim)
1058 }
1059
1060 static enum ml_training_result
809 -ml_acquired_dimension_train(ml_acquired_dimension_t acq_dim, const ml_training_request_t &TR)
1061 +ml_acquired_dimension_train(ml_training_thread_t *training_thread, ml_acquired_dimension_t acq_dim, const ml_training_request_t &tr)
1062 {
1063 if (!acq_dim.dim)
1064 return TRAINING_RESULT_NULL_ACQUIRED_DIMENSION;
1065
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;
1066 + return ml_dimension_train_model(training_thread, acq_dim.dim, tr);
1067 }
1068
1069 static void *
@@ -923,25 +1077,55 @@ ml_detect_main(void *arg)
1077 worker_register_job_name(WORKER_JOB_DETECTION_HOST_CHART, "host chart");
1078 worker_register_job_name(WORKER_JOB_DETECTION_STATS, "training stats");
1079
926 - service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, NULL, NULL, true);
927 -
1080 heartbeat_t hb;
1081 heartbeat_init(&hb);
1082
931 - while (service_running((SERVICE_TYPE)(SERVICE_ML_PREDICTION | SERVICE_COLLECTORS))) {
1083 + while (!Cfg.detection_stop) {
1084 worker_is_idle();
1085 heartbeat_next(&hb, USEC_PER_SEC);
1086
935 - void *rhp;
936 - dfe_start_reentrant(rrdhost_root_index, rhp) {
937 - RRDHOST *rh = (RRDHOST *) rhp;
938 -
1087 + RRDHOST *rh;
1088 + rrd_rdlock();
1089 + rrdhost_foreach_read(rh) {
1090 if (!rh->ml_host)
1091 continue;
1092
1093 ml_host_detect_once((ml_host_t *) rh->ml_host);
1094 }
944 - dfe_done(rhp);
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 + }
1129 }
1130
1131 return NULL;
@@ -975,31 +1159,6 @@ bool ml_streaming_enabled()
1159 return Cfg.stream_anomaly_detection_charts;
1160 }
1161
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 -
1162 void ml_host_new(RRDHOST *rh)
1163 {
1164 if (!ml_enabled(rh))
@@ -1009,14 +1168,12 @@ void ml_host_new(RRDHOST *rh)
1168
1169 host->rh = rh;
1170 host->mls = ml_machine_learning_stats_t();
1012 - host->ts = ml_training_stats_t();
1171 + //host->ts = ml_training_stats_t();
1172
1014 - host->host_anomaly_rate = 0.0;
1015 - host->threads_running = false;
1016 - host->threads_cancelled = false;
1017 - host->threads_joined = false;
1173 + static std::atomic<size_t> times_called(0);
1174 + host->training_queue = Cfg.training_threads[times_called++ % Cfg.num_training_threads].training_queue;
1175
1019 - host->training_queue = ml_queue_init();
1176 + host->host_anomaly_rate = 0.0;
1177
1178 netdata_mutex_init(&host->mutex);
1179
@@ -1030,7 +1187,6 @@ void ml_host_delete(RRDHOST *rh)
1187 return;
1188
1189 netdata_mutex_destroy(&host->mutex);
1033 - ml_queue_destroy(host->training_queue);
1190
1191 delete host;
1192 rh->ml_host = NULL;
@@ -1097,69 +1253,6 @@ void ml_host_get_models(RRDHOST *rh, BUFFER *wb)
1253 error("Fetching KMeans models is not supported yet");
1254 }
1255
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 -
1256 void ml_chart_new(RRDSET *rs)
1257 {
1258 ml_host_t *host = (ml_host_t *) rs->rrdhost->ml_host;
@@ -1225,7 +1318,7 @@ void ml_dimension_new(RRDDIM *rd)
1318
1319 dim->last_training_time = 0;
1320
1228 - ml_kmeans_init(&dim->kmeans, 2, 1000);
1321 + ml_kmeans_init(&dim->kmeans);
1322
1323 if (simple_pattern_matches(Cfg.sp_charts_to_skip, rrdset_name(rd->rrdset)))
1324 dim->mls = MACHINE_LEARNING_STATUS_DISABLED_DUE_TO_EXCLUDED_CHART;
@@ -1264,3 +1357,225 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
1357
1358 return is_anomalous;
1359 }
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 + if (!Cfg.enable_anomaly_detection)
1513 + return;
1514 +
1515 + int rc = sqlite3_close_v2(db);
1516 + if (unlikely(rc != SQLITE_OK))
1517 + error_report("Error %d while closing the SQLite database, %s", rc, sqlite3_errstr(rc));
1518 +}
1519 +
1520 +void ml_start_threads() {
1521 + if (!Cfg.enable_anomaly_detection)
1522 + return;
1523 +
1524 + // start detection & training threads
1525 + Cfg.detection_stop = false;
1526 + Cfg.training_stop = false;
1527 +
1528 + char tag[NETDATA_THREAD_TAG_MAX + 1];
1529 +
1530 + snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "PREDICT");
1531 + netdata_thread_create(&Cfg.detection_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_detect_main, NULL);
1532 +
1533 + for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1534 + ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1535 + snprintfz(tag, NETDATA_THREAD_TAG_MAX, "TRAIN[%zu]", training_thread->id);
1536 + netdata_thread_create(&training_thread->nd_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_train_main, training_thread);
1537 + }
1538 +}
1539 +
1540 +void ml_stop_threads()
1541 +{
1542 + if (!Cfg.enable_anomaly_detection)
1543 + return;
1544 +
1545 + Cfg.detection_stop = true;
1546 + Cfg.training_stop = true;
1547 +
1548 + netdata_thread_cancel(Cfg.detection_thread);
1549 + netdata_thread_join(Cfg.detection_thread, NULL);
1550 +
1551 + // signal the training queue of each thread
1552 + for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1553 + ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1554 +
1555 + ml_queue_signal(training_thread->training_queue);
1556 + }
1557 +
1558 + // cancel training threads
1559 + for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1560 + ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1561 +
1562 + netdata_thread_cancel(training_thread->nd_thread);
1563 + }
1564 +
1565 + // join training threads
1566 + for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1567 + ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1568 +
1569 + netdata_thread_join(training_thread->nd_thread, NULL);
1570 + }
1571 +
1572 + // clear training thread data
1573 + for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1574 + ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1575 +
1576 + delete[] training_thread->training_cns;
1577 + delete[] training_thread->scratch_training_cns;
1578 + ml_queue_destroy(training_thread->training_queue);
1579 + netdata_mutex_destroy(&training_thread->nd_mutex);
1580 + }
1581 +}
ml/ml.h
+7 -4
@@ -13,7 +13,12 @@ extern "C" {
13 bool ml_capable();
14 bool ml_enabled(RRDHOST *rh);
15 bool ml_streaming_enabled();
16 +
17 void ml_init(void);
18 +void ml_fini(void);
19 +
20 +void ml_start_threads(void);
21 +void ml_stop_threads(void);
22
23 void ml_host_new(RRDHOST *rh);
24 void ml_host_delete(RRDHOST *rh);
@@ -22,10 +27,6 @@ void ml_host_get_info(RRDHOST *RH, BUFFER *wb);
27 void ml_host_get_detection_info(RRDHOST *RH, BUFFER *wb);
28 void ml_host_get_models(RRDHOST *RH, BUFFER *wb);
29
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 -
30 void ml_chart_new(RRDSET *rs);
31 void ml_chart_delete(RRDSET *rs);
32 bool ml_chart_update_begin(RRDSET *rs);
@@ -35,6 +36,8 @@ void ml_dimension_new(RRDDIM *rd);
36 void ml_dimension_delete(RRDDIM *rd);
37 bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool exists);
38
39 +void ml_update_global_statistics_charts(uint64_t models_consulted);
40 +
41 #ifdef __cplusplus
42 };
43 #endif