Revert "Use static thread-pool for training. (#14702)" (#14782)
This reverts commit 5046e034212c008557dd014196b6f6204eda24b2. Will re-apply once we investigate an issue that occurs during the shutdown of the agent.
vkalintiris committed
Mar 21, 2023 at 18:31 UTC
5321ca8d1ef8d974a6a2b2128ca8804de6acb693
11 files changed
+408
-436
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
-7
@@ -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,10 +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_fini();
342
-
343
delta_shutdown_time("disable maintenance, new queries, new web requests, new streaming connections and aclk");
344
345
service_signal_exit(
@@ -351,11 +351,12 @@ void netdata_cleanup_and_exit(int ret) {
351
| SERVICE_ACLKSYNC
352
);
353
354
- 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);
@@ -367,10 +368,11 @@ void netdata_cleanup_and_exit(int ret) {
368
| SERVICE_STREAMING
369
, 3 * USEC_PER_SEC);
370
370
- delta_shutdown_time("stop context thread");
371
+ delta_shutdown_time("stop ML prediction and context threads");
372
373
timeout = !service_wait_exit(
373
- SERVICE_CONTEXT
374
+ SERVICE_ML_PREDICTION
375
+ | SERVICE_CONTEXT
376
, 3 * USEC_PER_SEC);
377
378
delta_shutdown_time("stop maintenance thread");
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
-6
@@ -19,8 +19,6 @@ bool ml_streaming_enabled() {
19
20
void ml_init(void) {}
21
22
-void ml_fini(void) {}
23
-
22
void ml_host_new(RRDHOST *rh) {
23
UNUSED(rh);
24
}
@@ -88,8 +86,4 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
86
return false;
87
}
88
91
-void ml_update_global_statistics_charts(uint64_t models_consulted) {
92
- UNUSED(models_consulted);
93
-}
94
-
89
#endif
ml/ml-private.h
+10
-26
@@ -33,7 +33,6 @@ typedef struct {
33
/*
34
* KMeans
35
*/
36
-
36
typedef struct {
37
size_t num_clusters;
38
size_t max_iterations;
@@ -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
+273
-273
@@ -8,13 +8,14 @@
8
9
#include "ad_charts.h"
10
11
-#define WORKER_TRAIN_QUEUE_POP 0
12
-#define WORKER_TRAIN_ACQUIRE_DIMENSION 1
13
-#define WORKER_TRAIN_QUERY 2
14
-#define WORKER_TRAIN_KMEANS 3
15
-#define WORKER_TRAIN_UPDATE_MODELS 4
16
-#define WORKER_TRAIN_RELEASE_DIMENSION 5
17
-#define WORKER_TRAIN_UPDATE_HOST 6
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
+
18
+static thread_local ml_tls_data_t tls_data;
19
20
/*
21
* Functions to convert enums to strings
@@ -263,14 +264,7 @@ ml_queue_pop(ml_queue_t *q)
264
{
265
netdata_mutex_lock(&q->mutex);
266
266
- ml_training_request_t req = {
267
- NULL, // host_id
268
- NULL, // chart id
269
- NULL, // dimension id
270
- 0, // current time
271
- 0, // first entry
272
- 0 // last entry
273
- };
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);
@@ -313,7 +307,7 @@ ml_queue_signal(ml_queue_t *q)
307
*/
308
309
static std::pair<calculated_number_t *, ml_training_response_t>
316
-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
@@ -357,7 +351,7 @@ ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimens
351
STORAGE_PRIORITY_BEST_EFFORT);
352
353
size_t idx = 0;
360
- memset(training_thread->training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
354
+ memset(tls_data.training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
355
calculated_number_t last_value = std::numeric_limits<calculated_number_t>::quiet_NaN();
356
357
while (!ops->is_finished(&handle)) {
@@ -374,11 +368,11 @@ ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimens
368
training_response.db_after_t = timestamp;
369
training_response.db_before_t = timestamp;
370
377
- training_thread->training_cns[idx] = value;
378
- last_value = training_thread->training_cns[idx];
371
+ tls_data.training_cns[idx] = value;
372
+ last_value = tls_data.training_cns[idx];
373
training_response.collected_values++;
374
} else
381
- training_thread->training_cns[idx] = last_value;
375
+ tls_data.training_cns[idx] = last_value;
376
377
idx++;
378
}
@@ -393,21 +387,20 @@ ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimens
387
}
388
389
// Find first non-NaN value.
396
- for (idx = 0; std::isnan(training_thread->training_cns[idx]); idx++, training_response.total_values--) { }
390
+ for (idx = 0; std::isnan(tls_data.training_cns[idx]); idx++, training_response.total_values--) { }
391
392
// Overwrite NaN values.
393
if (idx != 0)
400
- memmove(training_thread->training_cns, &training_thread->training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
394
+ memmove(tls_data.training_cns, &tls_data.training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
395
396
training_response.result = TRAINING_RESULT_OK;
403
- return { training_thread->training_cns, training_response };
397
+ return { tls_data.training_cns, training_response };
398
}
399
400
static enum ml_training_result
407
-ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
401
+ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &training_request)
402
{
409
- worker_is_busy(WORKER_TRAIN_QUERY);
410
- auto P = ml_dimension_calculated_numbers(training_thread, dim, training_request);
403
+ auto P = ml_dimension_calculated_numbers(dim, training_request);
404
ml_training_response_t training_response = P.second;
405
406
if (training_response.result != TRAINING_RESULT_OK) {
@@ -436,16 +429,15 @@ ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *
429
}
430
431
// compute kmeans
439
- worker_is_busy(WORKER_TRAIN_KMEANS);
432
{
441
- memcpy(training_thread->scratch_training_cns, training_thread->training_cns,
433
+ memcpy(tls_data.scratch_training_cns, tls_data.training_cns,
434
training_response.total_values * sizeof(calculated_number_t));
435
436
ml_features_t features = {
437
Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
446
- training_thread->scratch_training_cns, training_response.total_values,
447
- training_thread->training_cns, training_response.total_values,
448
- training_thread->training_samples
438
+ tls_data.scratch_training_cns, training_response.total_values,
439
+ tls_data.training_cns, training_response.total_values,
440
+ tls_data.training_samples
441
};
442
ml_features_preprocess(&features);
443
@@ -454,7 +446,6 @@ ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *
446
}
447
448
// update kmeans models
457
- worker_is_busy(WORKER_TRAIN_UPDATE_MODELS);
449
{
450
netdata_mutex_lock(&dim->mutex);
451
@@ -506,16 +497,11 @@ ml_dimension_schedule_for_training(ml_dimension_t *dim, time_t curr_time)
497
}
498
499
if (schedule_for_training) {
500
+ ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
501
ml_training_request_t req = {
510
- string_dup(dim->rd->rrdset->rrdhost->hostname),
511
- string_dup(dim->rd->rrdset->id),
512
- string_dup(dim->rd->id),
513
- curr_time,
514
- rrddim_first_entry_s(dim->rd),
515
- rrddim_last_entry_s(dim->rd),
502
+ string_dup(dim->rd->rrdset->id), string_dup(dim->rd->id),
503
+ curr_time, rrddim_first_entry_s(dim->rd), rrddim_last_entry_s(dim->rd),
504
};
517
-
518
- ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
505
ml_queue_push(host->training_queue, req);
506
}
507
}
@@ -691,6 +677,7 @@ ml_host_detect_once(ml_host_t *host)
677
678
host->mls = {};
679
ml_machine_learning_stats_t mls_copy = {};
680
+ ml_training_stats_t ts_copy = {};
681
682
{
683
netdata_mutex_lock(&host->mutex);
@@ -734,14 +721,54 @@ ml_host_detect_once(ml_host_t *host)
721
722
mls_copy = host->mls;
723
724
+ /*
725
+ * training stats
726
+ */
727
+ ts_copy = host->ts;
728
+
729
+ host->ts.queue_size = 0;
730
+ host->ts.num_popped_items = 0;
731
+
732
+ host->ts.allotted_ut = 0;
733
+ host->ts.consumed_ut = 0;
734
+ host->ts.remaining_ut = 0;
735
+
736
+ host->ts.training_result_ok = 0;
737
+ host->ts.training_result_invalid_query_time_range = 0;
738
+ host->ts.training_result_not_enough_collected_values = 0;
739
+ host->ts.training_result_null_acquired_dimension = 0;
740
+ host->ts.training_result_chart_under_replication = 0;
741
+
742
netdata_mutex_unlock(&host->mutex);
743
}
744
745
+ // Calc the avg values
746
+ if (ts_copy.num_popped_items) {
747
+ ts_copy.queue_size /= ts_copy.num_popped_items;
748
+ ts_copy.allotted_ut /= ts_copy.num_popped_items;
749
+ ts_copy.consumed_ut /= ts_copy.num_popped_items;
750
+ ts_copy.remaining_ut /= ts_copy.num_popped_items;
751
+
752
+ ts_copy.training_result_ok /= ts_copy.num_popped_items;
753
+ ts_copy.training_result_invalid_query_time_range /= ts_copy.num_popped_items;
754
+ ts_copy.training_result_not_enough_collected_values /= ts_copy.num_popped_items;
755
+ ts_copy.training_result_null_acquired_dimension /= ts_copy.num_popped_items;
756
+ ts_copy.training_result_chart_under_replication /= ts_copy.num_popped_items;
757
+ } else {
758
+ ts_copy.queue_size = 0;
759
+ ts_copy.allotted_ut = 0;
760
+ ts_copy.consumed_ut = 0;
761
+ ts_copy.remaining_ut = 0;
762
+ }
763
+
764
worker_is_busy(WORKER_JOB_DETECTION_DIM_CHART);
765
ml_update_dimensions_chart(host, mls_copy);
766
767
worker_is_busy(WORKER_JOB_DETECTION_HOST_CHART);
768
ml_update_host_and_detection_rate_charts(host, host->host_anomaly_rate * 10000.0);
769
+
770
+ worker_is_busy(WORKER_JOB_DETECTION_STATS);
771
+ ml_update_training_statistics_chart(host, ts_copy);
772
}
773
774
typedef struct {
@@ -750,21 +777,18 @@ typedef struct {
777
} ml_acquired_dimension_t;
778
779
static ml_acquired_dimension_t
753
-ml_acquired_dimension_get(STRING *host_id, STRING *chart_id, STRING *dimension_id)
780
+ml_acquired_dimension_get(RRDHOST *rh, STRING *chart_id, STRING *dimension_id)
781
{
782
RRDDIM_ACQUIRED *acq_rd = NULL;
783
ml_dimension_t *dim = NULL;
784
758
- RRDHOST *rh = rrdhost_find_by_hostname(string2str(host_id));
759
- if (rh) {
760
- RRDSET *rs = rrdset_find(rh, string2str(chart_id));
761
- if (rs) {
762
- acq_rd = rrddim_find_and_acquire(rs, string2str(dimension_id));
763
- if (acq_rd) {
764
- RRDDIM *rd = rrddim_acquired_to_rrddim(acq_rd);
765
- if (rd)
766
- dim = (ml_dimension_t *) rd->ml_dimension;
767
- }
785
+ RRDSET *rs = rrdset_find(rh, string2str(chart_id));
786
+ if (rs) {
787
+ acq_rd = rrddim_find_and_acquire(rs, string2str(dimension_id));
788
+ if (acq_rd) {
789
+ RRDDIM *rd = rrddim_acquired_to_rrddim(acq_rd);
790
+ if (rd)
791
+ dim = (ml_dimension_t *) rd->ml_dimension;
792
}
793
}
794
@@ -785,12 +809,110 @@ ml_acquired_dimension_release(ml_acquired_dimension_t acq_dim)
809
}
810
811
static enum ml_training_result
788
-ml_acquired_dimension_train(ml_training_thread_t *training_thread, ml_acquired_dimension_t acq_dim, const ml_training_request_t &tr)
812
+ml_acquired_dimension_train(ml_acquired_dimension_t acq_dim, const ml_training_request_t &TR)
813
{
814
if (!acq_dim.dim)
815
return TRAINING_RESULT_NULL_ACQUIRED_DIMENSION;
816
793
- return ml_dimension_train_model(training_thread, acq_dim.dim, tr);
817
+ return ml_dimension_train_model(acq_dim.dim, TR);
818
+}
819
+
820
+#define WORKER_JOB_TRAINING_FIND 0
821
+#define WORKER_JOB_TRAINING_TRAIN 1
822
+#define WORKER_JOB_TRAINING_STATS 2
823
+
824
+static void
825
+ml_host_train(ml_host_t *host)
826
+{
827
+ worker_register("MLTRAIN");
828
+ worker_register_job_name(WORKER_JOB_TRAINING_FIND, "find");
829
+ worker_register_job_name(WORKER_JOB_TRAINING_TRAIN, "train");
830
+ worker_register_job_name(WORKER_JOB_TRAINING_STATS, "stats");
831
+
832
+ service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, (force_quit_t ) ml_host_cancel_training_thread, host->rh, true);
833
+
834
+ while (service_running(SERVICE_ML_TRAINING)) {
835
+ ml_training_request_t training_req = ml_queue_pop(host->training_queue);
836
+ size_t queue_size = ml_queue_size(host->training_queue) + 1;
837
+
838
+ if (host->threads_cancelled) {
839
+ info("Stopping training thread for host %s because it was cancelled", rrdhost_hostname(host->rh));
840
+ break;
841
+ }
842
+
843
+ usec_t allotted_ut = (Cfg.train_every * host->rh->rrd_update_every * USEC_PER_SEC) / queue_size;
844
+ if (allotted_ut > USEC_PER_SEC)
845
+ allotted_ut = USEC_PER_SEC;
846
+
847
+ usec_t start_ut = now_monotonic_usec();
848
+ enum ml_training_result training_res;
849
+ {
850
+ worker_is_busy(WORKER_JOB_TRAINING_FIND);
851
+ ml_acquired_dimension_t acq_dim = ml_acquired_dimension_get(host->rh, training_req.chart_id, training_req.dimension_id);
852
+
853
+ worker_is_busy(WORKER_JOB_TRAINING_TRAIN);
854
+ training_res = ml_acquired_dimension_train(acq_dim, training_req);
855
+
856
+ string_freez(training_req.chart_id);
857
+ string_freez(training_req.dimension_id);
858
+
859
+ ml_acquired_dimension_release(acq_dim);
860
+ }
861
+ usec_t consumed_ut = now_monotonic_usec() - start_ut;
862
+
863
+ worker_is_busy(WORKER_JOB_TRAINING_STATS);
864
+
865
+ usec_t remaining_ut = 0;
866
+ if (consumed_ut < allotted_ut)
867
+ remaining_ut = allotted_ut - consumed_ut;
868
+
869
+ {
870
+ netdata_mutex_lock(&host->mutex);
871
+
872
+ host->ts.queue_size += queue_size;
873
+ host->ts.num_popped_items += 1;
874
+
875
+ host->ts.allotted_ut += allotted_ut;
876
+ host->ts.consumed_ut += consumed_ut;
877
+ host->ts.remaining_ut += remaining_ut;
878
+
879
+ switch (training_res) {
880
+ case TRAINING_RESULT_OK:
881
+ host->ts.training_result_ok += 1;
882
+ break;
883
+ case TRAINING_RESULT_INVALID_QUERY_TIME_RANGE:
884
+ host->ts.training_result_invalid_query_time_range += 1;
885
+ break;
886
+ case TRAINING_RESULT_NOT_ENOUGH_COLLECTED_VALUES:
887
+ host->ts.training_result_not_enough_collected_values += 1;
888
+ break;
889
+ case TRAINING_RESULT_NULL_ACQUIRED_DIMENSION:
890
+ host->ts.training_result_null_acquired_dimension += 1;
891
+ break;
892
+ case TRAINING_RESULT_CHART_UNDER_REPLICATION:
893
+ host->ts.training_result_chart_under_replication += 1;
894
+ break;
895
+ }
896
+
897
+ netdata_mutex_unlock(&host->mutex);
898
+ }
899
+
900
+ worker_is_idle();
901
+ std::this_thread::sleep_for(std::chrono::microseconds{remaining_ut});
902
+ worker_is_busy(0);
903
+ }
904
+}
905
+
906
+static void *
907
+train_main(void *arg)
908
+{
909
+ size_t max_elements_needed_for_training = Cfg.max_train_samples * (Cfg.lag_n + 1);
910
+ tls_data.training_cns = new calculated_number_t[max_elements_needed_for_training]();
911
+ tls_data.scratch_training_cns = new calculated_number_t[max_elements_needed_for_training]();
912
+
913
+ ml_host_t *host = (ml_host_t *) arg;
914
+ ml_host_train(host);
915
+ return NULL;
916
}
917
918
static void *
@@ -804,10 +926,12 @@ ml_detect_main(void *arg)
926
worker_register_job_name(WORKER_JOB_DETECTION_HOST_CHART, "host chart");
927
worker_register_job_name(WORKER_JOB_DETECTION_STATS, "training stats");
928
929
+ service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, NULL, NULL, true);
930
+
931
heartbeat_t hb;
932
heartbeat_init(&hb);
933
810
- while (!Cfg.detection_stop) {
934
+ while (service_running((SERVICE_TYPE)(SERVICE_ML_PREDICTION | SERVICE_COLLECTORS))) {
935
worker_is_idle();
936
heartbeat_next(&hb, USEC_PER_SEC);
937
@@ -821,39 +945,6 @@ ml_detect_main(void *arg)
945
ml_host_detect_once((ml_host_t *) rh->ml_host);
946
}
947
dfe_done(rhp);
824
-
825
- if (Cfg.enable_statistics_charts) {
826
- // collect and update training thread stats
827
- for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
828
- ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
829
-
830
- netdata_mutex_lock(&training_thread->nd_mutex);
831
- ml_training_stats_t training_stats = training_thread->training_stats;
832
- training_thread->training_stats = {};
833
- netdata_mutex_unlock(&training_thread->nd_mutex);
834
-
835
- // calc the avg values
836
- if (training_stats.num_popped_items) {
837
- training_stats.queue_size /= training_stats.num_popped_items;
838
- training_stats.allotted_ut /= training_stats.num_popped_items;
839
- training_stats.consumed_ut /= training_stats.num_popped_items;
840
- training_stats.remaining_ut /= training_stats.num_popped_items;
841
- } else {
842
- training_stats.queue_size = 0;
843
- training_stats.allotted_ut = 0;
844
- training_stats.consumed_ut = 0;
845
- training_stats.remaining_ut = 0;
846
-
847
- training_stats.training_result_ok = 0;
848
- training_stats.training_result_invalid_query_time_range = 0;
849
- training_stats.training_result_not_enough_collected_values = 0;
850
- training_stats.training_result_null_acquired_dimension = 0;
851
- training_stats.training_result_chart_under_replication = 0;
852
- }
853
-
854
- ml_update_training_statistics_chart(training_thread, training_stats);
855
- }
856
- }
948
}
949
950
return NULL;
@@ -887,6 +978,31 @@ bool ml_streaming_enabled()
978
return Cfg.stream_anomaly_detection_charts;
979
}
980
981
+void ml_init()
982
+{
983
+ // Read config values
984
+ ml_config_load(&Cfg);
985
+
986
+ if (!Cfg.enable_anomaly_detection)
987
+ return;
988
+
989
+ // Generate random numbers to efficiently sample the features we need
990
+ // for KMeans clustering.
991
+ std::random_device RD;
992
+ std::mt19937 Gen(RD());
993
+
994
+ Cfg.random_nums.reserve(Cfg.max_train_samples);
995
+ for (size_t Idx = 0; Idx != Cfg.max_train_samples; Idx++)
996
+ Cfg.random_nums.push_back(Gen());
997
+
998
+
999
+ // start detection & training threads
1000
+ char tag[NETDATA_THREAD_TAG_MAX + 1];
1001
+
1002
+ snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "PREDICT");
1003
+ netdata_thread_create(&Cfg.detection_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_detect_main, NULL);
1004
+}
1005
+
1006
void ml_host_new(RRDHOST *rh)
1007
{
1008
if (!ml_enabled(rh))
@@ -896,12 +1012,14 @@ void ml_host_new(RRDHOST *rh)
1012
1013
host->rh = rh;
1014
host->mls = ml_machine_learning_stats_t();
899
- //host->ts = ml_training_stats_t();
900
-
901
- static std::atomic<size_t> times_called(0);
902
- host->training_queue = Cfg.training_threads[times_called++ % Cfg.num_training_threads].training_queue;
1015
+ host->ts = ml_training_stats_t();
1016
1017
host->host_anomaly_rate = 0.0;
1018
+ host->threads_running = false;
1019
+ host->threads_cancelled = false;
1020
+ host->threads_joined = false;
1021
+
1022
+ host->training_queue = ml_queue_init();
1023
1024
netdata_mutex_init(&host->mutex);
1025
@@ -915,6 +1033,7 @@ void ml_host_delete(RRDHOST *rh)
1033
return;
1034
1035
netdata_mutex_destroy(&host->mutex);
1036
+ ml_queue_destroy(host->training_queue);
1037
1038
delete host;
1039
rh->ml_host = NULL;
@@ -981,6 +1100,69 @@ void ml_host_get_models(RRDHOST *rh, BUFFER *wb)
1100
error("Fetching KMeans models is not supported yet");
1101
}
1102
1103
+void ml_host_start_training_thread(RRDHOST *rh)
1104
+{
1105
+ if (!rh || !rh->ml_host)
1106
+ return;
1107
+
1108
+ ml_host_t *host = (ml_host_t *) rh->ml_host;
1109
+
1110
+ if (host->threads_running) {
1111
+ error("Anomaly detections threads for host %s are already-up and running.", rrdhost_hostname(host->rh));
1112
+ return;
1113
+ }
1114
+
1115
+ host->threads_running = true;
1116
+ host->threads_cancelled = false;
1117
+ host->threads_joined = false;
1118
+
1119
+ char tag[NETDATA_THREAD_TAG_MAX + 1];
1120
+
1121
+ snprintfz(tag, NETDATA_THREAD_TAG_MAX, "MLTR[%s]", rrdhost_hostname(host->rh));
1122
+ netdata_thread_create(&host->training_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, train_main, static_cast<void *>(host));
1123
+}
1124
+
1125
+void ml_host_cancel_training_thread(RRDHOST *rh)
1126
+{
1127
+ if (!rh || !rh->ml_host)
1128
+ return;
1129
+
1130
+ ml_host_t *host = (ml_host_t *) rh->ml_host;
1131
+
1132
+ if (!host->threads_running) {
1133
+ error("Anomaly detections threads for host %s have already been stopped.", rrdhost_hostname(host->rh));
1134
+ return;
1135
+ }
1136
+
1137
+ if (!host->threads_cancelled) {
1138
+ host->threads_cancelled = true;
1139
+
1140
+ // Signal the training queue to stop popping-items
1141
+ ml_queue_signal(host->training_queue);
1142
+ netdata_thread_cancel(host->training_thread);
1143
+ }
1144
+}
1145
+
1146
+void ml_host_stop_training_thread(RRDHOST *rh)
1147
+{
1148
+ if (!rh || !rh->ml_host)
1149
+ return;
1150
+
1151
+ ml_host_cancel_training_thread(rh);
1152
+
1153
+ ml_host_t *host = (ml_host_t *) rh->ml_host;
1154
+
1155
+ if (!host->threads_joined) {
1156
+ host->threads_joined = true;
1157
+ host->threads_running = false;
1158
+
1159
+ delete[] tls_data.training_cns;
1160
+ delete[] tls_data.scratch_training_cns;
1161
+
1162
+ netdata_thread_join(host->training_thread, NULL);
1163
+ }
1164
+}
1165
+
1166
void ml_chart_new(RRDSET *rs)
1167
{
1168
ml_host_t *host = (ml_host_t *) rs->rrdhost->ml_host;
@@ -1085,185 +1267,3 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
1267
1268
return is_anomalous;
1269
}
1088
-
1089
-static void *ml_train_main(void *arg) {
1090
- ml_training_thread_t *training_thread = (ml_training_thread_t *) arg;
1091
-
1092
- char worker_name[1024];
1093
- snprintfz(worker_name, 1024, "training_thread_%zu", training_thread->id);
1094
- worker_register("MLTRAIN");
1095
-
1096
- worker_register_job_name(WORKER_TRAIN_QUEUE_POP, "pop queue");
1097
- worker_register_job_name(WORKER_TRAIN_ACQUIRE_DIMENSION, "acquire");
1098
- worker_register_job_name(WORKER_TRAIN_QUERY, "query");
1099
- worker_register_job_name(WORKER_TRAIN_KMEANS, "kmeans");
1100
- worker_register_job_name(WORKER_TRAIN_UPDATE_MODELS, "update models");
1101
- worker_register_job_name(WORKER_TRAIN_RELEASE_DIMENSION, "release");
1102
- worker_register_job_name(WORKER_TRAIN_UPDATE_HOST, "update host");
1103
-
1104
- while (!Cfg.training_stop) {
1105
- worker_is_busy(WORKER_TRAIN_QUEUE_POP);
1106
-
1107
- ml_training_request_t training_req = ml_queue_pop(training_thread->training_queue);
1108
-
1109
- // we know this thread has been cancelled, when the queue starts
1110
- // returning "null" requests without blocking on queue's pop().
1111
- if (training_req.host_id == NULL)
1112
- break;
1113
-
1114
- size_t queue_size = ml_queue_size(training_thread->training_queue) + 1;
1115
-
1116
- usec_t allotted_ut = (Cfg.train_every * USEC_PER_SEC) / queue_size;
1117
- if (allotted_ut > USEC_PER_SEC)
1118
- allotted_ut = USEC_PER_SEC;
1119
-
1120
- usec_t start_ut = now_monotonic_usec();
1121
-
1122
- enum ml_training_result training_res;
1123
- {
1124
- worker_is_busy(WORKER_TRAIN_ACQUIRE_DIMENSION);
1125
- ml_acquired_dimension_t acq_dim = ml_acquired_dimension_get(
1126
- training_req.host_id,
1127
- training_req.chart_id,
1128
- training_req.dimension_id);
1129
-
1130
- training_res = ml_acquired_dimension_train(training_thread, acq_dim, training_req);
1131
-
1132
- string_freez(training_req.host_id);
1133
- string_freez(training_req.chart_id);
1134
- string_freez(training_req.dimension_id);
1135
-
1136
- worker_is_busy(WORKER_TRAIN_RELEASE_DIMENSION);
1137
- ml_acquired_dimension_release(acq_dim);
1138
- }
1139
-
1140
- usec_t consumed_ut = now_monotonic_usec() - start_ut;
1141
-
1142
- usec_t remaining_ut = 0;
1143
- if (consumed_ut < allotted_ut)
1144
- remaining_ut = allotted_ut - consumed_ut;
1145
-
1146
- if (Cfg.enable_statistics_charts) {
1147
- worker_is_busy(WORKER_TRAIN_UPDATE_HOST);
1148
-
1149
- netdata_mutex_lock(&training_thread->nd_mutex);
1150
-
1151
- training_thread->training_stats.queue_size += queue_size;
1152
- training_thread->training_stats.num_popped_items += 1;
1153
-
1154
- training_thread->training_stats.allotted_ut += allotted_ut;
1155
- training_thread->training_stats.consumed_ut += consumed_ut;
1156
- training_thread->training_stats.remaining_ut += remaining_ut;
1157
-
1158
- switch (training_res) {
1159
- case TRAINING_RESULT_OK:
1160
- training_thread->training_stats.training_result_ok += 1;
1161
- break;
1162
- case TRAINING_RESULT_INVALID_QUERY_TIME_RANGE:
1163
- training_thread->training_stats.training_result_invalid_query_time_range += 1;
1164
- break;
1165
- case TRAINING_RESULT_NOT_ENOUGH_COLLECTED_VALUES:
1166
- training_thread->training_stats.training_result_not_enough_collected_values += 1;
1167
- break;
1168
- case TRAINING_RESULT_NULL_ACQUIRED_DIMENSION:
1169
- training_thread->training_stats.training_result_null_acquired_dimension += 1;
1170
- break;
1171
- case TRAINING_RESULT_CHART_UNDER_REPLICATION:
1172
- training_thread->training_stats.training_result_chart_under_replication += 1;
1173
- break;
1174
- }
1175
-
1176
- netdata_mutex_unlock(&training_thread->nd_mutex);
1177
- }
1178
-
1179
- worker_is_idle();
1180
- std::this_thread::sleep_for(std::chrono::microseconds{remaining_ut});
1181
- }
1182
-
1183
- return NULL;
1184
-}
1185
-
1186
-void ml_init()
1187
-{
1188
- // Read config values
1189
- ml_config_load(&Cfg);
1190
-
1191
- if (!Cfg.enable_anomaly_detection)
1192
- return;
1193
-
1194
- // Generate random numbers to efficiently sample the features we need
1195
- // for KMeans clustering.
1196
- std::random_device RD;
1197
- std::mt19937 Gen(RD());
1198
-
1199
- Cfg.random_nums.reserve(Cfg.max_train_samples);
1200
- for (size_t Idx = 0; Idx != Cfg.max_train_samples; Idx++)
1201
- Cfg.random_nums.push_back(Gen());
1202
-
1203
-
1204
- // start detection & training threads
1205
- Cfg.detection_stop = false;
1206
- Cfg.training_stop = false;
1207
-
1208
- char tag[NETDATA_THREAD_TAG_MAX + 1];
1209
-
1210
- snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "PREDICT");
1211
- netdata_thread_create(&Cfg.detection_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_detect_main, NULL);
1212
-
1213
- Cfg.training_threads.resize(Cfg.num_training_threads);
1214
- for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1215
- ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1216
-
1217
-
1218
- size_t max_elements_needed_for_training = Cfg.max_train_samples * (Cfg.lag_n + 1);
1219
- training_thread->training_cns = new calculated_number_t[max_elements_needed_for_training]();
1220
- training_thread->scratch_training_cns = new calculated_number_t[max_elements_needed_for_training]();
1221
-
1222
- training_thread->id = idx;
1223
- training_thread->training_queue = ml_queue_init();
1224
- netdata_mutex_init(&training_thread->nd_mutex);
1225
-
1226
- snprintfz(tag, NETDATA_THREAD_TAG_MAX, "TRAIN[%zu]", training_thread->id);
1227
- netdata_thread_create(&training_thread->nd_thread, tag, NETDATA_THREAD_OPTION_JOINABLE, ml_train_main, training_thread);
1228
- }
1229
-}
1230
-
1231
-void ml_fini()
1232
-{
1233
- Cfg.detection_stop = true;
1234
- Cfg.training_stop = true;
1235
-
1236
- netdata_thread_cancel(Cfg.detection_thread);
1237
- netdata_thread_join(Cfg.detection_thread, NULL);
1238
-
1239
- // signal the training queue of each thread
1240
- for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1241
- ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1242
-
1243
- ml_queue_signal(training_thread->training_queue);
1244
- }
1245
-
1246
- // cancel training threads
1247
- for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1248
- ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1249
-
1250
- netdata_thread_cancel(training_thread->nd_thread);
1251
- }
1252
-
1253
- // join training threads
1254
- for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1255
- ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1256
-
1257
- netdata_thread_join(training_thread->nd_thread, NULL);
1258
- }
1259
-
1260
- // clear training thread data
1261
- for (size_t idx = 0; idx != Cfg.num_training_threads; idx++) {
1262
- ml_training_thread_t *training_thread = &Cfg.training_threads[idx];
1263
-
1264
- delete[] training_thread->training_cns;
1265
- delete[] training_thread->scratch_training_cns;
1266
- ml_queue_destroy(training_thread->training_queue);
1267
- netdata_mutex_destroy(&training_thread->nd_mutex);
1268
- }
1269
-}
ml/ml.h
+4
-4
@@ -13,9 +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);
17
18
void ml_host_new(RRDHOST *rh);
19
void ml_host_delete(RRDHOST *rh);
@@ -24,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);
@@ -33,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
36
-void ml_update_global_statistics_charts(uint64_t models_consulted);
37
-
38
#ifdef __cplusplus
39
};
40
#endif