forked from kafka4beam/wolff
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathwolff_metrics.erl
More file actions
125 lines (101 loc) · 3.81 KB
/
Copy pathwolff_metrics.erl
File metadata and controls
125 lines (101 loc) · 3.81 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
-module(wolff_metrics).
-export([
inflight_set/2,
queuing_set/2,
queuing_bytes_set/2,
dropped_inc/1,
dropped_inc/2,
dropped_queue_full_inc/1,
dropped_queue_full_inc/2,
dropped_expired_inc/1,
dropped_expired_inc/2,
failed_inc/1,
failed_inc/2,
retried_inc/1,
retried_inc/2,
retried_failed_inc/1,
retried_failed_inc/2,
retried_success_inc/1,
retried_success_inc/2,
success_inc/1,
success_inc/2
]).
%% Gauges (value can go both up and down):
%% --------------------------------------
%% @doc Count of requests (batches of messages) that are currently queuing. [Gauge]
queuing_set(Config, Val) ->
telemetry:execute([wolff, queuing],
#{gauge_set => Val},
telemetry_meta_data(Config)).
%% @doc Number of bytes (RAM and/or disk) currently queuing. [Gauge]
queuing_bytes_set(Config, Val) ->
telemetry:execute([wolff, queuing_bytes],
#{gauge_set => Val},
telemetry_meta_data(Config)).
%% @doc Count of messages that were sent asynchronously but ACKs are not
%% received. [Gauge]
inflight_set(Config, Val) ->
telemetry:execute([wolff, inflight],
#{gauge_set => Val},
telemetry_meta_data(Config)).
%% Counters (value can only got up):
%% --------------------------------------
%% @doc Count of messages dropped
dropped_inc(Config) ->
dropped_inc(Config, 1).
dropped_inc(Config, Val) ->
telemetry:execute([wolff, dropped],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%% @doc Count of messages dropped because the queue was full
dropped_queue_full_inc(Config) ->
dropped_queue_full_inc(Config, 1).
dropped_queue_full_inc(Config, Val) ->
telemetry:execute([wolff, dropped_queue_full],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%% @doc Count of messages dropped because they stayed in the buffer longer
%% than `max_batch_age' before they could be (re)sent to Kafka.
dropped_expired_inc(Config) ->
dropped_expired_inc(Config, 1).
dropped_expired_inc(Config, Val) ->
telemetry:execute([wolff, dropped_expired],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%% @doc The number of times message sends have been retried
retried_inc(Config) ->
retried_inc(Config, 1).
retried_inc(Config, Val) ->
telemetry:execute([wolff, retried],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%% @doc Count of message sends that have failed
failed_inc(Config) ->
failed_inc(Config, 1).
failed_inc(Config, Val) ->
telemetry:execute([wolff, failed],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%%% @doc Count of message sends that have failed after having been retried
retried_failed_inc(Config) ->
retried_failed_inc(Config, 1).
retried_failed_inc(Config, Val) ->
telemetry:execute([wolff, retried_failed],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%% @doc Count messages that were sucessfully sent after at least one retry
retried_success_inc(Config) ->
retried_success_inc(Config, 1).
retried_success_inc(Config, Val) ->
telemetry:execute([wolff, retried_success],
#{counter_inc => Val},
telemetry_meta_data(Config)).
%% @doc Count of messages that have been sent successfully
success_inc(Config) ->
success_inc(Config, 1).
success_inc(Config, Val) ->
telemetry:execute([wolff, success],
#{counter_inc => Val},
telemetry_meta_data(Config)).
telemetry_meta_data(Config) ->
maps:get(telemetry_meta_data, Config, #{}).