Schedule retention message calculation to a worker thread (#13039)
* Move aclk_update_retention to the proper header file * Do a scan but avoid going through all the dimensions if we have too much to delete -- do not generate a retention message in that case * Schedule the retention calculation to a worker * Adjust messages in the access log * Fix compilation errors with --disable-cloud
Stelios Fragkakis committed
Jun 1, 2022 at 19:10 UTC
c261a771cc0c93fe4e9fbb83e1be141406d314be
6 files changed
+87
-14
database/sqlite/sqlite_aclk.c
+70
-5
@@ -34,6 +34,28 @@ const char *aclk_sync_config[] = {
34
35
uv_mutex_t aclk_async_lock;
36
struct aclk_database_worker_config *aclk_thread_head = NULL;
37
+int retention_running = 0;
38
+
39
+#ifdef ENABLE_NEW_CLOUD_PROTOCOL
40
+static void stop_retention_run()
41
+{
42
+ uv_mutex_lock(&aclk_async_lock);
43
+ retention_running = 0;
44
+ uv_mutex_unlock(&aclk_async_lock);
45
+}
46
+
47
+static int request_retention_run()
48
+{
49
+ int rc = 0;
50
+ uv_mutex_lock(&aclk_async_lock);
51
+ if (unlikely(retention_running))
52
+ rc = 1;
53
+ else
54
+ retention_running = 1;
55
+ uv_mutex_unlock(&aclk_async_lock);
56
+ return rc;
57
+}
58
+#endif
59
60
int claimed()
61
{
@@ -318,9 +340,6 @@ static void timer_cb(uv_timer_t* handle)
340
341
if (aclk_use_new_cloud_arch && aclk_connected) {
342
if (wc->rotation_after && wc->rotation_after < now) {
321
- cmd.opcode = ACLK_DATABASE_NODE_INFO;
322
- aclk_database_enq_cmd_noblock(wc, &cmd);
323
-
343
cmd.opcode = ACLK_DATABASE_UPD_RETENTION;
344
if (!aclk_database_enq_cmd_noblock(wc, &cmd))
345
wc->rotation_after += ACLK_DATABASE_ROTATION_INTERVAL;
@@ -353,6 +372,38 @@ static void timer_cb(uv_timer_t* handle)
372
#endif
373
}
374
375
+
376
+#ifdef ENABLE_NEW_CLOUD_PROTOCOL
377
+void after_send_retention(uv_work_t *req, int status)
378
+{
379
+ struct aclk_database_worker_config *wc = req->data;
380
+ (void)status;
381
+ stop_retention_run();
382
+ wc->retention_running = 0;
383
+
384
+ struct aclk_database_cmd cmd;
385
+ memset(&cmd, 0, sizeof(cmd));
386
+ cmd.opcode = ACLK_DATABASE_DIM_DELETION;
387
+ if (aclk_database_enq_cmd_noblock(wc, &cmd))
388
+ info("Failed to queue a dimension deletion message");
389
+
390
+ cmd.opcode = ACLK_DATABASE_NODE_INFO;
391
+ if (aclk_database_enq_cmd_noblock(wc, &cmd))
392
+ info("Failed to queue a node update info message");
393
+}
394
+
395
+
396
+static void send_retention(uv_work_t *req)
397
+{
398
+ struct aclk_database_worker_config *wc = req->data;
399
+
400
+ if (unlikely(wc->is_shutting_down))
401
+ return;
402
+
403
+ aclk_update_retention(wc);
404
+}
405
+#endif
406
+
407
#define MAX_CMD_BATCH_SIZE (256)
408
409
void aclk_database_worker(void *arg)
@@ -429,6 +480,7 @@ void aclk_database_worker(void *arg)
480
481
memset(&cmd, 0, sizeof(cmd));
482
#ifdef ENABLE_NEW_CLOUD_PROTOCOL
483
+ uv_work_t retention_work;
484
sql_get_last_chart_sequence(wc);
485
wc->chart_payload_count = sql_get_pending_count(wc);
486
if (!wc->chart_payload_count)
@@ -440,6 +492,7 @@ void aclk_database_worker(void *arg)
492
wc->rotation_after = wc->startup_time + ACLK_DATABASE_ROTATION_DELAY;
493
494
debug(D_ACLK_SYNC,"Node %s reports pending message count = %u", wc->node_id, wc->chart_payload_count);
495
+
496
while (likely(!netdata_exit)) {
497
worker_is_idle();
498
uv_run(loop, UV_RUN_DEFAULT);
@@ -538,9 +591,21 @@ void aclk_database_worker(void *arg)
591
aclk_process_dimension_deletion(wc, cmd);
592
break;
593
case ACLK_DATABASE_UPD_RETENTION:
594
+ if (unlikely(wc->retention_running))
595
+ break;
596
+
597
+ if (unlikely(request_retention_run())) {
598
+ wc->rotation_after = now_realtime_sec() + ACLK_DATABASE_RETENTION_RETRY;
599
+ break;
600
+ }
601
+
602
debug(D_ACLK_SYNC,"Sending retention info for %s", wc->uuid_str);
542
- aclk_update_retention(wc, cmd);
543
- aclk_process_dimension_deletion(wc, cmd);
603
+ retention_work.data = wc;
604
+ wc->retention_running = 1;
605
+ if (unlikely(uv_queue_work(loop, &retention_work, send_retention, after_send_retention))) {
606
+ wc->retention_running = 0;
607
+ stop_retention_run();
608
+ }
609
break;
610
611
// NODE_INSTANCE DETECTION
database/sqlite/sqlite_aclk.h
+2
@@ -17,6 +17,7 @@
17
#define ACLK_MAX_ALERT_UPDATES (5)
18
#define ACLK_DATABASE_CLEANUP_FIRST (60)
19
#define ACLK_DATABASE_ROTATION_DELAY (180)
20
+#define ACLK_DATABASE_RETENTION_RETRY (60)
21
#define ACLK_DATABASE_CLEANUP_INTERVAL (3600)
22
#define ACLK_DATABASE_ROTATION_INTERVAL (3600)
23
#define ACLK_DELETE_ACK_INTERNAL (600)
@@ -197,6 +198,7 @@ struct aclk_database_worker_config {
198
int node_info_send;
199
int chart_pending;
200
int chart_reset_count;
201
+ int retention_running;
202
volatile unsigned is_shutting_down;
203
volatile unsigned is_orphan;
204
struct aclk_database_worker_config *next;
database/sqlite/sqlite_aclk_alert.c
+1
-1
@@ -773,7 +773,7 @@ void sql_process_queue_removed_alerts_to_aclk(struct aclk_database_worker_config
773
774
db_execute(buffer_tostring(sql));
775
776
- log_access("ACLK STA [%s (%s)]: Queued removed alerts.", wc->node_id, wc->host ? wc->host->hostname : "N/A");
776
+ log_access("ACLK STA [%s (%s)]: QUEUED REMOVED ALERTS", wc->node_id, wc->host ? wc->host->hostname : "N/A");
777
778
buffer_free(sql);
779
database/sqlite/sqlite_aclk_chart.c
+13
-7
@@ -566,7 +566,7 @@ void aclk_receive_chart_ack(struct aclk_database_worker_config *wc, struct aclk_
566
error_report("Failed to ACK sequence id, rc = %d", rc);
567
else
568
log_access(
569
- "ACLK STA [%s (%s)]: CHARTS ACKNOWLEDGED in the database upto %" PRIu64,
569
+ "ACLK STA [%s (%s)]: CHARTS ACKNOWLEDGED IN THE DATABASE UP TO %" PRIu64,
570
wc->node_id,
571
wc->host ? wc->host->hostname : "N/A",
572
cmd.param1);
@@ -847,9 +847,8 @@ failed:
847
"SELECT distinct h.host_id, c.update_every, c.type||'.'||c.id FROM chart c, host h " \
848
"WHERE c.host_id = h.host_id AND c.host_id = @host_id ORDER BY c.update_every ASC;"
849
850
-void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
850
+void aclk_update_retention(struct aclk_database_worker_config *wc)
851
{
852
- UNUSED(cmd);
852
int rc;
853
854
if (!aclk_use_new_cloud_arch || !aclk_connected)
@@ -916,7 +915,9 @@ void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_d
915
rotate_data.node_id = strdupz(wc->node_id);
916
917
time_t now = now_realtime_sec();
919
- while (sqlite3_step(res) == SQLITE_ROW) {
918
+ while (sqlite3_step(res) == SQLITE_ROW && dimension_update_count < ACLK_MAX_DIMENSION_CLEANUP) {
919
+ if (unlikely(netdata_exit))
920
+ break;
921
if (!update_every || update_every != (uint32_t)sqlite3_column_int(res, 1)) {
922
if (update_every) {
923
debug(D_ACLK_SYNC, "Update %s for %u oldest time = %ld", wc->host_guid, update_every, start_time);
@@ -1003,8 +1004,12 @@ void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_d
1004
if (!wc->host)
1005
hostname = get_hostname_by_node_id(wc->node_id);
1006
1006
- log_access("ACLK STA [%s (%s)]: UPDATES %d RETENTION MESSAGE SENT. CHECKED %u DIMENSIONS. %u DELETED, %u STOPPED COLLECTING",
1007
- wc->node_id, wc->host ? wc->host->hostname : hostname ? hostname : "N/A", wc->chart_updates, total_checked, total_deleted, total_stopped);
1007
+ if (dimension_update_count < ACLK_MAX_DIMENSION_CLEANUP && !netdata_exit)
1008
+ log_access("ACLK STA [%s (%s)]: UPDATES %d RETENTION MESSAGE SENT. CHECKED %u DIMENSIONS. %u DELETED, %u STOPPED COLLECTING",
1009
+ wc->node_id, wc->host ? wc->host->hostname : hostname ? hostname : "N/A", wc->chart_updates, total_checked, total_deleted, total_stopped);
1010
+ else
1011
+ log_access("ACLK STA [%s (%s)]: UPDATES %d RETENTION MESSAGE NOT SENT. CHECKED %u DIMENSIONS. %u DELETED, %u STOPPED COLLECTING",
1012
+ wc->node_id, wc->host ? wc->host->hostname : hostname ? hostname : "N/A", wc->chart_updates, total_checked, total_deleted, total_stopped);
1013
freez(hostname);
1014
1015
#ifdef NETDATA_INTERNAL_CHECKS
@@ -1017,7 +1022,8 @@ void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_d
1022
rotate_data.interval_durations[i].update_every,
1023
rotate_data.interval_durations[i].retention);
1024
#endif
1020
- aclk_retention_updated(&rotate_data);
1025
+ if (dimension_update_count < ACLK_MAX_DIMENSION_CLEANUP && !netdata_exit)
1026
+ aclk_retention_updated(&rotate_data);
1027
freez(rotate_data.node_id);
1028
freez(rotate_data.interval_durations);
1029
database/sqlite/sqlite_aclk_chart.h
+1
@@ -67,4 +67,5 @@ uint32_t sql_get_pending_count(struct aclk_database_worker_config *wc);
67
void aclk_send_dimension_update(RRDDIM *rd);
68
struct aclk_chart_sync_stats *aclk_get_chart_sync_stats(RRDHOST *host);
69
void sql_check_chart_liveness(RRDSET *st);
70
+void aclk_update_retention(struct aclk_database_worker_config *wc);
71
#endif //NETDATA_SQLITE_ACLK_CHART_H
database/sqlite/sqlite_aclk_node.h
-1
@@ -4,5 +4,4 @@
4
#define NETDATA_SQLITE_ACLK_NODE_H
5
6
void sql_build_node_info(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
7
-void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
7
#endif //NETDATA_SQLITE_ACLK_NODE_H