Initialize reusable buffers when streaming ML to parent (#21279)
Refactor kmeans streaming to use reusable buffers in worker
Stelios Fragkakis committed
Nov 10, 2025 at 19:44 UTC
257d073607e69476c6b4bde05f7bb804b83c3cce
3 files changed
+19
-4
src/ml/ml.cc
+7
-4
@@ -585,7 +585,7 @@ ml_dimension_deserialize_kmeans(const char *json_str)
585
return true;
586
}
587
588
-static void ml_dimension_stream_kmeans(const ml_dimension_t *dim)
588
+static void ml_dimension_stream_kmeans(ml_worker_t *worker, const ml_dimension_t *dim)
589
{
590
struct sender_state *s = dim->rd->rrdset->rrdhost->sender;
591
if (!s)
@@ -596,10 +596,13 @@ static void ml_dimension_stream_kmeans(const ml_dimension_t *dim)
596
!rrddim_check_upstream_exposed(dim->rd))
597
return;
598
599
- CLEAN_BUFFER *payload = buffer_create(0, NULL);
599
+ // Reuse worker's buffers instead of allocating new ones
600
+ BUFFER *payload = worker->stream_payload_buffer;
601
+ buffer_flush(payload);
602
ml_dimension_serialize_kmeans(dim, payload);
603
602
- CLEAN_BUFFER *wb = buffer_create(0, NULL);
604
+ BUFFER *wb = worker->stream_wb_buffer;
605
+ buffer_flush(wb);
606
607
buffer_sprintf(
608
wb, PLUGINSD_KEYWORD_JSON " " PLUGINSD_KEYWORD_JSON_CMD_ML_MODEL "\n%s\n" PLUGINSD_KEYWORD_JSON_END "\n",
@@ -650,7 +653,7 @@ static void ml_dimension_update_models(ml_worker_t *worker, ml_dimension_t *dim)
653
model_info.inlined_kmeans = dim->km_contexts.back();
654
worker->pending_model_info.push_back(model_info);
655
653
- ml_dimension_stream_kmeans(dim);
656
+ ml_dimension_stream_kmeans(worker, dim);
657
658
// Clear the training in progress flag
659
dim->training_in_progress = false;
src/ml/ml_public.cc
+8
@@ -408,6 +408,10 @@ void ml_init()
408
worker->queue = ml_queue_init();
409
worker->pending_model_info.reserve(Cfg.flush_models_batch_size);
410
netdata_mutex_init(&worker->nd_mutex);
411
+
412
+ // Initialize reusable buffers for streaming kmeans models
413
+ worker->stream_payload_buffer = buffer_create(0, NULL);
414
+ worker->stream_wb_buffer = buffer_create(0, NULL);
415
}
416
417
// open sqlite db
@@ -509,6 +513,10 @@ void ml_stop_threads()
513
delete[] worker->scratch_training_cns;
514
ml_queue_destroy(worker->queue);
515
netdata_mutex_destroy(&worker->nd_mutex);
516
+
517
+ // Free reusable buffers
518
+ buffer_free(worker->stream_payload_buffer);
519
+ buffer_free(worker->stream_wb_buffer);
520
}
521
}
522
src/ml/ml_worker.h
+4
@@ -24,6 +24,10 @@ typedef struct {
24
25
std::vector<ml_model_info_t> pending_model_info;
26
27
+ // Reusable buffers for streaming kmeans models
28
+ BUFFER *stream_payload_buffer;
29
+ BUFFER *stream_wb_buffer;
30
+
31
RRDSET *queue_stats_rs;
32
RRDDIM *queue_stats_num_create_new_model_requests_rd;
33
RRDDIM *queue_stats_num_add_existing_model_requests_rd;