Improve dimension ML model load (#16262)
* Prepare metadata sync thread cleanup earlier in the shutdown process * Set flag for the dimensions that need ML MODEL load instead of queueing a message in the event loop * Process the dimension ML load during the normal dimension metadata save loop * Use spinlock for cmd queue / dequeue instead of mutex Cleanup queue structure * Remove old ML model load code * Rebase and cleanup
Stelios Fragkakis committed
Oct 31, 2023 at 09:57 UTC
661f2eb6c52f4866b2c6c12985066c111fdc7245
5 files changed
+49
-137
daemon/main.c
+4
-4
@@ -371,6 +371,10 @@ void netdata_cleanup_and_exit(int ret) {
371
SERVICE_REPLICATION // replication has to be stopped after STREAMING, because it cleans up ARAL
372
, 3 * USEC_PER_SEC);
373
374
+ delta_shutdown_time("prepare metasync shutdown");
375
+
376
+ metadata_sync_shutdown_prepare();
377
+
378
delta_shutdown_time("disable ML detection and training threads");
379
380
ml_stop_threads();
@@ -396,10 +400,6 @@ void netdata_cleanup_and_exit(int ret) {
400
401
rrdhost_cleanup_all();
402
399
- delta_shutdown_time("prepare metasync shutdown");
400
-
401
- metadata_sync_shutdown_prepare();
402
-
403
delta_shutdown_time("stop aclk threads");
404
405
timeout = !service_wait_exit(
database/rrd.h
+2
-2
@@ -258,8 +258,8 @@ typedef enum __attribute__ ((__packed__)) rrddim_flags {
258
RRDDIM_FLAG_ARCHIVED = (1 << 2),
259
RRDDIM_FLAG_METADATA_UPDATE = (1 << 3), // Metadata needs to go to the database
260
261
- RRDDIM_FLAG_META_HIDDEN = (1 << 4), // Status of hidden option in the metadata database
262
-
261
+ RRDDIM_FLAG_META_HIDDEN = (1 << 4), // Status of hidden option in the metadata database
262
+ RRDDIM_FLAG_ML_MODEL_LOAD = (1 << 5), // Do ML LOAD for this dimension
263
264
// this is 8 bit
265
} RRDDIM_FLAGS;
database/sqlite/sqlite_functions.c
-1
@@ -85,7 +85,6 @@ sqlite3 *db_meta = NULL;
85
86
#define MAX_PREPARED_STATEMENTS (32)
87
pthread_key_t key_pool[MAX_PREPARED_STATEMENTS];
88
-pthread_key_t plugin_key;
88
89
SQLITE_API int sqlite3_exec_monitored(
90
sqlite3 *db, /* An open database */
database/sqlite/sqlite_metadata.c
+43
-127
@@ -83,7 +83,6 @@ enum metadata_opcode {
83
METADATA_MAINTENANCE,
84
METADATA_SYNC_SHUTDOWN,
85
METADATA_UNITTEST,
86
- METADATA_ML_LOAD_MODELS,
86
// leave this last
87
// we need it to check for worker utilization
88
METADATA_MAX_ENUMERATIONS_DEFINED
@@ -97,14 +96,9 @@ struct metadata_cmd {
96
struct metadata_cmd *prev, *next;
97
};
98
100
-struct metadata_database_cmdqueue {
101
- struct metadata_cmd *cmd_base;
102
-};
103
-
99
typedef enum {
100
METADATA_FLAG_PROCESSING = (1 << 0), // store or cleanup
101
METADATA_FLAG_SHUTDOWN = (1 << 1), // Shutting down
107
- METADATA_FLAG_ML_LOADING = (1 << 2), // ML model load in progress
102
} METADATA_FLAG;
103
104
struct metadata_wc {
@@ -113,13 +107,12 @@ struct metadata_wc {
107
uv_async_t async;
108
uv_timer_t timer_req;
109
time_t metadata_check_after;
116
- volatile unsigned queue_size;
110
METADATA_FLAG flags;
118
- struct completion init_complete;
111
+ struct completion start_stop_complete;
112
struct completion *scan_complete;
113
/* FIFO command queue */
121
- uv_mutex_t cmd_mutex;
122
- struct metadata_database_cmdqueue cmd_queue;
114
+ SPINLOCK cmd_queue_lock;
115
+ struct metadata_cmd *cmd_base;
116
};
117
118
#define metadata_flag_check(target_flags, flag) (__atomic_load_n(&((target_flags)->flags), __ATOMIC_SEQ_CST) & (flag))
@@ -1058,21 +1051,15 @@ static void cleanup_health_log(struct metadata_wc *wc)
1051
// EVENT LOOP STARTS HERE
1052
//
1053
1061
-static void metadata_init_cmd_queue(struct metadata_wc *wc)
1062
-{
1063
- wc->cmd_queue.cmd_base = NULL;
1064
- fatal_assert(0 == uv_mutex_init(&wc->cmd_mutex));
1065
-}
1066
-
1054
static void metadata_free_cmd_queue(struct metadata_wc *wc)
1055
{
1069
- uv_mutex_lock(&wc->cmd_mutex);
1070
- while(wc->cmd_queue.cmd_base) {
1071
- struct metadata_cmd *t = wc->cmd_queue.cmd_base;
1072
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wc->cmd_queue.cmd_base, t, prev, next);
1056
+ spinlock_lock(&wc->cmd_queue_lock);
1057
+ while(wc->cmd_base) {
1058
+ struct metadata_cmd *t = wc->cmd_base;
1059
+ DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wc->cmd_base, t, prev, next);
1060
freez(t);
1061
}
1075
- uv_mutex_unlock(&wc->cmd_mutex);
1062
+ spinlock_unlock(&wc->cmd_queue_lock);
1063
}
1064
1065
static void metadata_enq_cmd(struct metadata_wc *wc, struct metadata_cmd *cmd)
@@ -1089,9 +1076,9 @@ static void metadata_enq_cmd(struct metadata_wc *wc, struct metadata_cmd *cmd)
1076
*t = *cmd;
1077
t->prev = t->next = NULL;
1078
1092
- uv_mutex_lock(&wc->cmd_mutex);
1093
- DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(wc->cmd_queue.cmd_base, t, prev, next);
1094
- uv_mutex_unlock(&wc->cmd_mutex);
1079
+ spinlock_lock(&wc->cmd_queue_lock);
1080
+ DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(wc->cmd_base, t, prev, next);
1081
+ spinlock_unlock(&wc->cmd_queue_lock);
1082
1083
wakeup_event_loop:
1084
(void) uv_async_send(&wc->async);
@@ -1101,10 +1088,10 @@ static struct metadata_cmd metadata_deq_cmd(struct metadata_wc *wc)
1088
{
1089
struct metadata_cmd ret;
1090
1104
- uv_mutex_lock(&wc->cmd_mutex);
1105
- if(wc->cmd_queue.cmd_base) {
1106
- struct metadata_cmd *t = wc->cmd_queue.cmd_base;
1107
- DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wc->cmd_queue.cmd_base, t, prev, next);
1091
+ spinlock_lock(&wc->cmd_queue_lock);
1092
+ if(wc->cmd_base) {
1093
+ struct metadata_cmd *t = wc->cmd_base;
1094
+ DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wc->cmd_base, t, prev, next);
1095
ret = *t;
1096
freez(t);
1097
}
@@ -1112,7 +1099,7 @@ static struct metadata_cmd metadata_deq_cmd(struct metadata_wc *wc)
1099
ret.opcode = METADATA_DATABASE_NOOP;
1100
ret.completion = NULL;
1101
}
1115
- uv_mutex_unlock(&wc->cmd_mutex);
1102
+ spinlock_unlock(&wc->cmd_queue_lock);
1103
1104
return ret;
1105
}
@@ -1183,13 +1170,6 @@ void run_metadata_cleanup(struct metadata_wc *wc)
1170
(void) sqlite3_wal_checkpoint(db_meta, NULL);
1171
}
1172
1186
-struct ml_model_payload {
1187
- uv_work_t request;
1188
- struct metadata_wc *wc;
1189
- Pvoid_t JudyL;
1190
- size_t count;
1191
-};
1192
-
1173
struct scan_metadata_payload {
1174
uv_work_t request;
1175
struct metadata_wc *wc;
@@ -1338,6 +1318,10 @@ static bool metadata_scan_host(RRDHOST *host, uint32_t max_count, bool use_trans
1318
bool more_to_do = false;
1319
uint32_t scan_count = 1;
1320
1321
+ sqlite3_stmt *ml_load_stmt = NULL;
1322
+
1323
+ bool load_ml_models = max_count;
1324
+
1325
if (use_transaction)
1326
(void)db_execute(db_meta, "BEGIN TRANSACTION");
1327
@@ -1382,6 +1366,14 @@ static bool metadata_scan_host(RRDHOST *host, uint32_t max_count, bool use_trans
1366
rrdhost_hostname(host), rrdset_name(st),
1367
rrddim_name(rd));
1368
}
1369
+
1370
+ if(rrddim_flag_check(rd, RRDDIM_FLAG_ML_MODEL_LOAD)) {
1371
+ rrddim_flag_clear(rd, RRDDIM_FLAG_ML_MODEL_LOAD);
1372
+ if (likely(load_ml_models))
1373
+ (void) ml_dimension_load_models(rd, &ml_load_stmt);
1374
+ }
1375
+
1376
+ worker_is_idle();
1377
}
1378
rrddim_foreach_done(rd);
1379
}
@@ -1390,6 +1382,11 @@ static bool metadata_scan_host(RRDHOST *host, uint32_t max_count, bool use_trans
1382
if (use_transaction)
1383
(void)db_execute(db_meta, "COMMIT TRANSACTION");
1384
1385
+ if (ml_load_stmt) {
1386
+ sqlite3_finalize(ml_load_stmt);
1387
+ ml_load_stmt = NULL;
1388
+ }
1389
+
1390
return more_to_do;
1391
}
1392
@@ -1519,49 +1516,6 @@ static void start_metadata_hosts(uv_work_t *req __maybe_unused)
1516
worker_is_idle();
1517
}
1518
1522
-// Callback after scan of hosts is done
1523
-static void after_start_ml_model_load(uv_work_t *req, int status __maybe_unused)
1524
-{
1525
- struct ml_model_payload *ml_data = req->data;
1526
- struct metadata_wc *wc = ml_data->wc;
1527
- metadata_flag_clear(wc, METADATA_FLAG_ML_LOADING);
1528
- JudyLFreeArray(&ml_data->JudyL, PJE0);
1529
- freez(ml_data);
1530
-}
1531
-
1532
-static void start_ml_model_load(uv_work_t *req __maybe_unused)
1533
-{
1534
- register_libuv_worker_jobs();
1535
-
1536
- struct ml_model_payload *ml_data = req->data;
1537
-
1538
- worker_is_busy(UV_EVENT_METADATA_ML_LOAD);
1539
-
1540
- Pvoid_t *PValue;
1541
- Word_t Index = 0;
1542
- bool first = true;
1543
- RRDDIM *rd;
1544
- RRDDIM_ACQUIRED *rda;
1545
- internal_error(true, "Batch ML load loader, %zu items", ml_data->count);
1546
-
1547
- sqlite3_stmt *ml_load_stmt = NULL;
1548
- while((PValue = JudyLFirstThenNext(ml_data->JudyL, &Index, &first))) {
1549
- UNUSED(PValue);
1550
- rda = (RRDDIM_ACQUIRED *) Index;
1551
- rd = rrddim_acquired_to_rrddim(rda);
1552
- ml_dimension_load_models(rd, &ml_load_stmt);
1553
- rrddim_acquired_release(rda);
1554
- }
1555
-
1556
- if (ml_load_stmt) {
1557
- sqlite3_finalize(ml_load_stmt);
1558
- ml_load_stmt = NULL;
1559
- }
1560
- worker_is_idle();
1561
-}
1562
-
1563
-
1564
-
1519
static void metadata_event_loop(void *arg)
1520
{
1521
worker_register("METASYNC");
@@ -1571,7 +1525,6 @@ static void metadata_event_loop(void *arg)
1525
worker_register_job_name(METADATA_STORE_CLAIM_ID, "add claim id");
1526
worker_register_job_name(METADATA_ADD_HOST_INFO, "add host info");
1527
worker_register_job_name(METADATA_MAINTENANCE, "maintenance");
1574
- worker_register_job_name(METADATA_ML_LOAD_MODELS, "ml load models");
1528
1529
int ret;
1530
uv_loop_t *loop;
@@ -1612,11 +1565,10 @@ static void metadata_event_loop(void *arg)
1565
wc->metadata_check_after = now_realtime_sec() + METADATA_HOST_CHECK_FIRST_CHECK;
1566
1567
int shutdown = 0;
1615
- completion_mark_complete(&wc->init_complete);
1568
+ completion_mark_complete(&wc->start_stop_complete);
1569
BUFFER *work_buffer = buffer_create(1024, &netdata_buffers_statistics.buffers_sqlite);
1570
struct scan_metadata_payload *data;
1571
1619
- struct ml_model_payload *ml_data = NULL;
1572
while (shutdown == 0 || (wc->flags & METADATA_FLAG_PROCESSING)) {
1573
uuid_t *uuid;
1574
RRDHOST *host = NULL;
@@ -1643,43 +1595,10 @@ static void metadata_event_loop(void *arg)
1595
if (likely(opcode != METADATA_DATABASE_NOOP))
1596
worker_is_busy(opcode);
1597
1646
- // Have pending ML models to load?
1647
- if (opcode != METADATA_ML_LOAD_MODELS && ml_data && ml_data->count) {
1648
- static usec_t ml_submit_last = 0;
1649
- usec_t now = now_monotonic_usec();
1650
- if (!ml_submit_last)
1651
- ml_submit_last = now;
1652
-
1653
- if (!metadata_flag_check(wc, METADATA_FLAG_ML_LOADING) && (now - ml_submit_last > 150 * USEC_PER_MS)) {
1654
- metadata_flag_set(wc, METADATA_FLAG_ML_LOADING);
1655
- if (unlikely(uv_queue_work(loop, &ml_data->request, start_ml_model_load, after_start_ml_model_load)))
1656
- metadata_flag_clear(wc, METADATA_FLAG_ML_LOADING);
1657
- else {
1658
- ml_submit_last = now;
1659
- ml_data = NULL;
1660
- }
1661
- }
1662
- }
1663
-
1598
switch (opcode) {
1599
case METADATA_DATABASE_NOOP:
1600
case METADATA_DATABASE_TIMER:
1601
break;
1668
-
1669
- case METADATA_ML_LOAD_MODELS: {
1670
- RRDDIM *rd = (RRDDIM *) cmd.param[0];
1671
- RRDDIM_ACQUIRED *rda = rrddim_find_and_acquire(rd->rrdset, rrddim_id(rd));
1672
- if (likely(rda)) {
1673
- if (!ml_data) {
1674
- ml_data = callocz(1,sizeof(*ml_data));
1675
- ml_data->request.data = ml_data;
1676
- ml_data->wc = wc;
1677
- }
1678
- JudyLIns(&ml_data->JudyL, (Word_t)rda, PJE0);
1679
- ml_data->count++;
1680
- }
1681
- break;
1682
- }
1602
case METADATA_DEL_DIMENSION:
1603
uuid = (uuid_t *) cmd.param[0];
1604
if (likely(dimension_can_be_deleted(uuid, NULL, false)))
@@ -1765,8 +1684,8 @@ static void metadata_event_loop(void *arg)
1684
freez(loop);
1685
worker_unregister();
1686
1768
- netdata_log_info("METADATA: Shutting down event loop");
1769
- completion_mark_complete(&wc->init_complete);
1687
+ netdata_log_info("Shutting down event loop");
1688
+ completion_mark_complete(&wc->start_stop_complete);
1689
if (wc->scan_complete) {
1690
completion_destroy(wc->scan_complete);
1691
freez(wc->scan_complete);
@@ -1787,7 +1706,7 @@ struct metadata_wc metasync_worker = {.loop = NULL};
1706
1707
void metadata_sync_shutdown(void)
1708
{
1790
- completion_init(&metasync_worker.init_complete);
1709
+ completion_init(&metasync_worker.start_stop_complete);
1710
1711
struct metadata_cmd cmd;
1712
memset(&cmd, 0, sizeof(cmd));
@@ -1797,8 +1716,8 @@ void metadata_sync_shutdown(void)
1716
1717
/* wait for metadata thread to shut down */
1718
netdata_log_info("METADATA: Waiting for shutdown ACK");
1800
- completion_wait_for(&metasync_worker.init_complete);
1801
- completion_destroy(&metasync_worker.init_complete);
1719
+ completion_wait_for(&metasync_worker.start_stop_complete);
1720
+ completion_destroy(&metasync_worker.start_stop_complete);
1721
netdata_log_info("METADATA: Shutdown complete");
1722
}
1723
@@ -1840,13 +1759,12 @@ void metadata_sync_init(void)
1759
struct metadata_wc *wc = &metasync_worker;
1760
1761
memset(wc, 0, sizeof(*wc));
1843
- metadata_init_cmd_queue(wc);
1844
- completion_init(&wc->init_complete);
1762
+ completion_init(&wc->start_stop_complete);
1763
1764
fatal_assert(0 == uv_thread_create(&(wc->thread), metadata_event_loop, wc));
1765
1848
- completion_wait_for(&wc->init_complete);
1849
- completion_destroy(&wc->init_complete);
1766
+ completion_wait_for(&wc->start_stop_complete);
1767
+ completion_destroy(&wc->start_stop_complete);
1768
1769
netdata_log_info("SQLite metadata sync initialization complete");
1770
}
@@ -1899,9 +1817,7 @@ void metaqueue_host_update_info(RRDHOST *host)
1817
1818
void metaqueue_ml_load_models(RRDDIM *rd)
1819
{
1902
- if (unlikely(!metasync_worker.loop))
1903
- return;
1904
- queue_metadata_cmd(METADATA_ML_LOAD_MODELS, rd, NULL);
1820
+ rrddim_flag_set(rd, RRDDIM_FLAG_ML_MODEL_LOAD);
1821
}
1822
1823
void metadata_queue_load_host_context(RRDHOST *host)
streaming/receiver.c
-3
@@ -321,9 +321,6 @@ static inline bool receiver_should_stop(struct receiver_state *rpt) {
321
return false;
322
}
323
324
-extern pthread_key_t plugin_key;
325
-struct plugin_data;
326
-
324
static size_t streaming_parser(struct receiver_state *rpt, struct plugind *cd, int fd, void *ssl) {
325
size_t result = 0;
326