Journal sampling (#16433)
* initial work to support sampling in journal logs * always sample and entry that is useful as an item to return * add [unsampled] to the dashboard * facets sampling updates only the histogram * add unsampled to facets items * updated default query
Costa Tsaousis committed
Nov 20, 2023 at 10:13 UTC
8d83b5897b9ef5747c8c9a73479edaa042ea51be
4 files changed
+418
-72
collectors/systemd-journal.plugin/systemd-journal.c
+297
-50
@@ -22,6 +22,9 @@
22
#define SYSTEMD_JOURNAL_MAX_PARAMS 1000
23
#define SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION (1 * 3600)
24
#define SYSTEMD_JOURNAL_DEFAULT_ITEMS_PER_QUERY 200
25
+#define SYSTEMD_JOURNAL_DEFAULT_ITEMS_SAMPLING 1000000
26
+#define SYSTEMD_JOURNAL_SAMPLING_SLOTS 1000
27
+#define SYSTEMD_JOURNAL_SAMPLING_RECALIBRATE 10000
28
29
#define JOURNAL_PARAMETER_HELP "help"
30
#define JOURNAL_PARAMETER_AFTER "after"
@@ -41,6 +44,7 @@
44
#define JOURNAL_PARAMETER_SLICE "slice"
45
#define JOURNAL_PARAMETER_DELTA "delta"
46
#define JOURNAL_PARAMETER_TAIL "tail"
47
+#define JOURNAL_PARAMETER_SAMPLING "sampling"
48
49
#define JOURNAL_KEY_ND_JOURNAL_FILE "ND_JOURNAL_FILE"
50
#define JOURNAL_KEY_ND_JOURNAL_PROCESS "ND_JOURNAL_PROCESS"
@@ -189,11 +193,37 @@ typedef struct function_query_status {
193
bool tail;
194
bool data_only;
195
bool slice;
196
+ size_t sampling;
197
size_t filters;
198
usec_t last_modified;
199
const char *query;
200
const char *histogram;
201
202
+ struct {
203
+ uint32_t enable_after_samples;
204
+ uint32_t slots;
205
+ uint32_t sampled;
206
+ uint32_t unsampled;
207
+ } samples;
208
+
209
+ struct {
210
+ uint32_t enable_after_samples;
211
+ uint32_t every;
212
+ uint32_t skipped;
213
+ uint32_t recalibrate;
214
+ uint32_t sampled;
215
+ uint32_t unsampled;
216
+ } samples_per_file;
217
+
218
+ struct {
219
+ usec_t start_ut;
220
+ usec_t end_ut;
221
+ usec_t step_ut;
222
+ uint32_t enable_after_samples;
223
+ uint32_t sampled[SYSTEMD_JOURNAL_SAMPLING_SLOTS];
224
+ uint32_t unsampled[SYSTEMD_JOURNAL_SAMPLING_SLOTS];
225
+ } samples_per_time_slot;
226
+
227
// per file progress info
228
// size_t cached_count;
229
@@ -236,6 +266,172 @@ static inline bool netdata_systemd_journal_seek_to(sd_journal *j, usec_t timesta
266
267
#define JD_SOURCE_REALTIME_TIMESTAMP "_SOURCE_REALTIME_TIMESTAMP"
268
269
+// ----------------------------------------------------------------------------
270
+// sampling support
271
+
272
+static void sampling_query_init(FUNCTION_QUERY_STATUS *fqs, FACETS *facets) {
273
+ if(!fqs->sampling)
274
+ return;
275
+
276
+ if(!fqs->slice) {
277
+ // the user is doing a full data query
278
+ // disable sampling
279
+ fqs->sampling = 0;
280
+ return;
281
+ }
282
+
283
+ if(fqs->data_only) {
284
+ // the user is doing a data query
285
+ // disable sampling
286
+ fqs->sampling = 0;
287
+ return;
288
+ }
289
+
290
+ if(!fqs->files_matched) {
291
+ // no files have been matched
292
+ // disable sampling
293
+ fqs->sampling = 0;
294
+ return;
295
+ }
296
+
297
+ fqs->samples.slots = facets_histogram_slots(facets);
298
+ if(fqs->samples.slots < 2) fqs->samples.slots = 2;
299
+ if(fqs->samples.slots > SYSTEMD_JOURNAL_SAMPLING_SLOTS)
300
+ fqs->samples.slots = SYSTEMD_JOURNAL_SAMPLING_SLOTS;
301
+
302
+ if(!fqs->after_ut || !fqs->before_ut || fqs->after_ut >= fqs->before_ut) {
303
+ // we don't have enough information for sampling
304
+ fqs->sampling = 0;
305
+ return;
306
+ }
307
+
308
+ usec_t delta = fqs->before_ut - fqs->after_ut;
309
+ usec_t step = delta / facets_histogram_slots(facets) - 1;
310
+ if(step < 1) step = 1;
311
+
312
+ fqs->samples_per_time_slot.start_ut = fqs->after_ut;
313
+ fqs->samples_per_time_slot.end_ut = fqs->before_ut;
314
+ fqs->samples_per_time_slot.step_ut = step;
315
+
316
+ // the minimum number of rows to enable sampling
317
+ fqs->samples.enable_after_samples = fqs->sampling / 2;
318
+
319
+ // the minimum number of rows per file to enable sampling
320
+ fqs->samples_per_file.enable_after_samples = (fqs->sampling / 4) / fqs->files_matched;
321
+ if(fqs->samples_per_file.enable_after_samples < fqs->entries * 2)
322
+ fqs->samples_per_file.enable_after_samples = fqs->entries * 2;
323
+
324
+ // the minimum number of rows per time slot to enable sampling
325
+ fqs->samples_per_time_slot.enable_after_samples = (fqs->sampling / 4) / fqs->samples.slots;
326
+ if(fqs->samples_per_time_slot.enable_after_samples < fqs->entries)
327
+ fqs->samples_per_time_slot.enable_after_samples = fqs->entries;
328
+}
329
+
330
+static void sampling_file_init(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf __maybe_unused) {
331
+ fqs->samples_per_file.sampled = 0;
332
+ fqs->samples_per_file.unsampled = 0;
333
+ fqs->samples_per_file.every = 0;
334
+ fqs->samples_per_file.skipped = 0;
335
+ fqs->samples_per_file.recalibrate = 0;
336
+}
337
+
338
+static void sampling_decide_file_sampling_every(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
339
+ size_t sampled = fqs->samples_per_file.sampled + fqs->samples_per_file.unsampled;
340
+
341
+ if(!sampled)
342
+ sampled = 1;
343
+
344
+ // find the common duration
345
+ usec_t after_ut = fqs->samples_per_time_slot.start_ut < jf->msg_first_ut ? jf->msg_first_ut : fqs->samples_per_time_slot.start_ut;
346
+ usec_t before_ut = fqs->samples_per_time_slot.end_ut > jf->msg_last_ut ? jf->msg_last_ut : fqs->samples_per_time_slot.end_ut;
347
+
348
+ if(after_ut > before_ut) {
349
+ usec_t t = after_ut;
350
+ after_ut = before_ut;
351
+ before_ut = t;
352
+ }
353
+
354
+ if(after_ut == before_ut)
355
+ after_ut = before_ut - 1;
356
+
357
+ if(msg_ut <= after_ut)
358
+ msg_ut = after_ut + 1;
359
+
360
+ if(msg_ut >= before_ut)
361
+ msg_ut = before_ut - 1;
362
+
363
+ size_t expected_lines;
364
+
365
+ if(direction == FACETS_ANCHOR_DIRECTION_FORWARD)
366
+ expected_lines = sampled * (before_ut - after_ut) / (msg_ut - after_ut);
367
+ else
368
+ expected_lines = sampled * (before_ut - after_ut) / (before_ut - msg_ut);
369
+
370
+ if(expected_lines < 1)
371
+ expected_lines = 1;
372
+
373
+ size_t wanted_samples = (fqs->sampling / 2) / fqs->files_matched;
374
+
375
+ fqs->samples_per_file.every = expected_lines / wanted_samples;
376
+
377
+ if(fqs->samples_per_file.every < 1)
378
+ fqs->samples_per_file.every = 1;
379
+}
380
+
381
+static inline bool is_row_in_sample(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, usec_t msg_ut, FACETS_ANCHOR_DIRECTION direction, bool candidate_to_keep) {
382
+ if(!fqs->sampling)
383
+ return true;
384
+
385
+ if(unlikely(msg_ut < fqs->samples_per_time_slot.start_ut))
386
+ msg_ut = fqs->samples_per_time_slot.start_ut;
387
+ if(unlikely(msg_ut > fqs->samples_per_time_slot.end_ut))
388
+ msg_ut = fqs->samples_per_time_slot.end_ut;
389
+
390
+ size_t slot = (msg_ut - fqs->samples_per_time_slot.start_ut) / fqs->samples_per_time_slot.step_ut;
391
+ if(slot >= fqs->samples.slots)
392
+ slot = fqs->samples.slots - 1;
393
+
394
+ bool should_sample = candidate_to_keep;
395
+
396
+ if(fqs->samples.sampled < fqs->samples.enable_after_samples ||
397
+ fqs->samples_per_file.sampled < fqs->samples_per_file.enable_after_samples ||
398
+ fqs->samples_per_time_slot.sampled[slot] < fqs->samples_per_time_slot.enable_after_samples)
399
+ should_sample = true;
400
+
401
+ else if(fqs->samples_per_file.recalibrate >= SYSTEMD_JOURNAL_SAMPLING_RECALIBRATE || !fqs->samples_per_file.every) {
402
+ // this is the first to be unsampled for this file
403
+ sampling_decide_file_sampling_every(fqs, jf, direction, msg_ut);
404
+ fqs->samples_per_file.recalibrate = 0;
405
+ should_sample = true;
406
+ }
407
+ else {
408
+ // we sample 1 every fqs->samples_per_file.every
409
+ if(fqs->samples_per_file.skipped >= fqs->samples_per_file.every) {
410
+ fqs->samples_per_file.skipped = 0;
411
+ should_sample = true;
412
+ }
413
+ else
414
+ fqs->samples_per_file.skipped++;
415
+ }
416
+
417
+ fqs->samples_per_file.recalibrate++;
418
+
419
+ if(should_sample) {
420
+ fqs->samples.sampled++;
421
+ fqs->samples_per_file.sampled++;
422
+ fqs->samples_per_time_slot.sampled[slot]++;
423
+ }
424
+ else {
425
+ fqs->samples.unsampled++;
426
+ fqs->samples_per_file.unsampled++;
427
+ fqs->samples_per_time_slot.unsampled[slot]++;
428
+ }
429
+
430
+ return should_sample;
431
+}
432
+
433
+// ----------------------------------------------------------------------------
434
+
435
static inline bool parse_journal_field(const char *data, size_t data_length, const char **key, size_t *key_length, const char **value, size_t *value_length) {
436
const char *k = data;
437
const char *equal = strchr(k, '=');
@@ -361,36 +557,44 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
557
if (unlikely(msg_ut < stop_ut))
558
break;
559
364
- bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
560
+ bool to_sample = is_row_in_sample(fqs, jf, msg_ut,
561
+ FACETS_ANCHOR_DIRECTION_BACKWARD,
562
+ facets_row_candidate_to_keep(facets, msg_ut));
563
366
- // make sure each line gets a unique timestamp
367
- if(unlikely(msg_ut >= last_usec_from && msg_ut <= last_usec_to))
368
- msg_ut = --last_usec_from;
369
- else
370
- last_usec_from = last_usec_to = msg_ut;
371
-
372
- if(facets_row_finished(facets, msg_ut))
373
- rows_useful++;
374
-
375
- row_counter++;
376
- if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
377
- stop_when_full &&
378
- facets_rows(facets) >= fqs->entries)) {
379
- // stop the data only query
380
- usec_t oldest = facets_row_oldest_ut(facets);
381
- if(oldest && msg_ut < (oldest - anchor_delta))
382
- break;
383
- }
564
+ if(to_sample) {
565
+ bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
566
+
567
+ // make sure each line gets a unique timestamp
568
+ if(unlikely(msg_ut >= last_usec_from && msg_ut <= last_usec_to))
569
+ msg_ut = --last_usec_from;
570
+ else
571
+ last_usec_from = last_usec_to = msg_ut;
572
+
573
+ if(facets_row_finished(facets, msg_ut))
574
+ rows_useful++;
575
+
576
+ row_counter++;
577
+ if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
578
+ stop_when_full &&
579
+ facets_rows(facets) >= fqs->entries)) {
580
+ // stop the data only query
581
+ usec_t oldest = facets_row_oldest_ut(facets);
582
+ if(oldest && msg_ut < (oldest - anchor_delta))
583
+ break;
584
+ }
585
385
- if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
386
- FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
387
- last_row_counter = row_counter;
586
+ if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
587
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
588
+ last_row_counter = row_counter;
589
389
- FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
390
- last_bytes = bytes;
590
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
591
+ last_bytes = bytes;
592
392
- status = check_stop(fqs->cancelled, &fqs->stop_monotonic_ut);
593
+ status = check_stop(fqs->cancelled, &fqs->stop_monotonic_ut);
594
+ }
595
}
596
+ else
597
+ facets_row_finished_unsampled(facets, msg_ut);
598
}
599
600
FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
@@ -447,36 +651,44 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
651
if (unlikely(msg_ut > stop_ut))
652
break;
653
450
- bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
654
+ bool to_sample = is_row_in_sample(fqs, jf, msg_ut,
655
+ FACETS_ANCHOR_DIRECTION_FORWARD,
656
+ facets_row_candidate_to_keep(facets, msg_ut));
657
452
- // make sure each line gets a unique timestamp
453
- if(unlikely(msg_ut >= last_usec_from && msg_ut <= last_usec_to))
454
- msg_ut = ++last_usec_to;
455
- else
456
- last_usec_from = last_usec_to = msg_ut;
457
-
458
- if(facets_row_finished(facets, msg_ut))
459
- rows_useful++;
460
-
461
- row_counter++;
462
- if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
463
- stop_when_full &&
464
- facets_rows(facets) >= fqs->entries)) {
465
- // stop the data only query
466
- usec_t newest = facets_row_newest_ut(facets);
467
- if(newest && msg_ut > (newest + anchor_delta))
468
- break;
469
- }
658
+ if(to_sample) {
659
+ bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
660
+
661
+ // make sure each line gets a unique timestamp
662
+ if(unlikely(msg_ut >= last_usec_from && msg_ut <= last_usec_to))
663
+ msg_ut = ++last_usec_to;
664
+ else
665
+ last_usec_from = last_usec_to = msg_ut;
666
+
667
+ if(facets_row_finished(facets, msg_ut))
668
+ rows_useful++;
669
+
670
+ row_counter++;
671
+ if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
672
+ stop_when_full &&
673
+ facets_rows(facets) >= fqs->entries)) {
674
+ // stop the data only query
675
+ usec_t newest = facets_row_newest_ut(facets);
676
+ if(newest && msg_ut > (newest + anchor_delta))
677
+ break;
678
+ }
679
471
- if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
472
- FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
473
- last_row_counter = row_counter;
680
+ if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
681
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
682
+ last_row_counter = row_counter;
683
475
- FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
476
- last_bytes = bytes;
684
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
685
+ last_bytes = bytes;
686
478
- status = check_stop(fqs->cancelled, &fqs->stop_monotonic_ut);
687
+ status = check_stop(fqs->cancelled, &fqs->stop_monotonic_ut);
688
+ }
689
}
690
+ else
691
+ facets_row_finished_unsampled(facets, msg_ut);
692
}
693
694
FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
@@ -698,6 +910,8 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
910
usec_t started_ut;
911
usec_t ended_ut = now_monotonic_usec();
912
913
+ sampling_query_init(fqs, facets);
914
+
915
buffer_json_member_add_array(wb, "_journal_files");
916
for(size_t f = 0; f < files_used ;f++) {
917
const char *filename = dictionary_acquired_item_name(file_items[f]);
@@ -716,6 +930,8 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
930
size_t bytes_read = fqs->bytes_read;
931
size_t matches_setup_ut = fqs->matches_setup_ut;
932
933
+ sampling_file_init(fqs, jf);
934
+
935
ND_SD_JOURNAL_STATUS tmp_status = netdata_systemd_journal_query_one_file(filename, wb, facets, jf, fqs);
936
937
rows_useful = fqs->rows_useful - rows_useful;
@@ -750,6 +966,15 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
966
buffer_json_member_add_uint64(wb, "duration_matches_ut", matches_setup_ut);
967
buffer_json_member_add_uint64(wb, "fstat_query_calls", fs_calls);
968
buffer_json_member_add_uint64(wb, "fstat_query_cached_responses", fs_cached);
969
+
970
+ if(fqs->sampling) {
971
+ buffer_json_member_add_object(wb, "_sampling");
972
+ {
973
+ buffer_json_member_add_uint64(wb, "sampled", fqs->samples_per_file.sampled);
974
+ buffer_json_member_add_uint64(wb, "unsampled", fqs->samples_per_file.unsampled);
975
+ }
976
+ buffer_json_object_close(wb); // _sampling
977
+ }
978
}
979
buffer_json_object_close(wb); // journal file
980
@@ -838,6 +1063,16 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1063
buffer_json_member_add_uint64(wb, "cached", fstat_thread_cached_responses);
1064
}
1065
buffer_json_object_close(wb); // _fstat_caching
1066
+
1067
+ if(fqs->sampling) {
1068
+ buffer_json_member_add_object(wb, "_sampling");
1069
+ {
1070
+ buffer_json_member_add_uint64(wb, "sampled", fqs->samples.sampled);
1071
+ buffer_json_member_add_uint64(wb, "unsampled", fqs->samples.unsampled);
1072
+ }
1073
+ buffer_json_object_close(wb); // _sampling
1074
+ }
1075
+
1076
buffer_json_finalize(wb);
1077
1078
return HTTP_RESP_OK;
@@ -906,6 +1141,10 @@ static void netdata_systemd_journal_function_help(const char *transaction) {
1141
" The number of items to return.\n"
1142
" The default is %d.\n"
1143
"\n"
1144
+ " "JOURNAL_PARAMETER_SAMPLING":ITEMS\n"
1145
+ " The number of log entries to sample to estimate facets counters and histogram.\n"
1146
+ " The default is %d.\n"
1147
+ "\n"
1148
" "JOURNAL_PARAMETER_ANCHOR":TIMESTAMP_IN_MICROSECONDS\n"
1149
" Return items relative to this timestamp.\n"
1150
" The exact items to be returned depend on the query `"JOURNAL_PARAMETER_DIRECTION"`.\n"
@@ -946,6 +1185,7 @@ static void netdata_systemd_journal_function_help(const char *transaction) {
1185
, JOURNAL_DEFAULT_SLICE_MODE ? "true" : "false" // slice
1186
, -SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION
1187
, SYSTEMD_JOURNAL_DEFAULT_ITEMS_PER_QUERY
1188
+ , SYSTEMD_JOURNAL_DEFAULT_ITEMS_SAMPLING
1189
, JOURNAL_DEFAULT_DIRECTION == FACETS_ANCHOR_DIRECTION_BACKWARD ? "backward" : "forward"
1190
);
1191
@@ -1055,6 +1295,7 @@ void function_systemd_journal(const char *transaction, char *function, int timeo
1295
facets_accepted_param(facets, JOURNAL_PARAMETER_PROGRESS);
1296
facets_accepted_param(facets, JOURNAL_PARAMETER_DELTA);
1297
facets_accepted_param(facets, JOURNAL_PARAMETER_TAIL);
1298
+ facets_accepted_param(facets, JOURNAL_PARAMETER_SAMPLING);
1299
1300
#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
1301
facets_accepted_param(facets, JOURNAL_PARAMETER_SLICE);
@@ -1170,6 +1411,7 @@ void function_systemd_journal(const char *transaction, char *function, int timeo
1411
const char *progress_id = NULL;
1412
SD_JOURNAL_FILE_SOURCE_TYPE source_type = SDJF_ALL;
1413
size_t filters = 0;
1414
+ size_t sampling = SYSTEMD_JOURNAL_DEFAULT_ITEMS_SAMPLING;
1415
1416
buffer_json_member_add_object(wb, "_request");
1417
@@ -1205,6 +1447,9 @@ void function_systemd_journal(const char *transaction, char *function, int timeo
1447
else
1448
tail = true;
1449
}
1450
+ else if(strncmp(keyword, JOURNAL_PARAMETER_SAMPLING ":", sizeof(JOURNAL_PARAMETER_SAMPLING ":") - 1) == 0) {
1451
+ sampling = str2ul(&keyword[sizeof(JOURNAL_PARAMETER_SAMPLING ":") - 1]);
1452
+ }
1453
else if(strncmp(keyword, JOURNAL_PARAMETER_DATA_ONLY ":", sizeof(JOURNAL_PARAMETER_DATA_ONLY ":") - 1) == 0) {
1454
char *v = &keyword[sizeof(JOURNAL_PARAMETER_DATA_ONLY ":") - 1];
1455
@@ -1415,6 +1660,7 @@ void function_systemd_journal(const char *transaction, char *function, int timeo
1660
fqs->direction = direction;
1661
fqs->anchor.start_ut = anchor;
1662
fqs->anchor.stop_ut = 0;
1663
+ fqs->sampling = sampling;
1664
1665
if(fqs->anchor.start_ut && fqs->tail) {
1666
// a tail request
@@ -1477,6 +1723,7 @@ void function_systemd_journal(const char *transaction, char *function, int timeo
1723
buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_PROGRESS, false);
1724
buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_DELTA, fqs->delta);
1725
buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_TAIL, fqs->tail);
1726
+ buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_SAMPLING, fqs->sampling);
1727
buffer_json_member_add_string(wb, JOURNAL_PARAMETER_ID, progress_id);
1728
buffer_json_member_add_uint64(wb, "source_type", fqs->source_type);
1729
buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_AFTER, fqs->after_ut / USEC_PER_SEC);
collectors/systemd-journal.plugin/systemd-main.c
+1
-1
@@ -38,7 +38,7 @@ int main(int argc __maybe_unused, char **argv __maybe_unused) {
38
39
if(argc == 2 && strcmp(argv[1], "debug") == 0) {
40
bool cancelled = false;
41
- char buf[] = "systemd-journal after:-16000000 before:0 last:1";
41
+ char buf[] = "systemd-journal after:-8640000 before:0 direction:backward last:200 data_only:false slice:true source:all";
42
// char buf[] = "systemd-journal after:1695332964 before:1695937764 direction:backward last:100 slice:true source:all DHKucpqUoe1:PtVoyIuX.MU";
43
// char buf[] = "systemd-journal after:1694511062 before:1694514662 anchor:1694514122024403";
44
function_systemd_journal("123", buf, 600, &cancelled);
libnetdata/facets/facets.c
+115
-21
@@ -1,13 +1,15 @@
1
// SPDX-License-Identifier: GPL-3.0-or-later
2
#include "facets.h"
3
4
-#define HISTOGRAM_COLUMNS 150 // the target number of points in a histogram
4
+#define FACETS_HISTOGRAM_COLUMNS 150 // the target number of points in a histogram
5
#define FACETS_KEYS_WITH_VALUES_MAX 200 // the max number of keys that can be facets
6
#define FACETS_KEYS_IN_ROW_MAX 500 // the max number of keys in a row
7
8
#define FACETS_KEYS_HASHTABLE_ENTRIES 15
9
#define FACETS_VALUES_HASHTABLE_ENTRIES 15
10
11
+static inline void facets_reset_key(FACET_KEY *k);
12
+
13
// ----------------------------------------------------------------------------
14
15
static const char id_encoding_characters[64 + 1] = "ABCDEFGHIJKLMNOPQRSTUVWXYZ.abcdefghijklmnopqrstuvwxyz_0123456789";
@@ -34,6 +36,7 @@ static const uint8_t id_encoding_characters_reverse[256] = {
36
#define FACETS_HASH XXH64_hash_t
37
#define FACETS_HASH_FUNCTION(src, len) XXH3_64bits(src, len)
38
#define FACETS_HASH_ZERO (FACETS_HASH)0
39
+#define FACETS_HASH_UNSAMPLED (FACETS_HASH)UINT64_MAX
40
41
static inline void facets_hash_to_str(FACETS_HASH num, char *out) {
42
out[11] = '\0';
@@ -192,6 +195,7 @@ typedef struct facet_value {
195
196
bool selected;
197
bool empty;
198
+ bool unsampled;
199
200
uint32_t rows_matching_facet_value;
201
uint32_t final_facet_value_counter;
@@ -204,14 +208,17 @@ typedef struct facet_value {
208
} FACET_VALUE;
209
210
typedef enum {
207
- FACET_KEY_VALUE_NONE = 0,
208
- FACET_KEY_VALUE_UPDATED = (1 << 0),
209
- FACET_KEY_VALUE_EMPTY = (1 << 1),
210
- FACET_KEY_VALUE_COPIED = (1 << 2),
211
+ FACET_KEY_VALUE_NONE = 0,
212
+ FACET_KEY_VALUE_UPDATED = (1 << 0),
213
+ FACET_KEY_VALUE_EMPTY = (1 << 1),
214
+ FACET_KEY_VALUE_UNSAMPLED = (1 << 2),
215
+ FACET_KEY_VALUE_COPIED = (1 << 3),
216
} FACET_KEY_VALUE_FLAGS;
217
218
#define facet_key_value_updated(k) ((k)->current_value.flags & FACET_KEY_VALUE_UPDATED)
219
#define facet_key_value_empty(k) ((k)->current_value.flags & FACET_KEY_VALUE_EMPTY)
220
+#define facet_key_value_unsampled(k) ((k)->current_value.flags & FACET_KEY_VALUE_UNSAMPLED)
221
+#define facet_key_value_empty_or_unsampled(k) ((k)->current_value.flags & (FACET_KEY_VALUE_EMPTY|FACET_KEY_VALUE_UNSAMPLED))
222
#define facet_key_value_copied(k) ((k)->current_value.flags & FACET_KEY_VALUE_COPIED)
223
224
struct facet_key {
@@ -249,6 +256,10 @@ struct facet_key {
256
FACET_VALUE *v;
257
} empty_value;
258
259
+ struct {
260
+ FACET_VALUE *v;
261
+ } unsampled_value;
262
+
263
struct {
264
facet_dynamic_row_t cb;
265
void *data;
@@ -347,6 +358,7 @@ struct facets {
358
struct {
359
size_t evaluated;
360
size_t matched;
361
+ size_t unsampled;
362
size_t created;
363
size_t reused;
364
} rows;
@@ -361,6 +373,7 @@ struct facets {
373
size_t transformed;
374
size_t dynamic;
375
size_t empty;
376
+ size_t unsampled;
377
size_t indexed;
378
size_t inserts;
379
size_t conflicts;
@@ -515,6 +528,27 @@ static inline FACET_VALUE *FACET_VALUE_ADD_TO_INDEX(FACET_KEY *k, const FACET_VA
528
return v;
529
}
530
531
+static inline void FACET_VALUE_ADD_UNSAMPLED_VALUE_TO_INDEX(FACET_KEY *k) {
532
+ static const FACET_VALUE tv = {
533
+ .hash = FACETS_HASH_UNSAMPLED,
534
+ .name = FACET_VALUE_UNSAMPLED,
535
+ .name_len = sizeof(FACET_VALUE_UNSAMPLED) - 1,
536
+ };
537
+
538
+ k->current_value.hash = FACETS_HASH_UNSAMPLED;
539
+
540
+ if(k->unsampled_value.v) {
541
+ FACET_VALUE_ADD_CONFLICT(k, k->unsampled_value.v, &tv);
542
+ k->current_value.v = k->unsampled_value.v;
543
+ }
544
+ else {
545
+ FACET_VALUE *v = FACET_VALUE_ADD_TO_INDEX(k, &tv);
546
+ v->unsampled = true;
547
+ k->unsampled_value.v = v;
548
+ k->current_value.v = v;
549
+ }
550
+}
551
+
552
static inline void FACET_VALUE_ADD_EMPTY_VALUE_TO_INDEX(FACET_KEY *k) {
553
static const FACET_VALUE tv = {
554
.hash = FACETS_HASH_ZERO,
@@ -744,7 +778,7 @@ static usec_t calculate_histogram_bar_width(usec_t after_ut, usec_t before_ut) {
778
usec_t bar_width_ut = 1 * USEC_PER_SEC;
779
780
for (int i = array_size - 1; i >= 0; --i) {
747
- if (duration_ut / (valid_durations_s[i] * USEC_PER_SEC) >= HISTOGRAM_COLUMNS) {
781
+ if (duration_ut / (valid_durations_s[i] * USEC_PER_SEC) >= FACETS_HISTOGRAM_COLUMNS) {
782
bar_width_ut = valid_durations_s[i] * USEC_PER_SEC;
783
break;
784
}
@@ -758,6 +792,10 @@ static inline usec_t facets_histogram_slot_baseline_ut(FACETS *facets, usec_t ut
792
return ut - delta_ut;
793
}
794
795
+size_t facets_histogram_slots(FACETS *facets) {
796
+ return facets->histogram.slots;
797
+}
798
+
799
void facets_set_timeframe_and_histogram_by_id(FACETS *facets, const char *key_id, usec_t after_ut, usec_t before_ut) {
800
if(after_ut > before_ut) {
801
usec_t t = after_ut;
@@ -801,17 +839,7 @@ void facets_set_timeframe_and_histogram_by_name(FACETS *facets, const char *key_
839
facets_set_timeframe_and_histogram_by_id(facets, hash_str, after_ut, before_ut);
840
}
841
804
-static inline void facets_histogram_update_value(FACETS *facets, usec_t usec) {
805
- if(!facets->histogram.enabled ||
806
- !facets->histogram.key ||
807
- !facets->histogram.key->values.enabled ||
808
- !facet_key_value_updated(facets->histogram.key) ||
809
- usec < facets->histogram.after_ut ||
810
- usec > facets->histogram.before_ut)
811
- return;
812
-
813
- FACET_VALUE *v = facets->histogram.key->current_value.v;
814
-
842
+static inline void facets_histogram_update_value_slot(FACETS *facets, usec_t usec, FACET_VALUE *v) {
843
if(unlikely(!v->histogram))
844
v->histogram = callocz(facets->histogram.slots, sizeof(*v->histogram));
845
@@ -831,6 +859,40 @@ static inline void facets_histogram_update_value(FACETS *facets, usec_t usec) {
859
v->histogram[slot]++;
860
}
861
862
+static inline void facets_histogram_update_value(FACETS *facets, usec_t usec) {
863
+ if(!facets->histogram.enabled ||
864
+ !facets->histogram.key ||
865
+ !facets->histogram.key->values.enabled ||
866
+ !facet_key_value_updated(facets->histogram.key) ||
867
+ usec < facets->histogram.after_ut ||
868
+ usec > facets->histogram.before_ut)
869
+ return;
870
+
871
+ FACET_VALUE *v = facets->histogram.key->current_value.v;
872
+ facets_histogram_update_value_slot(facets, usec, v);
873
+}
874
+
875
+void facets_row_finished_unsampled(FACETS *facets, usec_t usec) {
876
+ facets->operations.rows.evaluated++;
877
+ facets->operations.rows.matched++;
878
+ facets->operations.rows.unsampled++;
879
+
880
+ if(!facets->histogram.enabled ||
881
+ !facets->histogram.key ||
882
+ !facets->histogram.key->values.enabled ||
883
+ usec < facets->histogram.after_ut ||
884
+ usec > facets->histogram.before_ut)
885
+ return;
886
+
887
+ if(!facets->histogram.key->unsampled_value.v)
888
+ FACET_VALUE_ADD_UNSAMPLED_VALUE_TO_INDEX(facets->histogram.key);
889
+
890
+ FACET_VALUE *v = facets->histogram.key->unsampled_value.v;
891
+ facets_histogram_update_value_slot(facets, usec, v);
892
+
893
+ facets_reset_key(facets->histogram.key);
894
+}
895
+
896
static inline void facets_histogram_value_names(BUFFER *wb, FACETS *facets __maybe_unused, FACET_KEY *k, const char *key, const char *first_key) {
897
BUFFER *tb = NULL;
898
@@ -1538,6 +1600,30 @@ void facets_set_additional_options(FACETS *facets, FACETS_OPTIONS options) {
1600
1601
// ----------------------------------------------------------------------------
1602
1603
+static inline void facets_key_set_unsampled_value(FACETS *facets, FACET_KEY *k) {
1604
+ if(likely(!facet_key_value_updated(k) && facets->keys_in_row.used < FACETS_KEYS_IN_ROW_MAX))
1605
+ facets->keys_in_row.array[facets->keys_in_row.used++] = k;
1606
+
1607
+ k->current_value.flags |= FACET_KEY_VALUE_UPDATED | FACET_KEY_VALUE_UNSAMPLED;
1608
+
1609
+ facets->operations.values.registered++;
1610
+ facets->operations.values.unsampled++;
1611
+
1612
+ // no need to copy the UNSET value
1613
+ // empty values are exported as empty
1614
+ k->current_value.raw = NULL;
1615
+ k->current_value.raw_len = 0;
1616
+ k->current_value.b->len = 0;
1617
+ k->current_value.flags &= ~FACET_KEY_VALUE_COPIED;
1618
+
1619
+ if(unlikely(k->values.enabled))
1620
+ FACET_VALUE_ADD_UNSAMPLED_VALUE_TO_INDEX(k);
1621
+ else {
1622
+ k->key_found_in_row++;
1623
+ k->key_values_selected_in_row++;
1624
+ }
1625
+}
1626
+
1627
static inline void facets_key_set_empty_value(FACETS *facets, FACET_KEY *k) {
1628
if(likely(!facet_key_value_updated(k) && facets->keys_in_row.used < FACETS_KEYS_IN_ROW_MAX))
1629
facets->keys_in_row.array[facets->keys_in_row.used++] = k;
@@ -1567,7 +1653,7 @@ static inline void facets_key_check_value(FACETS *facets, FACET_KEY *k) {
1653
facets->keys_in_row.array[facets->keys_in_row.used++] = k;
1654
1655
k->current_value.flags |= FACET_KEY_VALUE_UPDATED;
1570
- k->current_value.flags &= ~FACET_KEY_VALUE_EMPTY;
1656
+ k->current_value.flags &= ~(FACET_KEY_VALUE_EMPTY|FACET_KEY_VALUE_UNSAMPLED);
1657
1658
facets->operations.values.registered++;
1659
@@ -1581,7 +1667,7 @@ static inline void facets_key_check_value(FACETS *facets, FACET_KEY *k) {
1667
// if(strstr(buffer_tostring(k->current_value), "fprintd") != NULL)
1668
// found = true;
1669
1584
- if(facets->query && !facet_key_value_empty(k) && ((k->options & FACET_KEY_OPTION_FTS) || facets->options & FACETS_OPTION_ALL_KEYS_FTS)) {
1670
+ if(facets->query && !facet_key_value_empty_or_unsampled(k) && ((k->options & FACET_KEY_OPTION_FTS) || facets->options & FACETS_OPTION_ALL_KEYS_FTS)) {
1671
facets->operations.fts.searches++;
1672
facets_key_value_copy_to_buffer(k);
1673
switch(simple_pattern_matches_extract(facets->query, buffer_tostring(k->current_value.b), NULL, 0)) {
@@ -1692,7 +1778,7 @@ static FACET_ROW *facets_row_create(FACETS *facets, usec_t usec, FACET_ROW *into
1778
.empty = true,
1779
};
1780
1695
- if(facet_key_value_updated(k) && !facet_key_value_empty(k)) {
1781
+ if(facet_key_value_updated(k) && !facet_key_value_empty_or_unsampled(k)) {
1782
t.tmp = facets_key_get_value(k);
1783
t.tmp_len = facets_key_get_value_length(k);
1784
t.empty = false;
@@ -1771,6 +1857,12 @@ static inline bool facets_is_entry_within_anchor(FACETS *facets, usec_t usec) {
1857
return true;
1858
}
1859
1860
+bool facets_row_candidate_to_keep(FACETS *facets, usec_t usec) {
1861
+ return !facets->base ||
1862
+ (usec >= facets->base->prev->usec && usec <= facets->base->usec && facets_is_entry_within_anchor(facets, usec)) ||
1863
+ facets->items_to_return < facets->max_items_to_return;
1864
+}
1865
+
1866
static void facets_row_keep(FACETS *facets, usec_t usec) {
1867
facets->operations.rows.matched++;
1868
@@ -1898,9 +1990,10 @@ bool facets_row_finished(FACETS *facets, usec_t usec) {
1990
for(size_t p = 0; p < entries ;p++) {
1991
FACET_KEY *k = facets->keys_with_values.array[p];
1992
1901
- if(!facet_key_value_updated(k))
1993
+ if(!facet_key_value_updated(k)) {
1994
// put the FACET_VALUE_UNSET value into it
1995
facets_key_set_empty_value(facets, k);
1996
+ }
1997
1998
total_keys++;
1999
@@ -2518,6 +2611,7 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2611
if(show_items) {
2612
buffer_json_member_add_uint64(wb, "evaluated", facets->operations.rows.evaluated);
2613
buffer_json_member_add_uint64(wb, "matched", facets->operations.rows.matched);
2614
+ buffer_json_member_add_uint64(wb, "unsampled", facets->operations.rows.unsampled);
2615
buffer_json_member_add_uint64(wb, "returned", facets->items_to_return);
2616
buffer_json_member_add_uint64(wb, "max_to_return", facets->max_items_to_return);
2617
buffer_json_member_add_uint64(wb, "before", facets->operations.skips_before);
libnetdata/facets/facets.h
+5
@@ -6,6 +6,7 @@
6
#include "../libnetdata.h"
7
8
#define FACET_VALUE_UNSET "-"
9
+#define FACET_VALUE_UNSAMPLED "[unsampled]"
10
11
typedef enum __attribute__((packed)) {
12
FACETS_ANCHOR_DIRECTION_FORWARD,
@@ -85,11 +86,15 @@ void facets_accepted_param(FACETS *facets, const char *param);
86
void facets_rows_begin(FACETS *facets);
87
bool facets_row_finished(FACETS *facets, usec_t usec);
88
89
+void facets_row_finished_unsampled(FACETS *facets, usec_t usec);
90
+size_t facets_histogram_slots(FACETS *facets);
91
+
92
FACET_KEY *facets_register_key_name(FACETS *facets, const char *key, FACET_KEY_OPTIONS options);
93
void facets_set_query(FACETS *facets, const char *query);
94
void facets_set_items(FACETS *facets, uint32_t items);
95
void facets_set_anchor(FACETS *facets, usec_t start_ut, usec_t stop_ut, FACETS_ANCHOR_DIRECTION direction);
96
void facets_enable_slice_mode(FACETS *facets);
97
+bool facets_row_candidate_to_keep(FACETS *facets, usec_t usec);
98
99
FACET_KEY *facets_register_facet_id(FACETS *facets, const char *key_id, FACET_KEY_OPTIONS options);
100
void facets_register_facet_id_filter(FACETS *facets, const char *key_id, char *value_id, FACET_KEY_OPTIONS options);