Use static thread-pool for training. (#14702)
* Use static thread-pool for training. * Add missing function definition * disable training stats chart * Add config option to explicitly enable ML stats charts. --------- Co-authored-by: Costa Tsaousis <costa@netdata.cloud>
vkalintiris committed
Mar 21, 2023 at 11:24 UTC
5046e034212c008557dd014196b6f6204eda24b2
11 files changed
+436
-408
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
+7
-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,10 @@ void netdata_cleanup_and_exit(int ret) {
336
}
337
#endif
338
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,12 +351,11 @@ void netdata_cleanup_and_exit(int ret) {
351
| SERVICE_ACLKSYNC
352
);
353
354
- delta_shutdown_time("stop replication, exporters, ML training, health and web servers threads");
354
+ delta_shutdown_time("stop replication, exporters, health and web servers threads");
355
356
timeout = !service_wait_exit(
357
SERVICE_REPLICATION
358
| SERVICE_EXPORTERS
359
- | SERVICE_ML_TRAINING
359
| SERVICE_HEALTH
360
| SERVICE_WEB_SERVER
361
, 3 * USEC_PER_SEC);
@@ -368,11 +367,10 @@ void netdata_cleanup_and_exit(int ret) {
367
| SERVICE_STREAMING
368
, 3 * USEC_PER_SEC);
369
371
- delta_shutdown_time("stop ML prediction and context threads");
370
+ delta_shutdown_time("stop context thread");
371
372
timeout = !service_wait_exit(
374
- SERVICE_ML_PREDICTION
375
- | SERVICE_CONTEXT
373
+ SERVICE_CONTEXT
374
, 3 * USEC_PER_SEC);
375
376
delta_shutdown_time("stop maintenance thread");
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
+6
@@ -19,6 +19,8 @@ bool ml_streaming_enabled() {
19
20
void ml_init(void) {}
21
22
+void ml_fini(void) {}
23
+
24
void ml_host_new(RRDHOST *rh) {
25
UNUSED(rh);
26
}
@@ -86,4 +88,8 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
88
return false;
89
}
90
91
+void ml_update_global_statistics_charts(uint64_t models_consulted) {
92
+ UNUSED(models_consulted);
93
+}
94
+
95
#endif
ml/ml-private.h
+26
-10
@@ -33,6 +33,7 @@ typedef struct {
33
/*
34
* KMeans
35
*/
36
+
37
typedef struct {
38
size_t num_clusters;
39
size_t max_iterations;
@@ -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
+273
-273
@@ -8,14 +8,13 @@
8
9
#include "ad_charts.h"
10
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;
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
18
19
/*
20
* Functions to convert enums to strings
@@ -264,7 +263,14 @@ ml_queue_pop(ml_queue_t *q)
263
{
264
netdata_mutex_lock(&q->mutex);
265
267
- ml_training_request_t req = { NULL, NULL, 0, 0, 0 };
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
+ };
274
275
while (q->internal.empty()) {
276
pthread_cond_wait(&q->cond_var, &q->mutex);
@@ -307,7 +313,7 @@ ml_queue_signal(ml_queue_t *q)
313
*/
314
315
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)
316
+ml_dimension_calculated_numbers(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
317
{
318
ml_training_response_t training_response = {};
319
@@ -351,7 +357,7 @@ ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t
357
STORAGE_PRIORITY_BEST_EFFORT);
358
359
size_t idx = 0;
354
- memset(tls_data.training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
360
+ memset(training_thread->training_cns, 0, sizeof(calculated_number_t) * max_n * (Cfg.lag_n + 1));
361
calculated_number_t last_value = std::numeric_limits<calculated_number_t>::quiet_NaN();
362
363
while (!ops->is_finished(&handle)) {
@@ -368,11 +374,11 @@ ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t
374
training_response.db_after_t = timestamp;
375
training_response.db_before_t = timestamp;
376
371
- tls_data.training_cns[idx] = value;
372
- last_value = tls_data.training_cns[idx];
377
+ training_thread->training_cns[idx] = value;
378
+ last_value = training_thread->training_cns[idx];
379
training_response.collected_values++;
380
} else
375
- tls_data.training_cns[idx] = last_value;
381
+ training_thread->training_cns[idx] = last_value;
382
383
idx++;
384
}
@@ -387,20 +393,21 @@ ml_dimension_calculated_numbers(ml_dimension_t *dim, const ml_training_request_t
393
}
394
395
// Find first non-NaN value.
390
- for (idx = 0; std::isnan(tls_data.training_cns[idx]); idx++, training_response.total_values--) { }
396
+ for (idx = 0; std::isnan(training_thread->training_cns[idx]); idx++, training_response.total_values--) { }
397
398
// Overwrite NaN values.
399
if (idx != 0)
394
- memmove(tls_data.training_cns, &tls_data.training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
400
+ memmove(training_thread->training_cns, &training_thread->training_cns[idx], sizeof(calculated_number_t) * training_response.total_values);
401
402
training_response.result = TRAINING_RESULT_OK;
397
- return { tls_data.training_cns, training_response };
403
+ return { training_thread->training_cns, training_response };
404
}
405
406
static enum ml_training_result
401
-ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &training_request)
407
+ml_dimension_train_model(ml_training_thread_t *training_thread, ml_dimension_t *dim, const ml_training_request_t &training_request)
408
{
403
- auto P = ml_dimension_calculated_numbers(dim, training_request);
409
+ worker_is_busy(WORKER_TRAIN_QUERY);
410
+ auto P = ml_dimension_calculated_numbers(training_thread, dim, training_request);
411
ml_training_response_t training_response = P.second;
412
413
if (training_response.result != TRAINING_RESULT_OK) {
@@ -429,15 +436,16 @@ ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &train
436
}
437
438
// compute kmeans
439
+ worker_is_busy(WORKER_TRAIN_KMEANS);
440
{
433
- memcpy(tls_data.scratch_training_cns, tls_data.training_cns,
441
+ memcpy(training_thread->scratch_training_cns, training_thread->training_cns,
442
training_response.total_values * sizeof(calculated_number_t));
443
444
ml_features_t features = {
445
Cfg.diff_n, Cfg.smooth_n, Cfg.lag_n,
438
- tls_data.scratch_training_cns, training_response.total_values,
439
- tls_data.training_cns, training_response.total_values,
440
- tls_data.training_samples
446
+ training_thread->scratch_training_cns, training_response.total_values,
447
+ training_thread->training_cns, training_response.total_values,
448
+ training_thread->training_samples
449
};
450
ml_features_preprocess(&features);
451
@@ -446,6 +454,7 @@ ml_dimension_train_model(ml_dimension_t *dim, const ml_training_request_t &train
454
}
455
456
// update kmeans models
457
+ worker_is_busy(WORKER_TRAIN_UPDATE_MODELS);
458
{
459
netdata_mutex_lock(&dim->mutex);
460
@@ -497,11 +506,16 @@ ml_dimension_schedule_for_training(ml_dimension_t *dim, time_t curr_time)
506
}
507
508
if (schedule_for_training) {
500
- ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
509
ml_training_request_t req = {
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),
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),
516
};
517
+
518
+ ml_host_t *host = (ml_host_t *) dim->rd->rrdset->rrdhost->ml_host;
519
ml_queue_push(host->training_queue, req);
520
}
521
}
@@ -677,7 +691,6 @@ ml_host_detect_once(ml_host_t *host)
691
692
host->mls = {};
693
ml_machine_learning_stats_t mls_copy = {};
680
- ml_training_stats_t ts_copy = {};
694
695
{
696
netdata_mutex_lock(&host->mutex);
@@ -721,54 +734,14 @@ ml_host_detect_once(ml_host_t *host)
734
735
mls_copy = host->mls;
736
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
-
737
netdata_mutex_unlock(&host->mutex);
738
}
739
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
-
740
worker_is_busy(WORKER_JOB_DETECTION_DIM_CHART);
741
ml_update_dimensions_chart(host, mls_copy);
742
743
worker_is_busy(WORKER_JOB_DETECTION_HOST_CHART);
744
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);
745
}
746
747
typedef struct {
@@ -777,18 +750,21 @@ typedef struct {
750
} ml_acquired_dimension_t;
751
752
static ml_acquired_dimension_t
780
-ml_acquired_dimension_get(RRDHOST *rh, STRING *chart_id, STRING *dimension_id)
753
+ml_acquired_dimension_get(STRING *host_id, STRING *chart_id, STRING *dimension_id)
754
{
755
RRDDIM_ACQUIRED *acq_rd = NULL;
756
ml_dimension_t *dim = NULL;
757
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;
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
+ }
768
}
769
}
770
@@ -809,110 +785,12 @@ ml_acquired_dimension_release(ml_acquired_dimension_t acq_dim)
785
}
786
787
static enum ml_training_result
812
-ml_acquired_dimension_train(ml_acquired_dimension_t acq_dim, const ml_training_request_t &TR)
788
+ml_acquired_dimension_train(ml_training_thread_t *training_thread, ml_acquired_dimension_t acq_dim, const ml_training_request_t &tr)
789
{
790
if (!acq_dim.dim)
791
return TRAINING_RESULT_NULL_ACQUIRED_DIMENSION;
792
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;
793
+ return ml_dimension_train_model(training_thread, acq_dim.dim, tr);
794
}
795
796
static void *
@@ -926,12 +804,10 @@ ml_detect_main(void *arg)
804
worker_register_job_name(WORKER_JOB_DETECTION_HOST_CHART, "host chart");
805
worker_register_job_name(WORKER_JOB_DETECTION_STATS, "training stats");
806
929
- service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, NULL, NULL, true);
930
-
807
heartbeat_t hb;
808
heartbeat_init(&hb);
809
934
- while (service_running((SERVICE_TYPE)(SERVICE_ML_PREDICTION | SERVICE_COLLECTORS))) {
810
+ while (!Cfg.detection_stop) {
811
worker_is_idle();
812
heartbeat_next(&hb, USEC_PER_SEC);
813
@@ -945,6 +821,39 @@ ml_detect_main(void *arg)
821
ml_host_detect_once((ml_host_t *) rh->ml_host);
822
}
823
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
+ }
857
}
858
859
return NULL;
@@ -978,31 +887,6 @@ bool ml_streaming_enabled()
887
return Cfg.stream_anomaly_detection_charts;
888
}
889
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
-
890
void ml_host_new(RRDHOST *rh)
891
{
892
if (!ml_enabled(rh))
@@ -1012,14 +896,12 @@ void ml_host_new(RRDHOST *rh)
896
897
host->rh = rh;
898
host->mls = ml_machine_learning_stats_t();
1015
- host->ts = ml_training_stats_t();
899
+ //host->ts = ml_training_stats_t();
900
1017
- host->host_anomaly_rate = 0.0;
1018
- host->threads_running = false;
1019
- host->threads_cancelled = false;
1020
- host->threads_joined = false;
901
+ static std::atomic<size_t> times_called(0);
902
+ host->training_queue = Cfg.training_threads[times_called++ % Cfg.num_training_threads].training_queue;
903
1022
- host->training_queue = ml_queue_init();
904
+ host->host_anomaly_rate = 0.0;
905
906
netdata_mutex_init(&host->mutex);
907
@@ -1033,7 +915,6 @@ void ml_host_delete(RRDHOST *rh)
915
return;
916
917
netdata_mutex_destroy(&host->mutex);
1036
- ml_queue_destroy(host->training_queue);
918
919
delete host;
920
rh->ml_host = NULL;
@@ -1100,69 +981,6 @@ void ml_host_get_models(RRDHOST *rh, BUFFER *wb)
981
error("Fetching KMeans models is not supported yet");
982
}
983
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
-
984
void ml_chart_new(RRDSET *rs)
985
{
986
ml_host_t *host = (ml_host_t *) rs->rrdhost->ml_host;
@@ -1267,3 +1085,185 @@ bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool
1085
1086
return is_anomalous;
1087
}
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,7 +13,9 @@ 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_host_new(RRDHOST *rh);
21
void ml_host_delete(RRDHOST *rh);
@@ -22,10 +24,6 @@ void ml_host_get_info(RRDHOST *RH, BUFFER *wb);
24
void ml_host_get_detection_info(RRDHOST *RH, BUFFER *wb);
25
void ml_host_get_models(RRDHOST *RH, BUFFER *wb);
26
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
-
27
void ml_chart_new(RRDSET *rs);
28
void ml_chart_delete(RRDSET *rs);
29
bool ml_chart_update_begin(RRDSET *rs);
@@ -35,6 +33,8 @@ void ml_dimension_new(RRDDIM *rd);
33
void ml_dimension_delete(RRDDIM *rd);
34
bool ml_dimension_is_anomalous(RRDDIM *rd, time_t curr_time, double value, bool exists);
35
36
+void ml_update_global_statistics_charts(uint64_t models_consulted);
37
+
38
#ifdef __cplusplus
39
};
40
#endif