Batch ML model load commands (#16155)
Batch ML model load
Stelios Fragkakis committed
Oct 10, 2023 at 19:43 UTC
e8c770d1a8eb0ad75169eef6faa73ff125f250c3
3 files changed
+76
-3
daemon/event_loop.c
+1
@@ -52,6 +52,7 @@ void register_libuv_worker_jobs() {
52
worker_register_job_name(UV_EVENT_HOST_CONTEXT_LOAD, "metadata load host context");
53
worker_register_job_name(UV_EVENT_METADATA_STORE, "metadata store host");
54
worker_register_job_name(UV_EVENT_METADATA_CLEANUP, "metadata cleanup");
55
+ worker_register_job_name(UV_EVENT_METADATA_ML_LOAD, "metadata load ml models");
56
57
// netdatacli
58
worker_register_job_name(UV_EVENT_SCHEDULE_CMD, "schedule command");
daemon/event_loop.h
+1
@@ -44,6 +44,7 @@ enum event_loop_job {
44
UV_EVENT_HOST_CONTEXT_LOAD,
45
UV_EVENT_METADATA_STORE,
46
UV_EVENT_METADATA_CLEANUP,
47
+ UV_EVENT_METADATA_ML_LOAD,
48
49
// netdatacli
50
UV_EVENT_SCHEDULE_CMD,
database/sqlite/sqlite_metadata.c
+74
-3
@@ -55,7 +55,7 @@
55
#define DELETE_NON_EXISTING_LOCALHOST "DELETE FROM host WHERE hops = 0 AND host_id <> @host_id;"
56
#define DELETE_MISSING_NODE_INSTANCES "DELETE FROM node_instance WHERE host_id NOT IN (SELECT host_id FROM host);"
57
58
-#define METADATA_CMD_Q_MAX_SIZE (1024) // Max queue size; callers will block until there is room
58
+#define METADATA_CMD_Q_MAX_SIZE (2048) // Max queue size; callers will block until there is room
59
#define METADATA_MAINTENANCE_FIRST_CHECK (1800) // Maintenance first run after agent startup in seconds
60
#define METADATA_MAINTENANCE_REPEAT (60) // Repeat if last run for dimensions, charts, labels needs more work
61
#define METADATA_HEALTH_LOG_INTERVAL (3600) // Repeat maintenance for health
@@ -105,6 +105,7 @@ struct metadata_database_cmdqueue {
105
typedef enum {
106
METADATA_FLAG_PROCESSING = (1 << 0), // store or cleanup
107
METADATA_FLAG_SHUTDOWN = (1 << 1), // Shutting down
108
+ METADATA_FLAG_ML_LOADING = (1 << 2), // ML model load in progress
109
} METADATA_FLAG;
110
111
struct metadata_wc {
@@ -1231,6 +1232,13 @@ void run_metadata_cleanup(struct metadata_wc *wc)
1232
(void) sqlite3_wal_checkpoint(db_meta, NULL);
1233
}
1234
1235
+struct ml_model_payload {
1236
+ uv_work_t request;
1237
+ struct metadata_wc *wc;
1238
+ Pvoid_t JudyL;
1239
+ size_t count;
1240
+};
1241
+
1242
struct scan_metadata_payload {
1243
uv_work_t request;
1244
struct metadata_wc *wc;
@@ -1556,6 +1564,42 @@ static void start_metadata_hosts(uv_work_t *req __maybe_unused)
1564
worker_is_idle();
1565
}
1566
1567
+// Callback after scan of hosts is done
1568
+static void after_start_ml_model_load(uv_work_t *req, int status __maybe_unused)
1569
+{
1570
+ struct ml_model_payload *ml_data = req->data;
1571
+ struct metadata_wc *wc = ml_data->wc;
1572
+ metadata_flag_clear(wc, METADATA_FLAG_ML_LOADING);
1573
+ JudyLFreeArray(&ml_data->JudyL, PJE0);
1574
+ freez(ml_data);
1575
+}
1576
+
1577
+static void start_ml_model_load(uv_work_t *req __maybe_unused)
1578
+{
1579
+ register_libuv_worker_jobs();
1580
+
1581
+ struct ml_model_payload *ml_data = req->data;
1582
+
1583
+ worker_is_busy(UV_EVENT_METADATA_ML_LOAD);
1584
+
1585
+ Pvoid_t *PValue;
1586
+ Word_t Index = 0;
1587
+ bool first = true;
1588
+ RRDDIM *rd;
1589
+ RRDDIM_ACQUIRED *rda;
1590
+ internal_error(true, "Batch ML load loader, %zu items", ml_data->count);
1591
+ while((PValue = JudyLFirstThenNext(ml_data->JudyL, &Index, &first))) {
1592
+ UNUSED(PValue);
1593
+ rda = (RRDDIM_ACQUIRED *) Index;
1594
+ rd = rrddim_acquired_to_rrddim(rda);
1595
+ ml_dimension_load_models(rd);
1596
+ rrddim_acquired_release(rda);
1597
+ }
1598
+ worker_is_idle();
1599
+}
1600
+
1601
+
1602
+
1603
static void metadata_event_loop(void *arg)
1604
{
1605
worker_register("METASYNC");
@@ -1610,6 +1654,7 @@ static void metadata_event_loop(void *arg)
1654
BUFFER *work_buffer = buffer_create(1024, &netdata_buffers_statistics.buffers_sqlite);
1655
struct scan_metadata_payload *data;
1656
1657
+ struct ml_model_payload *ml_data = NULL;
1658
while (shutdown == 0 || (wc->flags & METADATA_FLAG_PROCESSING)) {
1659
uuid_t *uuid;
1660
RRDHOST *host = NULL;
@@ -1636,6 +1681,24 @@ static void metadata_event_loop(void *arg)
1681
if (likely(opcode != METADATA_DATABASE_NOOP))
1682
worker_is_busy(opcode);
1683
1684
+ // Have pending ML models to load?
1685
+ if (opcode != METADATA_ML_LOAD_MODELS && ml_data && ml_data->count) {
1686
+ static usec_t ml_submit_last = 0;
1687
+ usec_t now = now_monotonic_usec();
1688
+ if (!ml_submit_last)
1689
+ ml_submit_last = now;
1690
+
1691
+ if (!metadata_flag_check(wc, METADATA_FLAG_ML_LOADING) && (now - ml_submit_last > 150 * USEC_PER_MS)) {
1692
+ metadata_flag_set(wc, METADATA_FLAG_ML_LOADING);
1693
+ if (unlikely(uv_queue_work(loop, &ml_data->request, start_ml_model_load, after_start_ml_model_load)))
1694
+ metadata_flag_clear(wc, METADATA_FLAG_ML_LOADING);
1695
+ else {
1696
+ ml_submit_last = now;
1697
+ ml_data = NULL;
1698
+ }
1699
+ }
1700
+ }
1701
+
1702
switch (opcode) {
1703
case METADATA_DATABASE_NOOP:
1704
case METADATA_DATABASE_TIMER:
@@ -1643,8 +1706,16 @@ static void metadata_event_loop(void *arg)
1706
1707
case METADATA_ML_LOAD_MODELS: {
1708
RRDDIM *rd = (RRDDIM *) cmd.param[0];
1646
- if (!shutdown)
1647
- ml_dimension_load_models(rd);
1709
+ RRDDIM_ACQUIRED *rda = rrddim_find_and_acquire(rd->rrdset, rrddim_id(rd));
1710
+ if (likely(rda)) {
1711
+ if (!ml_data) {
1712
+ ml_data = callocz(1,sizeof(*ml_data));
1713
+ ml_data->request.data = ml_data;
1714
+ ml_data->wc = wc;
1715
+ }
1716
+ JudyLIns(&ml_data->JudyL, (Word_t)rda, PJE0);
1717
+ ml_data->count++;
1718
+ }
1719
break;
1720
}
1721
case METADATA_DEL_DIMENSION: