@cryptotaxi247 / netdata-1 / commits / 12cd39889

Faster streaming by 25% on the child (#13708)

* faster printing of BEGIN, SET, END * fewer conditions * faster buffer_fast_strcat() * faster buffer_fast_strcat() fix * eliminate atomic operations and conditions in the BEGIN, SET, END flow * removed unecessary condition

Costa Tsaousis committed Sep 24, 2022 at 14:10 UTC 12cd3988949b604af74c28c6424b7753cf3a5069
4 files changed +76 -45
database/rrddim.c
+3 -3
@@ -363,8 +363,8 @@ inline int rrddim_set_algorithm(RRDSET *st, RRDDIM *rd, RRD_ALGORITHM algorithm)
363 debug(D_RRD_CALLS, "Updating algorithm of dimension '%s/%s' from %s to %s", rrdset_id(st), rrddim_name(rd), rrd_algorithm_name(rd->algorithm), rrd_algorithm_name(algorithm));
364 rd->algorithm = algorithm;
365 rd->exposed = 0;
366 - rrdset_flag_set(st, RRDSET_FLAG_HOMOGENEOUS_CHECK);
366 rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
367 + rrdset_flag_set(st, RRDSET_FLAG_HOMOGENEOUS_CHECK);
368 rrdcontext_updated_rrddim_algorithm(rd);
369 return 1;
370 }
@@ -376,8 +376,8 @@ inline int rrddim_set_multiplier(RRDSET *st, RRDDIM *rd, collected_number multip
376 debug(D_RRD_CALLS, "Updating multiplier of dimension '%s/%s' from " COLLECTED_NUMBER_FORMAT " to " COLLECTED_NUMBER_FORMAT, rrdset_id(st), rrddim_name(rd), rd->multiplier, multiplier);
377 rd->multiplier = multiplier;
378 rd->exposed = 0;
379 - rrdset_flag_set(st, RRDSET_FLAG_HOMOGENEOUS_CHECK);
379 rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
380 + rrdset_flag_set(st, RRDSET_FLAG_HOMOGENEOUS_CHECK);
381 rrdcontext_updated_rrddim_multiplier(rd);
382 return 1;
383 }
@@ -389,8 +389,8 @@ inline int rrddim_set_divisor(RRDSET *st, RRDDIM *rd, collected_number divisor)
389 debug(D_RRD_CALLS, "Updating divisor of dimension '%s/%s' from " COLLECTED_NUMBER_FORMAT " to " COLLECTED_NUMBER_FORMAT, rrdset_id(st), rrddim_name(rd), rd->divisor, divisor);
390 rd->divisor = divisor;
391 rd->exposed = 0;
392 - rrdset_flag_set(st, RRDSET_FLAG_HOMOGENEOUS_CHECK);
392 rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
393 + rrdset_flag_set(st, RRDSET_FLAG_HOMOGENEOUS_CHECK);
394 rrdcontext_updated_rrddim_divisor(rd);
395 return 1;
396 }
libnetdata/buffer/buffer.c
+12
@@ -136,6 +136,18 @@ void buffer_print_llu(BUFFER *wb, unsigned long long uvalue)
136 wb->len += wstr - str;
137 }
138
139 +void buffer_print_ll(BUFFER *wb, long long value)
140 +{
141 + buffer_need_bytes(wb, 50);
142 +
143 + if(value < 0) {
144 + buffer_fast_strcat(wb, "-", 1);
145 + value = -value;
146 + }
147 +
148 + buffer_print_llu(wb, value);
149 +}
150 +
151 void buffer_fast_strcat(BUFFER *wb, const char *txt, size_t len) {
152 if(unlikely(!txt || !*txt)) return;
153
libnetdata/buffer/buffer.h
+1
@@ -78,6 +78,7 @@ extern char *print_number_llu_r(char *str, unsigned long long uvalue);
78 extern char *print_number_llu_r_smart(char *str, unsigned long long uvalue);
79
80 extern void buffer_print_llu(BUFFER *wb, unsigned long long uvalue);
81 +extern void buffer_print_ll(BUFFER *wb, long long value);
82
83 static inline void buffer_need_bytes(BUFFER *buffer, size_t needed_free_size) {
84 if(unlikely(buffer->size - buffer->len < needed_free_size))
streaming/rrdpush.c
+60 -42
@@ -129,29 +129,44 @@ int rrdpush_init() {
129 unsigned int remote_clock_resync_iterations = 60;
130
131
132 -static inline int should_send_chart_matching(RRDSET *st) {
133 - // Do not stream anomaly rates charts.
134 - if (unlikely(rrdset_is_ar_chart(st)))
135 - return false;
136 -
137 - if (rrdset_flag_check(st, RRDSET_FLAG_ANOMALY_DETECTION))
138 - return ml_streaming_enabled();
132 +static inline bool should_send_chart_matching(RRDSET *st) {
133 + RRDSET_FLAGS flags = rrdset_flag_check(st, RRDSET_FLAG_UPSTREAM_SEND|RRDSET_FLAG_UPSTREAM_IGNORE);
134
140 - if(!rrdset_flag_check(st, RRDSET_FLAG_UPSTREAM_SEND|RRDSET_FLAG_UPSTREAM_IGNORE)) {
135 + if(unlikely(!flags)) {
136 RRDHOST *host = st->rrdhost;
137
143 - if(simple_pattern_matches(host->rrdpush_send_charts_matching, rrdset_id(st)) ||
138 + // Do not stream anomaly rates charts.
139 + if (unlikely(rrdset_is_ar_chart(st))) {
140 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_SEND);
141 + rrdset_flag_set(st, RRDSET_FLAG_UPSTREAM_IGNORE);
142 + flags = RRDSET_FLAG_UPSTREAM_IGNORE;
143 + }
144 + else if (rrdset_flag_check(st, RRDSET_FLAG_ANOMALY_DETECTION)) {
145 + if(ml_streaming_enabled()) {
146 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_IGNORE);
147 + rrdset_flag_set(st, RRDSET_FLAG_UPSTREAM_SEND);
148 + flags = RRDSET_FLAG_UPSTREAM_SEND;
149 + }
150 + else {
151 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_SEND);
152 + rrdset_flag_set(st, RRDSET_FLAG_UPSTREAM_IGNORE);
153 + flags = RRDSET_FLAG_UPSTREAM_IGNORE;
154 + }
155 + }
156 + else if(simple_pattern_matches(host->rrdpush_send_charts_matching, rrdset_id(st)) ||
157 simple_pattern_matches(host->rrdpush_send_charts_matching, rrdset_name(st))) {
158 rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_IGNORE);
159 rrdset_flag_set(st, RRDSET_FLAG_UPSTREAM_SEND);
160 + flags = RRDSET_FLAG_UPSTREAM_SEND;
161 }
162 else {
163 rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_SEND);
164 rrdset_flag_set(st, RRDSET_FLAG_UPSTREAM_IGNORE);
165 + flags = RRDSET_FLAG_UPSTREAM_IGNORE;
166 }
167 }
168
154 - return(rrdset_flag_check(st, RRDSET_FLAG_UPSTREAM_SEND));
169 + return flags & RRDSET_FLAG_UPSTREAM_SEND;
170 }
171
172 int configured_as_parent() {
@@ -173,22 +188,7 @@ int configured_as_parent() {
188 return is_parent;
189 }
190
176 -// checks if the current chart definition has been sent
177 -static inline int need_to_send_chart_definition(RRDSET *st) {
178 - if(unlikely(!(rrdset_flag_check(st, RRDSET_FLAG_UPSTREAM_EXPOSED))))
179 - return 1;
180 -
181 - RRDDIM *rd;
182 - dfe_start_read(st->rrddim_root_index, rd) {
183 - if(unlikely(!rd->exposed)) {
184 - internal_error(true, "host '%s', chart '%s', dimension '%s' flag 'exposed' triggered chart refresh to upstream", rrdhost_hostname(st->rrdhost), rrdset_id(st), rrddim_id(rd));
185 - return 1;
186 - }
187 - }
188 - dfe_done(rd);
189 -
190 - return 0;
191 -}
191 +#define need_to_send_chart_definition(st) (!rrdset_flag_check(st, RRDSET_FLAG_UPSTREAM_EXPOSED))
192
193 // chart labels
194 static int send_clabels_callback(const char *name, const char *value, RRDLABEL_SRC ls, void *data) {
@@ -276,22 +276,42 @@ static inline void rrdpush_send_chart_definition(RRDSET *st) {
276 // sends the current chart dimensions
277 static inline bool rrdpush_send_chart_metrics_nolock(RRDSET *st, struct sender_state *s) {
278 RRDHOST *host = st->rrdhost;
279 - buffer_sprintf(host->sender->build, "BEGIN \"%s\" %llu", rrdset_id(st), (st->last_collected_time.tv_sec > st->upstream_resync_time)?st->usec_since_last_update:0);
280 - if (s->version >= VERSION_GAP_FILLING)
281 - buffer_sprintf(host->sender->build, " %"PRId64"\n", (int64_t)st->last_collected_time.tv_sec);
282 - else
283 - buffer_strcat(host->sender->build, "\n");
279 + BUFFER *wb = host->sender->build;
280 +
281 + buffer_fast_strcat(wb, "BEGIN \"", 7);
282 + buffer_fast_strcat(wb, rrdset_id(st), string_strlen(st->id));
283 + buffer_fast_strcat(wb, "\" ", 2);
284 + buffer_print_llu(wb, (st->last_collected_time.tv_sec > st->upstream_resync_time)?st->usec_since_last_update:0);
285 +
286 + if (s->version >= VERSION_GAP_FILLING) {
287 + buffer_fast_strcat(wb, " ", 1);
288 + buffer_print_ll(wb, st->last_collected_time.tv_sec);
289 + }
290 +
291 + buffer_fast_strcat(wb, "\n", 1);
292
293 size_t count_of_dimensions_written = 0;
294 RRDDIM *rd;
295 rrddim_foreach_read(rd, st) {
288 - if(rd->updated && rd->exposed) {
289 - buffer_sprintf(host->sender->build, "SET \"%s\" = " COLLECTED_NUMBER_FORMAT "\n", rrddim_id(rd), rd->collected_value);
296 + if(unlikely(!rd->updated))
297 + continue;
298 +
299 + if(likely(rd->exposed)) {
300 + buffer_fast_strcat(wb, "SET \"", 5);
301 + buffer_fast_strcat(wb, rrddim_id(rd), string_strlen(rd->id));
302 + buffer_fast_strcat(wb, "\" = ", 4);
303 + buffer_print_ll(wb, rd->collected_value);
304 + buffer_fast_strcat(wb, "\n", 1);
305 count_of_dimensions_written++;
306 }
307 + else {
308 + internal_error(true, "host '%s', chart '%s', dimension '%s' flag 'exposed' is updated but not exposed", rrdhost_hostname(st->rrdhost), rrdset_id(st), rrddim_id(rd));
309 + // we will include it in the next iteration
310 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
311 + }
312 }
313 rrddim_foreach_done(rd);
294 - buffer_strcat(host->sender->build, "END\n");
314 + buffer_fast_strcat(wb, "END\n", 4);
315
316 return count_of_dimensions_written != 0;
317 }
@@ -350,15 +370,16 @@ void rrdset_done_push(RRDSET *st) {
370
371 RRDHOST *host = st->rrdhost;
372
353 - if(unlikely(host->rrdpush_send_enabled && !host->rrdpush_sender_spawn))
354 - rrdpush_sender_thread_spawn(host);
355 -
373 // Handle non-connected case
374 if(unlikely(!__atomic_load_n(&host->rrdpush_sender_connected, __ATOMIC_SEQ_CST)
375 || !rrdhost_flag_check(host, RRDHOST_FLAG_STREAM_COLLECTED_METRICS))) {
376
377 + if(unlikely(host->rrdpush_send_enabled && !host->rrdpush_sender_spawn))
378 + rrdpush_sender_thread_spawn(host);
379 +
380 if(unlikely(!host->rrdpush_sender_error_shown))
381 error("STREAM %s [send]: not ready - collected metrics are not sent to parent.", rrdhost_hostname(host));
382 +
383 host->rrdpush_sender_error_shown = 1;
384
385 return;
@@ -368,15 +389,12 @@ void rrdset_done_push(RRDSET *st) {
389 host->rrdpush_sender_error_shown = 0;
390 }
391
371 - if(dictionary_entries(st->rrddim_root_index) == 0)
372 - return;
373 -
392 sender_start(host->sender);
393
376 - if(need_to_send_chart_definition(st))
394 + if(unlikely(need_to_send_chart_definition(st)))
395 rrdpush_send_chart_definition(st);
396
379 - if(rrdpush_send_chart_metrics_nolock(st, host->sender)) {
397 + if(likely(rrdpush_send_chart_metrics_nolock(st, host->sender))) {
398 // signal the sender there are more data
399 if (host->rrdpush_sender_pipe[PIPE_WRITE] != -1 && write(host->rrdpush_sender_pipe[PIPE_WRITE], " ", 1) == -1)
400 error("STREAM %s [send]: cannot write to internal pipe", rrdhost_hostname(host));