3
#include "query.h"
4
#include "web/api/formatters/rrd2json.h"
5
#include "rrdr.h"
6
-#include "database/ram/rrddim_mem.h"
6
7
#include "average/average.h"
8
#include "countif/countif.h"
49
// continue after a flush as if nothing changed, for others a
50
// cleanup of the internal structures may be required).
51
NETDATA_DOUBLE (*flush)(struct rrdresult *r, RRDR_VALUE_FLAGS *rrdr_value_options_ptr);
52
+
53
+ TIER_QUERY_FETCH tier_query_fetch;
54
} api_v1_data_groups[] = {
55
{.name = "average",
56
.hash = 0,
60
.reset = grouping_reset_average,
61
.free = grouping_free_average,
62
.add = grouping_add_average,
62
- .flush = grouping_flush_average
63
+ .flush = grouping_flush_average,
64
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
65
},
66
{.name = "mean", // alias on 'average'
67
.hash = 0,
71
.reset = grouping_reset_average,
72
.free = grouping_free_average,
73
.add = grouping_add_average,
72
- .flush = grouping_flush_average
74
+ .flush = grouping_flush_average,
75
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
76
},
77
{.name = "incremental_sum",
78
.hash = 0,
82
.reset = grouping_reset_incremental_sum,
83
.free = grouping_free_incremental_sum,
84
.add = grouping_add_incremental_sum,
82
- .flush = grouping_flush_incremental_sum
85
+ .flush = grouping_flush_incremental_sum,
86
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
87
},
88
{.name = "incremental-sum",
89
.hash = 0,
93
.reset = grouping_reset_incremental_sum,
94
.free = grouping_free_incremental_sum,
95
.add = grouping_add_incremental_sum,
92
- .flush = grouping_flush_incremental_sum
96
+ .flush = grouping_flush_incremental_sum,
97
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
98
},
99
{.name = "median",
100
.hash = 0,
104
.reset = grouping_reset_median,
105
.free = grouping_free_median,
106
.add = grouping_add_median,
102
- .flush = grouping_flush_median
107
+ .flush = grouping_flush_median,
108
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
109
},
110
{.name = "min",
111
.hash = 0,
115
.reset = grouping_reset_min,
116
.free = grouping_free_min,
117
.add = grouping_add_min,
112
- .flush = grouping_flush_min
118
+ .flush = grouping_flush_min,
119
+ .tier_query_fetch = TIER_QUERY_FETCH_MIN
120
},
121
{.name = "max",
122
.hash = 0,
126
.reset = grouping_reset_max,
127
.free = grouping_free_max,
128
.add = grouping_add_max,
122
- .flush = grouping_flush_max
129
+ .flush = grouping_flush_max,
130
+ .tier_query_fetch = TIER_QUERY_FETCH_MAX
131
},
132
{.name = "sum",
133
.hash = 0,
137
.reset = grouping_reset_sum,
138
.free = grouping_free_sum,
139
.add = grouping_add_sum,
132
- .flush = grouping_flush_sum
140
+ .flush = grouping_flush_sum,
141
+ .tier_query_fetch = TIER_QUERY_FETCH_SUM
142
},
143
144
// standard deviation
150
.reset = grouping_reset_stddev,
151
.free = grouping_free_stddev,
152
.add = grouping_add_stddev,
144
- .flush = grouping_flush_stddev
153
+ .flush = grouping_flush_stddev,
154
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
155
},
156
{.name = "cv", // coefficient variation is calculated by stddev
157
.hash = 0,
161
.reset = grouping_reset_stddev, // not an error, stddev calculates this too
162
.free = grouping_free_stddev, // not an error, stddev calculates this too
163
.add = grouping_add_stddev, // not an error, stddev calculates this too
154
- .flush = grouping_flush_coefficient_of_variation
164
+ .flush = grouping_flush_coefficient_of_variation,
165
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
166
},
167
{.name = "rsd", // alias of 'cv'
168
.hash = 0,
172
.reset = grouping_reset_stddev, // not an error, stddev calculates this too
173
.free = grouping_free_stddev, // not an error, stddev calculates this too
174
.add = grouping_add_stddev, // not an error, stddev calculates this too
164
- .flush = grouping_flush_coefficient_of_variation
175
+ .flush = grouping_flush_coefficient_of_variation,
176
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
177
},
178
179
/*
185
.reset = grouping_reset_stddev,
186
.free = grouping_free_stddev,
187
.add = grouping_add_stddev,
176
- .flush = grouping_flush_mean
188
+ .flush = grouping_flush_mean,
189
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
190
},
191
*/
192
199
.reset = grouping_reset_stddev,
200
.free = grouping_free_stddev,
201
.add = grouping_add_stddev,
189
- .flush = grouping_flush_variance
202
+ .flush = grouping_flush_variance,
203
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
204
},
205
*/
206
208
{.name = "ses",
209
.hash = 0,
210
.value = RRDR_GROUPING_SES,
197
- .init = grouping_init_ses,
211
+ .init = grouping_init_ses,
212
.create= grouping_create_ses,
213
.reset = grouping_reset_ses,
214
.free = grouping_free_ses,
215
.add = grouping_add_ses,
202
- .flush = grouping_flush_ses
216
+ .flush = grouping_flush_ses,
217
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
218
},
219
{.name = "ema", // alias for 'ses'
220
.hash = 0,
221
.value = RRDR_GROUPING_SES,
207
- .init = NULL,
222
+ .init = NULL,
223
.create= grouping_create_ses,
224
.reset = grouping_reset_ses,
225
.free = grouping_free_ses,
226
.add = grouping_add_ses,
212
- .flush = grouping_flush_ses
227
+ .flush = grouping_flush_ses,
228
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
229
},
230
{.name = "ewma", // alias for ses
231
.hash = 0,
232
.value = RRDR_GROUPING_SES,
217
- .init = NULL,
233
+ .init = NULL,
234
.create= grouping_create_ses,
235
.reset = grouping_reset_ses,
236
.free = grouping_free_ses,
237
.add = grouping_add_ses,
222
- .flush = grouping_flush_ses
238
+ .flush = grouping_flush_ses,
239
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
240
},
241
242
// double exponential smoothing
243
{.name = "des",
244
.hash = 0,
245
.value = RRDR_GROUPING_DES,
229
- .init = grouping_init_des,
246
+ .init = grouping_init_des,
247
.create= grouping_create_des,
248
.reset = grouping_reset_des,
249
.free = grouping_free_des,
250
.add = grouping_add_des,
234
- .flush = grouping_flush_des
251
+ .flush = grouping_flush_des,
252
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
253
},
254
255
{.name = "countif",
260
.reset = grouping_reset_countif,
261
.free = grouping_free_countif,
262
.add = grouping_add_countif,
245
- .flush = grouping_flush_countif
263
+ .flush = grouping_flush_countif,
264
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
265
},
266
267
// terminator
273
.reset = grouping_reset_average,
274
.free = grouping_free_average,
275
.add = grouping_add_average,
257
- .flush = grouping_flush_average
276
+ .flush = grouping_flush_average,
277
+ .tier_query_fetch = TIER_QUERY_FETCH_AVERAGE
278
}
279
};
280
326
int i, found = 0;
327
for(i = 0; !found && api_v1_data_groups[i].name ;i++) {
328
if(api_v1_data_groups[i].value == group_method) {
309
- r->internal.grouping_create= api_v1_data_groups[i].create;
310
- r->internal.grouping_reset = api_v1_data_groups[i].reset;
311
- r->internal.grouping_free = api_v1_data_groups[i].free;
312
- r->internal.grouping_add = api_v1_data_groups[i].add;
313
- r->internal.grouping_flush = api_v1_data_groups[i].flush;
329
+ r->internal.grouping_create = api_v1_data_groups[i].create;
330
+ r->internal.grouping_reset = api_v1_data_groups[i].reset;
331
+ r->internal.grouping_free = api_v1_data_groups[i].free;
332
+ r->internal.grouping_add = api_v1_data_groups[i].add;
333
+ r->internal.grouping_flush = api_v1_data_groups[i].flush;
334
+ r->internal.tier_query_fetch = api_v1_data_groups[i].tier_query_fetch;
335
found = 1;
336
}
337
}
338
if(!found) {
339
errno = 0;
340
internal_error(true, "QUERY: grouping method %u not found. Using 'average'", (unsigned int)group_method);
320
- r->internal.grouping_create = grouping_create_average;
321
- r->internal.grouping_reset = grouping_reset_average;
322
- r->internal.grouping_free = grouping_free_average;
323
- r->internal.grouping_add = grouping_add_average;
324
- r->internal.grouping_flush = grouping_flush_average;
341
+ r->internal.grouping_create = grouping_create_average;
342
+ r->internal.grouping_reset = grouping_reset_average;
343
+ r->internal.grouping_free = grouping_free_average;
344
+ r->internal.grouping_add = grouping_add_average;
345
+ r->internal.grouping_flush = grouping_flush_average;
346
+ r->internal.tier_query_fetch = TIER_QUERY_FETCH_AVERAGE;
347
}
348
}
349
402
403
// check if all dimensions are hidden
404
if(unlikely(!dims_not_hidden_not_zero && dims_selected)) {
383
- // there are a few selected dimensions
405
+ // there are a few selected dimensions,
406
// but they are all zero
407
// enable the selected ones
408
// to avoid returning an empty chart
446
447
448
// ----------------------------------------------------------------------------
427
-// fill RRDR for a single dimension
449
+// tier management
450
+
451
+static int rrddim_find_best_tier_for_timeframe(RRDDIM *rd, time_t after_wanted, time_t before_wanted, long points_wanted) {
452
+ if(unlikely(storage_tiers < 2))
453
+ return 0;
454
+
455
+ if(unlikely(after_wanted == before_wanted || points_wanted <= 0 || !rd || !rd->rrdset)) {
456
+
457
+ if(!rd)
458
+ internal_error(true, "QUERY: NULL dimension - invalid params to tier calculation");
459
+ else
460
+ internal_error(true, "QUERY: chart '%s' dimension '%s' invalid params to tier calculation",
461
+ (rd->rrdset)?rd->rrdset->name:"unknown", rd->name);
462
+
463
+ return 0;
464
+ }
465
+
466
+ //BUFFER *wb = buffer_create(1000);
467
+ //buffer_sprintf(wb, "Best tier for chart '%s', dim '%s', from %ld to %ld (dur %ld, every %d), points %ld",
468
+ // rd->rrdset->name, rd->name, after_wanted, before_wanted, before_wanted - after_wanted, rd->update_every, points_wanted);
469
+
470
+ long weight[storage_tiers];
471
+
472
+ for(int tier = 0; tier < storage_tiers ; tier++) {
473
+ if(unlikely(!rd->tiers[tier])) {
474
+ internal_error(true, "QUERY: tier %d of chart '%s' dimension '%s' not initialized",
475
+ tier, rd->rrdset->name, rd->name);
476
+ // buffer_free(wb);
477
+ return 0;
478
+ }
479
+
480
+ time_t first_t = rd->tiers[tier]->query_ops.oldest_time(rd->tiers[tier]->db_metric_handle);
481
+ time_t last_t = rd->tiers[tier]->query_ops.latest_time(rd->tiers[tier]->db_metric_handle);
482
+
483
+ time_t common_after = MAX(first_t, after_wanted);
484
+ time_t common_before = MIN(last_t, before_wanted);
485
+
486
+ long time_coverage = (common_before - common_after) * 1000 / (before_wanted - after_wanted);
487
+ if(time_coverage < 0) time_coverage = 0;
488
+
489
+ int update_every = (int)rd->tiers[tier]->tier_grouping * (int)rd->update_every;
490
+ if(unlikely(update_every == 0)) {
491
+ internal_error(true, "QUERY: update_every of tier %d for chart '%s' dimension '%s' is zero. tg = %d, ue = %d",
492
+ tier, rd->rrdset->name, rd->name, rd->tiers[tier]->tier_grouping, rd->update_every);
493
+ // buffer_free(wb);
494
+ return 0;
495
+ }
496
+
497
+ long points_available = (before_wanted - after_wanted) / update_every;
498
+ long points_delta = points_available - points_wanted;
499
+ long points_coverage = (points_delta < 0) ? points_available * 1000 / points_wanted: 1000;
500
+
501
+ if(points_available <= 0)
502
+ weight[tier] = -LONG_MAX;
503
+ else
504
+ weight[tier] = points_coverage;
505
+
506
+ // buffer_sprintf(wb, ": tier %d, first %ld, last %ld (dur %ld, tg %d, every %d), points %ld, tcoverage %ld, pcoverage %ld, weight %ld",
507
+ // tier, first_t, last_t, last_t - first_t, rd->tiers[tier]->tier_grouping, update_every,
508
+ // points_available, time_coverage, points_coverage, weight[tier]);
509
+ }
510
+
511
+ int best_tier = 0;
512
+ for(int tier = 1; tier < storage_tiers ; tier++) {
513
+ if(weight[tier] >= weight[best_tier])
514
+ best_tier = tier;
515
+ }
516
429
-static inline NETDATA_DOUBLE interpolate_value(NETDATA_DOUBLE this_value, NETDATA_DOUBLE last_value, time_t last_value_end_t, time_t this_value_start_t, time_t now, time_t this_value_end_t) {
430
- if(unlikely(
431
- this_value_start_t + 1 == this_value_end_t ||
432
- !netdata_double_isnumber(this_value) ||
433
- !netdata_double_isnumber(last_value) ||
434
- last_value_end_t != this_value_start_t))
435
- return this_value;
517
+ if(weight[best_tier] == -LONG_MAX)
518
+ best_tier = 0;
519
437
- return last_value + (this_value - last_value) * ( 1.0 - (NETDATA_DOUBLE)(this_value_end_t - now) / (NETDATA_DOUBLE)(this_value_end_t - this_value_start_t) );
520
+ //buffer_sprintf(wb, ": final best tier %d", best_tier);
521
+ //internal_error(true, "%s", buffer_tostring(wb));
522
+ //buffer_free(wb);
523
+
524
+ return best_tier;
525
}
526
440
-static inline void rrd2rrdr_do_dimension(
441
- RRDR *r
442
- , long points_wanted
443
- , RRDDIM *rd
444
- , long dim_id_in_rrdr
445
- , time_t after_wanted
446
- , time_t before_wanted
447
- , RRDR_OPTIONS options
448
-){
449
- time_t now = after_wanted,
450
- query_granularity = r->update_every / r->group,
451
- max_date = 0,
452
- min_date = 0;
527
+static int rrdset_find_natural_update_every_for_timeframe(RRDSET *st, time_t after_wanted, time_t before_wanted, long points_wanted, RRDR_OPTIONS options, int tier) {
528
+ int ret = st->update_every;
529
+
530
+ if(unlikely(!st->dimensions))
531
+ return ret;
532
454
- bool interpolate = query_granularity < rd->update_every;
533
+ rrdset_rdlock(st);
534
+ int best_tier;
535
456
- long group_points_wanted = r->group,
457
- points_added = 0, group_points_added = 0, group_points_non_zero = 0,
458
- rrdr_line = -1;
536
+ if(options & RRDR_OPTION_SELECTED_TIER && tier >= 0 && tier < storage_tiers)
537
+ best_tier = tier;
538
+ else
539
+ best_tier = rrddim_find_best_tier_for_timeframe(st->dimensions, after_wanted, before_wanted, points_wanted);
540
+
541
+ if(!st->dimensions->tiers[best_tier]) {
542
+ internal_error(
543
+ true,
544
+ "QUERY: tier %d on chart '%s', is not initialized", best_tier, st->name);
545
+ }
546
+ else {
547
+ ret = (int)st->dimensions->tiers[best_tier]->tier_grouping * (int)st->update_every;
548
+ if(unlikely(!ret)) {
549
+ internal_error(
550
+ true,
551
+ "QUERY: update_every calculated to be zero on chart '%s', tier_grouping %d, update_every %d",
552
+ st->name, st->dimensions->tiers[best_tier]->tier_grouping, st->update_every);
553
+
554
+ ret = st->update_every;
555
+ }
556
+ }
557
+
558
+ rrdset_unlock(st);
559
460
- size_t group_anomaly_rate = 0;
560
+ return ret;
561
+}
562
+
563
+// ----------------------------------------------------------------------------
564
+// query ops
565
+
566
+typedef struct query_point {
567
+ NETDATA_DOUBLE value;
568
+ SN_FLAGS flags;
569
+ size_t anomaly;
570
+ time_t start_time;
571
+ time_t end_time;
572
+} QUERY_POINT;
573
+
574
+QUERY_POINT QUERY_POINT_EMPTY = {
575
+ .value = NAN,
576
+ .flags = SN_EMPTY_SLOT,
577
+ .anomaly = 0,
578
+ .start_time = 0,
579
+ .end_time = 0
580
+};
581
+
582
+typedef struct query_plan_entry {
583
+ size_t tier;
584
+ time_t after;
585
+ time_t before;
586
+} QUERY_PLAN_ENTRY;
587
462
- RRDR_VALUE_FLAGS group_value_flags = RRDR_VALUE_NOTHING;
588
+typedef struct query_plan {
589
+ size_t entries;
590
+ QUERY_PLAN_ENTRY data[RRD_STORAGE_TIERS*2];
591
+} QUERY_PLAN;
592
593
+typedef struct query_engine_ops {
594
+ // configuration
595
+ RRDR *r;
596
+ RRDDIM *rd;
597
+ time_t view_update_every;
598
+ time_t query_granularity;
599
+ TIER_QUERY_FETCH tier_query_fetch;
600
+
601
+ // query planer
602
+ QUERY_PLAN plan;
603
+ size_t current_plan;
604
+ time_t current_plan_expire_time;
605
+
606
+ // storage queries
607
+ size_t tier;
608
+ struct rrddim_tier *tier_ptr;
609
struct rrddim_query_handle handle;
610
+ STORAGE_POINT (*next_metric)(struct rrddim_query_handle *handle);
611
+ int (*is_finished)(struct rrddim_query_handle *handle);
612
+ void (*finalize)(struct rrddim_query_handle *handle);
613
466
- NETDATA_DOUBLE min = r->min, max = r->max;
467
- size_t db_points_read = 0;
468
-
469
- // cache the function pointers we need in the loop
470
- NETDATA_DOUBLE (*next_metric)(struct rrddim_query_handle *handle, time_t *current_time, time_t *end_time, SN_FLAGS *flags) = rd->state->query_ops.next_metric;
471
- void (*grouping_add)(struct rrdresult *r, NETDATA_DOUBLE value) = r->internal.grouping_add;
472
- NETDATA_DOUBLE (*grouping_flush)(struct rrdresult *r, RRDR_VALUE_FLAGS *rrdr_value_options_ptr) = r->internal.grouping_flush;
473
-
474
- NETDATA_DOUBLE last2_point_value;
475
- //SN_FLAGS last2_point_flags;
476
- //size_t last2_point_anomaly;
477
- //time_t last2_point_start_time;
478
- time_t last2_point_end_time;
479
-
480
- NETDATA_DOUBLE last1_point_value = NAN;
481
- SN_FLAGS last1_point_flags = SN_EMPTY_SLOT;
482
- size_t last1_point_anomaly = 0;
483
- time_t last1_point_start_time = 0;
484
- time_t last1_point_end_time = 0;
485
-
486
- NETDATA_DOUBLE new_point_value = NAN;
487
- SN_FLAGS new_point_flags = SN_EMPTY_SLOT;
488
- size_t new_point_anomaly = 0;
489
- time_t new_point_start_time = 0;
490
- time_t new_point_end_time = 0;
491
-
492
- for(rd->state->query_ops.init(rd, &handle, now, before_wanted) ; points_added < points_wanted ; now += query_granularity) {
493
-
494
- if(unlikely(now > before_wanted))
495
- break;
614
+ // aggregating points over time
615
+ void (*grouping_add)(struct rrdresult *r, NETDATA_DOUBLE value);
616
+ NETDATA_DOUBLE (*grouping_flush)(struct rrdresult *r, RRDR_VALUE_FLAGS *rrdr_value_options_ptr);
617
+ size_t group_points_non_zero;
618
+ size_t group_points_added;
619
+ size_t group_anomaly_rate;
620
+ RRDR_VALUE_FLAGS group_value_flags;
621
497
- last2_point_value = last1_point_value;
498
- //last2_point_flags = last1_point_flags;
499
- //last2_point_anomaly = last1_point_anomaly;
500
- //last2_point_start_time = last1_point_start_time;
501
- last2_point_end_time = last1_point_end_time;
622
+ // statistics
623
+ size_t db_total_points_read;
624
+ size_t db_points_read_per_tier[RRD_STORAGE_TIERS];
625
+} QUERY_ENGINE_OPS;
626
503
- last1_point_value = new_point_value;
504
- last1_point_flags = new_point_flags;
505
- last1_point_anomaly = new_point_anomaly;
506
- last1_point_start_time = new_point_start_time;
507
- last1_point_end_time = new_point_end_time;
627
509
- if(likely(!rd->state->query_ops.is_finished(&handle))) {
510
- // fetch the new point
511
- new_point_value = next_metric(&handle, &new_point_start_time, &new_point_end_time, &new_point_flags);
512
- db_points_read++;
513
-
514
- // dbengine does not take into account the starting time of points
515
- // and depending on the data collection frequency it may return
516
- // a point that is just before the wanted one.
517
- // So, here we fetch the next one.
518
- if(unlikely(new_point_end_time < now)) {
519
- internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu from %ld to %ld, before now (now = %ld, after_wanted = %ld, before_wanted = %ld, dt = %ld). Fetching the next one.",
520
- rd->rrdset->name, rd->name, db_points_read, new_point_start_time, new_point_end_time, now, after_wanted, before_wanted, query_granularity);
521
-
522
- new_point_value = next_metric(&handle, &new_point_start_time, &new_point_end_time, &new_point_flags);
523
- db_points_read++;
524
- }
628
+// ----------------------------------------------------------------------------
629
+// query planer
630
+
631
+#define query_plan_should_switch_plan(ops, now) ((now) >= (ops).current_plan_expire_time)
632
+
633
+static void query_planer_activate_plan(QUERY_ENGINE_OPS *ops, size_t plan_id, time_t overwrite_after) {
634
+ if(unlikely(plan_id >= ops->plan.entries))
635
+ plan_id = ops->plan.entries - 1;
636
+
637
+ time_t after = ops->plan.data[plan_id].after;
638
+ time_t before = ops->plan.data[plan_id].before;
639
+
640
+ if(overwrite_after > after && overwrite_after < before)
641
+ after = overwrite_after;
642
+
643
+ ops->tier = ops->plan.data[plan_id].tier;
644
+ ops->tier_ptr = ops->rd->tiers[ops->tier];
645
+ ops->tier_ptr->query_ops.init(ops->tier_ptr->db_metric_handle, &ops->handle, after, before, ops->r->internal.tier_query_fetch);
646
+ ops->next_metric = ops->tier_ptr->query_ops.next_metric;
647
+ ops->is_finished = ops->tier_ptr->query_ops.is_finished;
648
+ ops->finalize = ops->tier_ptr->query_ops.finalize;
649
+ ops->current_plan = plan_id;
650
+ ops->current_plan_expire_time = ops->plan.data[plan_id].before;
651
+}
652
+
653
+static void query_planer_next_plan(QUERY_ENGINE_OPS *ops, time_t now, time_t last_point_end_time) {
654
+ internal_error(now < ops->current_plan_expire_time && now < ops->plan.data[ops->current_plan].before,
655
+ "QUERY: switching query plan too early!");
656
+
657
+ time_t next_plan_before_time;
658
+ do {
659
+ ops->current_plan++;
660
+
661
+ if (ops->current_plan >= ops->plan.entries) {
662
+ ops->current_plan = ops->plan.entries - 1;
663
+ return;
664
+ }
665
+
666
+ next_plan_before_time = ops->plan.data[ops->current_plan].before;
667
+ } while(now >= next_plan_before_time || last_point_end_time >= next_plan_before_time);
668
+
669
+ if(ops->finalize) {
670
+ ops->finalize(&ops->handle);
671
+ ops->finalize = NULL;
672
+ }
673
+
674
+ query_planer_activate_plan(ops, ops->current_plan, MIN(now, last_point_end_time));
675
+
676
+ // internal_error(true, "QUERY: switched plan to %zu (all is %zu), previous expiration was %ld, this starts at %ld, now is %ld, last_point_end_time %ld", ops->current_plan, ops->plan.entries, ops->plan.data[ops->current_plan-1].before, ops->plan.data[ops->current_plan].after, now, last_point_end_time);
677
+}
678
+
679
+static int compare_query_plan_entries_on_start_time(const void *a, const void *b) {
680
+ QUERY_PLAN_ENTRY *p1 = (QUERY_PLAN_ENTRY *)a;
681
+ QUERY_PLAN_ENTRY *p2 = (QUERY_PLAN_ENTRY *)b;
682
+ return (p1->after < p2->after)?-1:1;
683
+}
684
+
685
+static void query_plan(QUERY_ENGINE_OPS *ops, time_t after_wanted, time_t before_wanted, long points_wanted) {
686
+ RRDDIM *rd = ops->rd;
687
526
- if(likely(netdata_double_isnumber(new_point_value))) {
527
- new_point_anomaly = (new_point_flags & SN_ANOMALY_BIT) ? 0 : 100;
688
+ //BUFFER *wb = buffer_create(1000);
689
+ //buffer_sprintf(wb, "QUERY PLAN for chart '%s' dimension '%s', from %ld to %ld:", rd->rrdset->name, rd->name, after_wanted, before_wanted);
690
529
- if(unlikely(options & RRDR_OPTION_ANOMALY_BIT))
530
- new_point_value = (NETDATA_DOUBLE)new_point_anomaly;
691
+ // put our selected tier as the first plan
692
+ size_t selected_tier;
693
+
694
+ if(ops->r->internal.query_options & RRDR_OPTION_SELECTED_TIER && ops->r->internal.query_tier >= 0 && ops->r->internal.query_tier < storage_tiers) {
695
+ selected_tier = ops->r->internal.query_tier;
696
+ }
697
+ else {
698
+
699
+ selected_tier = rrddim_find_best_tier_for_timeframe(rd, after_wanted, before_wanted, points_wanted);
700
+
701
+ if(ops->r->internal.query_options & RRDR_OPTION_SELECTED_TIER)
702
+ ops->r->internal.query_options &= ~RRDR_OPTION_SELECTED_TIER;
703
+ }
704
+
705
+ ops->plan.entries = 1;
706
+ ops->plan.data[0].tier = selected_tier;
707
+ ops->plan.data[0].after = rd->tiers[selected_tier]->query_ops.oldest_time(rd->tiers[selected_tier]->db_metric_handle);
708
+ ops->plan.data[0].before = rd->tiers[selected_tier]->query_ops.latest_time(rd->tiers[selected_tier]->db_metric_handle);
709
+
710
+ if(!(ops->r->internal.query_options & RRDR_OPTION_SELECTED_TIER)) {
711
+ // the selected tier
712
+ time_t selected_tier_first_time_t = ops->plan.data[0].after;
713
+ time_t selected_tier_last_time_t = ops->plan.data[0].before;
714
+
715
+ //buffer_sprintf(wb, ": SELECTED tier %zu, from %ld to %ld", selected_tier, ops->plan.data[0].after, ops->plan.data[0].before);
716
+
717
+ // check if our selected tier can start the query
718
+ if (selected_tier_first_time_t > after_wanted) {
719
+ // we need some help from other tiers
720
+ for (int tr = (int)selected_tier + 1; tr < storage_tiers; tr++) {
721
+ // find the first time of this tier
722
+ time_t first_time_t = rd->tiers[tr]->query_ops.oldest_time(rd->tiers[tr]->db_metric_handle);
723
+
724
+ //buffer_sprintf(wb, ": EVAL AFTER tier %d, %ld", tier, first_time_t);
725
+
726
+ // can it help?
727
+ if (first_time_t < selected_tier_first_time_t) {
728
+ // it can help us add detail at the beginning of the query
729
+ QUERY_PLAN_ENTRY t = {
730
+ .tier = tr,
731
+ .after = (first_time_t < after_wanted) ? after_wanted : first_time_t,
732
+ .before = selected_tier_first_time_t};
733
+ ops->plan.data[ops->plan.entries++] = t;
734
+
735
+ // prepare for the tier
736
+ selected_tier_first_time_t = t.after;
737
+
738
+ if (t.after <= after_wanted)
739
+ break;
740
+ }
741
}
532
- else {
533
- new_point_flags = SN_EMPTY_SLOT;
534
- new_point_value = NAN;
535
- new_point_anomaly = 0;
742
+ }
743
+
744
+ // check if our selected tier can finish the query
745
+ if (selected_tier_last_time_t < before_wanted) {
746
+ // we need some help from other tiers
747
+ for (int tr = (int)selected_tier - 1; tr >= 0; tr--) {
748
+ // find the last time of this tier
749
+ time_t last_time_t = rd->tiers[tr]->query_ops.latest_time(rd->tiers[tr]->db_metric_handle);
750
+
751
+ //buffer_sprintf(wb, ": EVAL BEFORE tier %d, %ld", tier, last_time_t);
752
+
753
+ // can it help?
754
+ if (last_time_t > selected_tier_last_time_t) {
755
+ // it can help us add detail at the end of the query
756
+ QUERY_PLAN_ENTRY t = {
757
+ .tier = tr,
758
+ .after = selected_tier_last_time_t,
759
+ .before = (last_time_t > before_wanted) ? before_wanted : last_time_t};
760
+ ops->plan.data[ops->plan.entries++] = t;
761
+
762
+ // prepare for the tier
763
+ selected_tier_last_time_t = t.before;
764
+
765
+ if (t.before >= before_wanted)
766
+ break;
767
+ }
768
}
769
+ }
770
+ }
771
538
- if(unlikely(new_point_start_time == new_point_end_time)) {
539
- internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu start time %ld, end time %ld, that are both equal",
540
- rd->rrdset->name, rd->name, db_points_read, new_point_start_time, new_point_end_time);
772
+ // sort the query plan
773
+ if(ops->plan.entries > 1)
774
+ qsort(&ops->plan.data, ops->plan.entries, sizeof(QUERY_PLAN_ENTRY), compare_query_plan_entries_on_start_time);
775
+
776
+ // make sure it has the whole timeframe we need
777
+ ops->plan.data[0].after = after_wanted;
778
+ ops->plan.data[ops->plan.entries - 1].before = before_wanted;
779
+
780
+ //buffer_sprintf(wb, ": FINAL STEPS %zu", ops->plan.entries);
781
+
782
+ //for(size_t i = 0; i < ops->plan.entries ;i++)
783
+ // buffer_sprintf(wb, ": STEP %zu = use tier %zu from %ld to %ld", i+1, ops->plan.data[i].tier, ops->plan.data[i].after, ops->plan.data[i].before);
784
+
785
+ //internal_error(true, "%s", buffer_tostring(wb));
786
+
787
+ query_planer_activate_plan(ops, 0, 0);
788
+}
789
+
790
+
791
+// ----------------------------------------------------------------------------
792
+// dimension level query engine
793
+
794
+#define query_interpolate_point(this_point, last_point, now) do { \
795
+ if(likely( \
796
+ /* the point to interpolate is more than 1s wide */ \
797
+ (this_point).end_time - (this_point).start_time > 1 \
798
+ \
799
+ /* the two points are exactly next to each other */ \
800
+ && (last_point).end_time == (this_point).start_time \
801
+ \
802
+ /* both points are valid numbers */ \
803
+ && netdata_double_isnumber((this_point).value) \
804
+ && netdata_double_isnumber((last_point).value) \
805
+ \
806
+ )) { \
807
+ (this_point).value = (last_point).value + ((this_point).value - (last_point).value) * (1.0 - (NETDATA_DOUBLE)((this_point).end_time - (now)) / (NETDATA_DOUBLE)((this_point).end_time - (this_point).start_time)); \
808
+ (this_point).end_time = now; \
809
+ } \
810
+} while(0)
811
+
812
+#define query_add_point_to_group(r, point, ops) do { \
813
+ if(likely(netdata_double_isnumber((point).value))) { \
814
+ if(likely((point).value != 0.0)) \
815
+ (ops).group_points_non_zero++; \
816
+ \
817
+ if(unlikely((point).flags & SN_EXISTS_RESET)) \
818
+ (ops).group_value_flags |= RRDR_VALUE_RESET; \
819
+ \
820
+ (ops).grouping_add(r, (point).value); \
821
+ } \
822
+ \
823
+ (ops).group_points_added++; \
824
+ (ops).group_anomaly_rate += (point).anomaly; \
825
+} while(0)
826
542
- new_point_start_time = new_point_end_time - rd->update_every;
827
+static inline void rrd2rrdr_do_dimension(
828
+ RRDR *r
829
+ , long points_wanted
830
+ , RRDDIM *rd
831
+ , long dim_id_in_rrdr
832
+ , time_t after_wanted
833
+ , time_t before_wanted
834
+){
835
+ time_t max_date = 0,
836
+ min_date = 0;
837
+
838
+ size_t points_added = 0;
839
+
840
+ QUERY_ENGINE_OPS ops = {
841
+ .r = r,
842
+ .rd = rd,
843
+ .grouping_add = r->internal.grouping_add,
844
+ .grouping_flush = r->internal.grouping_flush,
845
+ .tier_query_fetch = r->internal.tier_query_fetch,
846
+ .view_update_every = r->update_every,
847
+ .query_granularity = r->update_every / r->group,
848
+ .group_value_flags = RRDR_VALUE_NOTHING
849
+ };
850
+
851
+ long rrdr_line = -1;
852
+ bool use_anomaly_bit_as_value = (r->internal.query_options & RRDR_OPTION_ANOMALY_BIT) ? true : false;
853
+
854
+ query_plan(&ops, after_wanted, before_wanted, points_wanted);
855
+
856
+ NETDATA_DOUBLE min = r->min, max = r->max;
857
+
858
+ QUERY_POINT last2_point = QUERY_POINT_EMPTY;
859
+ QUERY_POINT last1_point = QUERY_POINT_EMPTY;
860
+ QUERY_POINT new_point = QUERY_POINT_EMPTY;
861
+
862
+ // The main loop, based on the query granularity we need
863
+ for(time_t now = after_wanted + ops.view_update_every - ops.query_granularity; (long)points_added < points_wanted ; now += ops.view_update_every) {
864
+
865
+ if(query_plan_should_switch_plan(ops, now))
866
+ query_planer_next_plan(&ops, now, new_point.end_time);
867
+
868
+ // real all the points of the db, prior to the time we need (now)
869
+ size_t count_same_end_time = 0;
870
+ while(count_same_end_time < 100) {
871
+ if(likely(count_same_end_time == 0)) {
872
+ last2_point = last1_point;
873
+ last1_point = new_point;
874
+ }
875
+
876
+ if(unlikely(ops.is_finished(&ops.handle))) {
877
+ if(count_same_end_time != 0) {
878
+ last2_point = last1_point;
879
+ last1_point = new_point;
880
+ }
881
+ new_point = QUERY_POINT_EMPTY;
882
+ new_point.start_time = last1_point.end_time;
883
+ new_point.end_time = now;
884
+ break;
885
}
886
545
- if(unlikely(new_point_start_time < last1_point_start_time && new_point_end_time < last1_point_end_time)) {
546
- internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu start time %ld, end time %ld, before the last point start time %ld, end time %ld",
547
- rd->rrdset->name, rd->name, db_points_read, new_point_start_time, new_point_end_time,
548
- last1_point_start_time,
549
- last1_point_end_time);
887
+ // fetch the new point
888
+ {
889
+ STORAGE_POINT sp = ops.next_metric(&ops.handle);
890
+
891
+ ops.db_points_read_per_tier[ops.tier]++;
892
+ ops.db_total_points_read++;
893
+
894
+ new_point.start_time = sp.start_time;
895
+ new_point.end_time = sp.end_time;
896
+ new_point.anomaly = sp.count ? sp.anomaly_count * 100 / sp.count : 0;
897
+
898
+ if(likely(!storage_point_is_unset(sp) && !storage_point_is_empty(sp))) {
899
+
900
+ if(unlikely(use_anomaly_bit_as_value))
901
+ new_point.value = (NETDATA_DOUBLE)new_point.anomaly;
902
+
903
+ else {
904
+ switch (ops.tier_query_fetch) {
905
+ default:
906
+ case TIER_QUERY_FETCH_AVERAGE:
907
+ new_point.value = sp.sum / sp.count;
908
+ break;
909
+
910
+ case TIER_QUERY_FETCH_MIN:
911
+ new_point.value = sp.min;
912
+ break;
913
+
914
+ case TIER_QUERY_FETCH_MAX:
915
+ new_point.value = sp.max;
916
+ break;
917
551
- new_point_value = last1_point_value;
552
- new_point_flags = last1_point_flags;
553
- new_point_start_time = last1_point_start_time;
554
- new_point_end_time = last1_point_end_time;
918
+ case TIER_QUERY_FETCH_SUM:
919
+ new_point.value = sp.sum;
920
+ break;
921
+ };
922
+ }
923
+ }
924
+ else {
925
+ new_point.value = NAN;
926
+ new_point.flags = SN_EMPTY_SLOT;
927
+ }
928
}
929
557
- if(unlikely(new_point_end_time < last1_point_end_time)) {
558
- internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu end time %ld, before the last point end time %ld",
559
- rd->rrdset->name, rd->name, db_points_read, new_point_end_time,
560
- last1_point_end_time);
930
+ if(unlikely(new_point.start_time == new_point.end_time)) {
931
+ internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu start time %ld, end time %ld, that are both equal",
932
+ rd->rrdset->name, rd->name, ops.db_total_points_read, new_point.start_time, new_point.end_time);
933
562
- new_point_value = last1_point_value;
563
- new_point_flags = last1_point_flags;
564
- new_point_start_time = last1_point_start_time;
565
- new_point_end_time = last1_point_end_time;
934
+ new_point.start_time = new_point.end_time - ((time_t)ops.tier_ptr->tier_grouping * (time_t)ops.rd->update_every);
935
}
936
568
- if(unlikely(new_point_end_time < now)) {
569
- internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu from %ld to %ld, before now (now = %ld, after_wanted = %ld, before_wanted = %ld, dt = %ld)",
570
- rd->rrdset->name, rd->name, db_points_read, new_point_start_time, new_point_end_time,
571
- now, after_wanted, before_wanted, query_granularity);
937
+ if(unlikely(new_point.end_time <= last1_point.end_time)) {
938
+ internal_error(true, "QUERY: next_metric(%s, %s) returned point %zu from %ld time %ld, before the last point end time %ld, now is %ld",
939
+ rd->rrdset->name, rd->name, ops.db_total_points_read, new_point.start_time, new_point.end_time, last1_point.end_time, now);
940
573
- new_point_end_time = now;
941
+ count_same_end_time++;
942
+ continue;
943
}
944
+
945
+ count_same_end_time = 0;
946
+
947
+ if(new_point.end_time < now)
948
+ query_add_point_to_group(r, new_point, ops);
949
+ else
950
+ break;
951
}
576
- else {
577
- new_point_value = NAN;
578
- new_point_flags = SN_EMPTY_SLOT;
579
- new_point_start_time = last1_point_end_time;
580
- new_point_end_time = now;
952
+
953
+ if(count_same_end_time) {
954
+ internal_error(true,
955
+ "QUERY: the database does not advance the query, it returned an end time less or equal to %ld, %zu times",
956
+ last1_point.end_time, count_same_end_time);
957
+
958
+ new_point.end_time = now;
959
}
960
961
// the inner loop
584
- // we have 3 points in memory: last, new, next
962
+ // we have 3 points in memory: last2, last1, new
963
// we select the one to use based on their timestamps
964
965
size_t iterations = 0;
588
- for ( ; now <= new_point_end_time && points_added < points_wanted; now += query_granularity, iterations++) {
966
+ for ( ; now <= new_point.end_time && (long)points_added < points_wanted; now += ops.view_update_every, iterations++) {
967
+ QUERY_POINT current_point;
968
590
- NETDATA_DOUBLE current_point_value;
591
- SN_FLAGS current_point_flags;
592
- size_t current_point_anomaly;
593
- //time_t current_point_start_time;
594
- //time_t current_point_end_time;
595
-
596
- if(likely(now > new_point_start_time)) {
969
+ if(likely(now > new_point.start_time)) {
970
// it is time for our NEW point to be used
598
- current_point_value = interpolate ? interpolate_value(new_point_value, last1_point_value, last1_point_end_time, new_point_start_time, now, new_point_end_time) : new_point_value;
599
- current_point_flags = new_point_flags;
600
- current_point_anomaly = new_point_anomaly;
601
- //current_point_start_time = new_point_start_time;
602
- //current_point_end_time = new_point_end_time;
971
+ current_point = new_point;
972
+ query_interpolate_point(current_point, last1_point, now);
973
}
604
- else if(likely(now <= last1_point_end_time)) {
974
+ else if(likely(now <= last1_point.end_time)) {
975
// our LAST point is still valid
606
- current_point_value = interpolate ? interpolate_value(last1_point_value, last2_point_value, last2_point_end_time, last1_point_start_time, now, last1_point_end_time) : last1_point_value;
607
- current_point_flags = last1_point_flags;
608
- current_point_anomaly = last1_point_anomaly;
609
- //current_point_start_time = last_point_start_time;
610
- //current_point_end_time = last_point_end_time;
976
+ current_point = last1_point;
977
+ query_interpolate_point(current_point, last2_point, now);
978
}
979
else {
980
// a GAP, we don't have a value this time
614
- current_point_value = NAN;
615
- current_point_flags = SN_EMPTY_SLOT;
616
- current_point_anomaly = 0;
617
- //current_point_start_time = now - dt;
618
- //current_point_end_time = now;
981
+ current_point = QUERY_POINT_EMPTY;
982
}
983
621
- if(likely(netdata_double_isnumber(current_point_value))) {
622
- if(likely(current_point_value != 0.0))
623
- group_points_non_zero++;
624
-
625
- if(unlikely(current_point_flags & SN_EXISTS_RESET))
626
- group_value_flags |= RRDR_VALUE_RESET;
984
+ query_add_point_to_group(r, current_point, ops);
985
628
- grouping_add(r, current_point_value);
629
- }
986
+ rrdr_line = rrdr_line_init(r, now, rrdr_line);
987
+ size_t rrdr_o_v_index = rrdr_line * r->d + dim_id_in_rrdr;
988
631
- // add this value for grouping
632
- group_points_added++;
633
- group_anomaly_rate += current_point_anomaly;
989
+ if(unlikely(!min_date)) min_date = now;
990
+ max_date = now;
991
635
- if(unlikely(group_points_added == group_points_wanted)) {
636
- rrdr_line = rrdr_line_init(r, now, rrdr_line);
637
- size_t rrdr_o_v_index = rrdr_line * r->d + dim_id_in_rrdr;
992
+ // find the place to store our values
993
+ RRDR_VALUE_FLAGS *rrdr_value_options_ptr = &r->o[rrdr_o_v_index];
994
639
- if(unlikely(!min_date)) min_date = now;
640
- max_date = now;
995
+ // update the dimension options
996
+ if(likely(ops.group_points_non_zero))
997
+ r->od[dim_id_in_rrdr] |= RRDR_DIMENSION_NONZERO;
998
642
- // find the place to store our values
643
- RRDR_VALUE_FLAGS *rrdr_value_options_ptr = &r->o[rrdr_o_v_index];
999
+ // store the specific point options
1000
+ *rrdr_value_options_ptr = ops.group_value_flags;
1001
645
- // update the dimension options
646
- if(likely(group_points_non_zero))
647
- r->od[dim_id_in_rrdr] |= RRDR_DIMENSION_NONZERO;
1002
+ // store the group value
1003
+ NETDATA_DOUBLE group_value = ops.grouping_flush(r, rrdr_value_options_ptr);
1004
+ r->v[rrdr_o_v_index] = group_value;
1005
649
- // store the specific point options
650
- *rrdr_value_options_ptr = group_value_flags;
1006
+ // we only store uint8_t anomaly rates,
1007
+ // so let's get double precision by storing
1008
+ // anomaly rates in the range 0 - 200
1009
+ ops.group_anomaly_rate = (ops.group_anomaly_rate << 1) / ops.group_points_added;
1010
+ r->ar[rrdr_o_v_index] = (uint8_t)ops.group_anomaly_rate;
1011
652
- // store the group value
653
- NETDATA_DOUBLE group_value = grouping_flush(r, rrdr_value_options_ptr);
654
- r->v[rrdr_o_v_index] = group_value;
1012
+ if(likely(points_added || dim_id_in_rrdr)) {
1013
+ // find the min/max across all dimensions
1014
656
- // we only store uint8_t anomaly rates,
657
- // so let's get double precision by storing
658
- // anomaly rates in the range 0 - 200
659
- group_anomaly_rate = (group_anomaly_rate << 1) / group_points_added;
660
- r->ar[rrdr_o_v_index] = (uint8_t)group_anomaly_rate;
1015
+ if(unlikely(group_value < min)) min = group_value;
1016
+ if(unlikely(group_value > max)) max = group_value;
1017
662
- if(likely(points_added || dim_id_in_rrdr)) {
663
- // find the min/max across all dimensions
664
-
665
- if(unlikely(group_value < min)) min = group_value;
666
- if(unlikely(group_value > max)) max = group_value;
667
-
668
- }
669
- else {
670
- // runs only when dim_id_in_rrdr == 0 && points_added == 0
671
- // so, on the first point added for the query.
672
- min = max = group_value;
673
- }
674
-
675
- points_added++;
676
- group_points_added = 0;
677
- group_value_flags = RRDR_VALUE_NOTHING;
678
- group_points_non_zero = 0;
679
- group_anomaly_rate = 0;
1018
}
1019
+ else {
1020
+ // runs only when dim_id_in_rrdr == 0 && points_added == 0
1021
+ // so, on the first point added for the query.
1022
+ min = max = group_value;
1023
+ }
1024
+
1025
+ points_added++;
1026
+ ops.group_points_added = 0;
1027
+ ops.group_value_flags = RRDR_VALUE_NOTHING;
1028
+ ops.group_points_non_zero = 0;
1029
+ ops.group_anomaly_rate = 0;
1030
}
682
- // the loop above increased "now" by dt,
683
- // but the main loop will increase it,
1031
+ // the loop above increased "now" by query_granularity,
1032
+ // but the main loop will increase it too,
1033
// so, let's undo the last iteration of this loop
1034
if(iterations)
686
- now -= query_granularity;
1035
+ now -= ops.view_update_every;
1036
}
688
- rd->state->query_ops.finalize(&handle);
1037
+ ops.finalize(&ops.handle);
1038
690
- r->internal.db_points_read += db_points_read;
1039
r->internal.result_points_generated += points_added;
1040
+ r->internal.db_points_read += ops.db_total_points_read;
1041
+ for(int tr = 0; tr < storage_tiers ; tr++)
1042
+ r->internal.tier_points_read[tr] += ops.db_points_read_per_tier[tr];
1043
1044
r->min = min;
1045
r->max = max;
1046
r->before = max_date;
696
- r->after = min_date - (r->group - 1) * query_granularity;
1047
+ r->after = min_date - ops.view_update_every + ops.query_granularity;
1048
rrdr_done(r, rrdr_line);
1049
699
- internal_error(points_wanted != points_added,
1050
+ internal_error((long)points_added != points_wanted,
1051
"QUERY: query on %s/%s requested %zu points, but RRDR added %zu (%zu db points read).",
701
- r->st->name, rd->name, (size_t)points_wanted, (size_t)points_added, db_points_read);
1052
+ r->st->name, rd->name, (size_t)points_wanted, (size_t)points_added, ops.db_total_points_read);
1053
+}
1054
+
1055
+// ----------------------------------------------------------------------------
1056
+// fill the gap of a tier
1057
+
1058
+extern void store_metric_at_tier(RRDDIM *rd, struct rrddim_tier *t, STORAGE_POINT sp, usec_t now_ut);
1059
+
1060
+void rrdr_fill_tier_gap_from_smaller_tiers(RRDDIM *rd, int tier, time_t now) {
1061
+ if(unlikely(tier < 0 || tier >= storage_tiers)) return;
1062
+ if(storage_tiers_backfill[tier] == RRD_BACKFILL_NONE) return;
1063
+
1064
+ struct rrddim_tier *t = rd->tiers[tier];
1065
+ if(unlikely(!t)) return;
1066
+
1067
+ time_t latest_time_t = t->query_ops.latest_time(t->db_metric_handle);
1068
+ time_t granularity = (time_t)t->tier_grouping * (time_t)rd->update_every;
1069
+ time_t time_diff = now - latest_time_t;
1070
+
1071
+ // if the user wants only NEW backfilling, and we don't have any data
1072
+ if(storage_tiers_backfill[tier] == RRD_BACKFILL_NEW && latest_time_t <= 0) return;
1073
+
1074
+ // there is really nothing we can do
1075
+ if(now <= latest_time_t || time_diff < granularity) return;
1076
+
1077
+ struct rrddim_query_handle handle;
1078
+
1079
+ size_t all_points_read = 0;
1080
+
1081
+ // for each lower tier
1082
+ for(int tr = tier - 1; tr >= 0 ;tr--){
1083
+ time_t smaller_tier_first_time = rd->tiers[tr]->query_ops.oldest_time(rd->tiers[tr]->db_metric_handle);
1084
+ time_t smaller_tier_last_time = rd->tiers[tr]->query_ops.latest_time(rd->tiers[tr]->db_metric_handle);
1085
+ if(smaller_tier_last_time <= latest_time_t) continue; // it is as bad as we are
1086
+
1087
+ long after_wanted = (latest_time_t < smaller_tier_first_time) ? smaller_tier_first_time : latest_time_t;
1088
+ long before_wanted = smaller_tier_last_time;
1089
+
1090
+ struct rrddim_tier *tmp = rd->tiers[tr];
1091
+ tmp->query_ops.init(tmp->db_metric_handle, &handle, after_wanted, before_wanted, TIER_QUERY_FETCH_AVERAGE);
1092
+
1093
+ size_t points = 0;
1094
+
1095
+ while(!tmp->query_ops.is_finished(&handle)) {
1096
+
1097
+ STORAGE_POINT sp = tmp->query_ops.next_metric(&handle);
1098
+
1099
+ if(sp.end_time > latest_time_t) {
1100
+ latest_time_t = sp.end_time;
1101
+ store_metric_at_tier(rd, t, sp, sp.end_time * USEC_PER_SEC);
1102
+ points++;
1103
+ }
1104
+ }
1105
+
1106
+ all_points_read += points;
1107
+ tmp->query_ops.finalize(&handle);
1108
+
1109
+ internal_error(true, "DBENGINE: backfilled chart '%s', dimension '%s', tier %d, from %ld to %ld, with %zu points from tier %d",
1110
+ rd->rrdset->name, rd->name, tier, after_wanted, before_wanted, points, tr);
1111
+ }
1112
+
1113
+ rrdr_query_completed(all_points_read, all_points_read);
1114
}
1115
1116
// ----------------------------------------------------------------------------
1118
1119
#ifdef NETDATA_INTERNAL_CHECKS
1120
static void rrd2rrdr_log_request_response_metadata(RRDR *r
1121
+ , RRDR_OPTIONS options __maybe_unused
1122
, RRDR_GROUPING group_method
1123
, bool aligned
1124
, long group
1136
) {
1137
netdata_rwlock_rdlock(&r->st->rrdset_rwlock);
1138
info("INTERNAL ERROR: rrd2rrdr() on %s update every %d with %s grouping %s (group: %ld, resampling_time: %ld, resampling_group: %ld), "
726
- "after (got: %zu, want: %zu, req: %zu, db: %zu), "
727
- "before (got: %zu, want: %zu, req: %zu, db: %zu), "
728
- "duration (got: %zu, want: %zu, req: %zu, db: %zu), "
1139
+ "after (got: %zu, want: %zu, req: %ld, db: %zu), "
1140
+ "before (got: %zu, want: %zu, req: %ld, db: %zu), "
1141
+ "duration (got: %zu, want: %zu, req: %ld, db: %zu), "
1142
//"slot (after: %zu, before: %zu, delta: %zu), "
1143
"points (got: %ld, want: %ld, req: %ld, db: %ld), "
1144
"%s"
1155
// after
1156
, (size_t)r->after
1157
, (size_t)after_wanted
745
- , (size_t)after_requested
1158
+ , after_requested
1159
, (size_t)rrdset_first_entry_t_nolock(r->st)
1160
1161
// before
1162
, (size_t)r->before
1163
, (size_t)before_wanted
751
- , (size_t)before_requested
1164
+ , before_requested
1165
, (size_t)rrdset_last_entry_t_nolock(r->st)
1166
1167
// duration
1168
, (size_t)(r->before - r->after + r->st->update_every)
1169
, (size_t)(before_wanted - after_wanted + r->st->update_every)
757
- , (size_t)(before_requested - after_requested)
1170
+ , before_requested - after_requested
1171
, (size_t)((rrdset_last_entry_t_nolock(r->st) - rrdset_first_entry_t_nolock(r->st)) + r->st->update_every)
1172
1173
// slot
1191
#endif // NETDATA_INTERNAL_CHECKS
1192
1193
// Returns 1 if an absolute period was requested or 0 if it was a relative period
781
-int rrdr_relative_window_to_absolute(long long *after, long long *before, int update_every, long points) {
1194
+int rrdr_relative_window_to_absolute(long long *after, long long *before) {
1195
time_t now = now_realtime_sec() - 1;
1196
1197
int absolute_period_requested = -1;
1219
// if the user didn't give an after, use the number of points
1220
// to give a sane default
1221
if(after_requested == 0)
809
- after_requested = -(points * update_every);
1222
+ after_requested = -600;
1223
1224
// since the query engine now returns inclusive timestamps
1225
// it is awkward to return 6 points when after=-5 is given
1264
buffer_free(debug_log); \
1265
debug_log = NULL; \
1266
}
1267
+#define query_debug_log_free() do { buffer_free(debug_log); } while(0)
1268
#else
1269
#define query_debug_log_init() debug_dummy()
1270
#define query_debug_log(args...) debug_dummy()
1271
#define query_debug_log_fin() debug_dummy()
1272
+#define query_debug_log_free() debug_dummy()
1273
#endif
1274
1275
RRDR *rrd2rrdr(
1285
, struct context_param *context_param_list
1286
, const char *group_options
1287
, int timeout
1288
+ , int tier
1289
) {
1290
// RULES
1291
// points_requested = 0
1326
query_debug_log(":relative+natural");
1327
}
1328
913
- // this is the update_every of the query
914
- // it may be different to the update_every of the database
915
- time_t query_granularity = (natural_points)?update_every:1;
916
- query_debug_log(":query_granularity %ld", query_granularity);
1329
+ // if the user wants virtual points, make sure we do it
1330
+ if(options & RRDR_OPTION_VIRTUAL_POINTS)
1331
+ natural_points = false;
1332
+
1333
+ // set the right flag about natural and virtual points
1334
+ if(natural_points) {
1335
+ options |= RRDR_OPTION_NATURAL_POINTS;
1336
+
1337
+ if(options & RRDR_OPTION_VIRTUAL_POINTS)
1338
+ options &= ~RRDR_OPTION_VIRTUAL_POINTS;
1339
+ }
1340
+ else {
1341
+ options |= RRDR_OPTION_VIRTUAL_POINTS;
1342
+
1343
+ if(options & RRDR_OPTION_NATURAL_POINTS)
1344
+ options &= ~RRDR_OPTION_NATURAL_POINTS;
1345
+ }
1346
1347
if(after_wanted == 0 || before_wanted == 0) {
1348
// for non-context queries we have to find the duration of the database
1356
time_t last_entry_t = rrdset_last_entry_t_nolock(st);
1357
rrdset_unlock(st);
1358
1359
+ if(first_entry_t == 0 || last_entry_t == 0) {
1360
+ internal_error(true, "QUERY: chart without data detected on '%s'", st->name);
1361
+ query_debug_log_free();
1362
+ return NULL;
1363
+ }
1364
+
1365
query_debug_log(":first_entry_t %ld, last_entry_t %ld", first_entry_t, last_entry_t);
1366
1367
if (after_wanted == 0) {
1387
after_wanted = -600;
1388
query_debug_log(":zero600 after_wanted %lld", after_wanted);
1389
}
1390
+ }
1391
956
- if(points_wanted == 0) {
957
- points_wanted = 600;
958
- query_debug_log(":zero600 points_wanted %ld", points_wanted);
959
- }
1392
+ if(points_wanted == 0) {
1393
+ points_wanted = 600;
1394
+ query_debug_log(":zero600 points_wanted %ld", points_wanted);
1395
}
1396
1397
// convert our before_wanted and after_wanted to absolute
963
- rrdr_relative_window_to_absolute(&after_wanted, &before_wanted, (int)query_granularity, points_wanted);
1398
+ rrdr_relative_window_to_absolute(&after_wanted, &before_wanted);
1399
query_debug_log(":relative2absolute after %lld, before %lld", after_wanted, before_wanted);
1400
1401
+ if(natural_points && (options & RRDR_OPTION_SELECTED_TIER) && tier > 0 && storage_tiers > 1) {
1402
+ update_every = rrdset_find_natural_update_every_for_timeframe(st, after_wanted, before_wanted, points_wanted, options, tier);
1403
+ if(update_every <= 0) update_every = st->update_every;
1404
+ query_debug_log(":natural update every %d", update_every);
1405
+ }
1406
+
1407
+ // this is the update_every of the query
1408
+ // it may be different to the update_every of the database
1409
+ time_t query_granularity = (natural_points)?update_every:1;
1410
+ if(query_granularity <= 0) query_granularity = 1;
1411
+ query_debug_log(":query_granularity %ld", query_granularity);
1412
+
1413
// align before_wanted and after_wanted to query_granularity
1414
if (before_wanted % query_granularity) {
1415
before_wanted -= before_wanted % query_granularity;
1424
// automatic_natural_points is set when the user wants all the points available in the database
1425
if(automatic_natural_points) {
1426
points_wanted = (before_wanted - after_wanted + 1) / query_granularity;
1427
+ if(unlikely(points_wanted <= 0)) points_wanted = 1;
1428
query_debug_log(":auto natural points_wanted %ld", points_wanted);
1429
}
1430
1451
1452
// the available points of the query
1453
long points_available = (duration + 1) / query_granularity;
1454
+ if(unlikely(points_available <= 0)) points_available = 1;
1455
query_debug_log(":points_available %ld", points_available);
1456
1457
if(points_wanted > points_available) {
1481
if(points_wanted * group < points_available)
1482
points_wanted++;
1483
1484
+ if(unlikely(points_wanted <= 0))
1485
+ points_wanted = 1;
1486
+
1487
query_debug_log(":optimal points %ld", points_wanted);
1488
}
1489
1584
r->internal.points_wanted = points_wanted;
1585
r->internal.resampling_group = resampling_group;
1586
r->internal.resampling_divisor = resampling_divisor;
1135
-
1587
+ r->internal.query_options = options;
1588
+ r->internal.query_tier = tier;
1589
1590
// -------------------------------------------------------------------------
1591
// assign the processor functions
1632
// reset the grouping for the new dimension
1633
r->internal.grouping_reset(r);
1634
1182
- rrd2rrdr_do_dimension(r, points_wanted, rd, c, after_wanted, before_wanted, options);
1635
+ rrd2rrdr_do_dimension(r, points_wanted, rd, c, after_wanted, before_wanted);
1636
if (timeout)
1637
now_realtime_timeval(&query_current_time);
1638
1669
}
1670
1671
dimensions_used++;
1219
- if (timeout && (dt_usec(&query_start_time, &query_current_time) / 1000.0) > timeout) {
1672
+ if (timeout && ((NETDATA_DOUBLE)dt_usec(&query_start_time, &query_current_time) / 1000.0) > timeout) {
1673
log_access("QUERY CANCELED RUNTIME EXCEEDED %0.2f ms (LIMIT %d ms)",
1221
- dt_usec(&query_start_time, &query_current_time) / 1000.0, timeout);
1674
+ (NETDATA_DOUBLE)dt_usec(&query_start_time, &query_current_time) / 1000.0, timeout);
1675
r->result_options |= RRDR_RESULT_OPTION_CANCEL;
1676
break;
1677
}
1680
#ifdef NETDATA_INTERNAL_CHECKS
1681
if (dimensions_used) {
1682
if(r->internal.log)
1230
- rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group,
1683
+ rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group,
1684
after_wanted, after_requested, before_wanted, before_requested,
1685
points_requested, points_wanted, /*after_slot, before_slot,*/
1686
r->internal.log);
1687
1688
if(r->rows != points_wanted)
1236
- rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group,
1689
+ rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group,
1690
after_wanted, after_requested, before_wanted, before_requested,
1691
points_requested, points_wanted, /*after_slot, before_slot,*/
1692
"got 'points' is not wanted 'points'");
1693
1694
if(aligned && (r->before % (group * query_granularity)) != 0)
1242
- rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group,
1695
+ rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group,
1696
after_wanted, after_requested, before_wanted,before_wanted,
1697
points_requested, points_wanted, /*after_slot, before_slot,*/
1698
"'before' is not aligned but alignment is required");
1699
1700
// 'after' should not be aligned, since we start inside the first group
1701
//if(aligned && (r->after % group) != 0)
1249
- // rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group, after_wanted, after_requested, before_wanted, before_requested, points_requested, points_wanted, after_slot, before_slot, "'after' is not aligned but alignment is required");
1702
+ // rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group, after_wanted, after_requested, before_wanted, before_requested, points_requested, points_wanted, after_slot, before_slot, "'after' is not aligned but alignment is required");
1703
1704
if(r->before != before_wanted)
1252
- rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group,
1705
+ rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group,
1706
after_wanted, after_requested, before_wanted, before_requested,
1707
points_requested, points_wanted, /*after_slot, before_slot,*/
1708
"chart is not aligned to requested 'before'");
1709
1710
if(r->before != before_wanted)
1258
- rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group,
1711
+ rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group,
1712
after_wanted, after_requested, before_wanted, before_requested,
1713
points_requested, points_wanted, /*after_slot, before_slot,*/
1714
"got 'before' is not wanted 'before'");
1715
1716
// reported 'after' varies, depending on group
1717
if(r->after != after_wanted)
1265
- rrd2rrdr_log_request_response_metadata(r, group_method, aligned, group, resampling_time_requested, resampling_group,
1718
+ rrd2rrdr_log_request_response_metadata(r, options, group_method, aligned, group, resampling_time_requested, resampling_group,
1719
after_wanted, after_requested, before_wanted, before_requested,
1720
points_requested, points_wanted, /*after_slot, before_slot,*/
1721
"got 'after' is not wanted 'after'");
1722
+
1723
}
1724
#endif
1725