Skip to content

Commit a93f79a

Browse files
committed
fix: improve log message to indicate detailed reason of request drop
1 parent 523bc97 commit a93f79a

2 files changed

Lines changed: 28 additions & 14 deletions

File tree

changelog.md

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,7 @@
1+
* 4.1.4
2+
- Optimize log message `replayq_overflow_dropped_produce_calls`.
3+
Changed to `dropped_produce_requests` with a descriptive `cause` to hint detailed reason: `buffer_size_limit` or `high_system_RAM_usage`.
4+
15
* 4.1.3
26
- Ensure `wolff_client_sup:ensure_absence` and `wolff_producers_sup:ensure_absence` will perform shutdown and cleanup atomically.
37
Previously, if the caller process is killed while waiting for shutdown, a terminated child may leak under the supervisor.

src/wolff_producer.erl

Lines changed: 24 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -101,6 +101,8 @@
101101
-define(EMPTY, empty).
102102
-define(SYNC_REF(Caller, Ref), {Caller, Ref}).
103103
-define(IS_SYNC_REF(Caller, Ref), ((is_reference(Caller) orelse is_pid(Caller)) andalso is_reference(Ref))).
104+
-define(DISCARD_OVERFLOW, "buffer_size_limit").
105+
-define(DISCARD_HIGH_MEM, "high_system_RAM_usage").
104106

105107
-type ack_fun() :: wolff:ack_fun().
106108
-type send_req() :: ?SEND_REQ({pid(), reference()}, [wolff:msg()], ack_fun()).
@@ -1004,7 +1006,7 @@ eval_ack_cb(?ACK_CB({Caller, Ref}, Partition), BaseOffset) when ?IS_SYNC_REF(Cal
10041006
ok.
10051007

10061008
handle_overflow(St, _IsHighMemOverflow, Overflow) when Overflow =< 0 ->
1007-
ok = maybe_log_discard(St, 0),
1009+
ok = maybe_log_discard(St, 0, ?DISCARD_OVERFLOW),
10081010
St;
10091011
handle_overflow(#{replayq := Q,
10101012
pending_acks := PendingAcks,
@@ -1013,10 +1015,17 @@ handle_overflow(#{replayq := Q,
10131015
IsHighMemOverflow,
10141016
Overflow) ->
10151017
BytesMode =
1016-
case IsHighMemOverflow of
1017-
true -> at_least;
1018-
false -> at_most
1019-
end,
1018+
case IsHighMemOverflow of
1019+
true -> at_least;
1020+
false -> at_most
1021+
end,
1022+
OverflowReason =
1023+
case IsHighMemOverflow of
1024+
true ->
1025+
?DISCARD_HIGH_MEM;
1026+
false ->
1027+
?DISCARD_OVERFLOW
1028+
end,
10201029
{NewQ, QAckRef, Items} =
10211030
replayq:pop(Q, #{bytes_limit => {BytesMode, Overflow}, count_limit => 999999999}),
10221031
ok = replayq:ack(NewQ, QAckRef),
@@ -1027,7 +1036,7 @@ handle_overflow(#{replayq := Q,
10271036
wolff_metrics:dropped_inc(Config, NrOfCalls),
10281037
wolff_metrics:queuing_set(Config, replayq:count(NewQ)),
10291038
wolff_metrics:queuing_bytes_set(Config, replayq:bytes(NewQ)),
1030-
ok = maybe_log_discard(St, NrOfCalls),
1039+
ok = maybe_log_discard(St, NrOfCalls, OverflowReason),
10311040
{CbList, NewPendingAcks} = wolff_pendack:drop_backlog(PendingAcks, CallIDs),
10321041
lists:foreach(fun(Cb) -> eval_ack_cb(Cb, ?buffer_overflow_discarded) end, CbList),
10331042
St#{replayq := NewQ, pending_acks := NewPendingAcks}.
@@ -1049,27 +1058,28 @@ reply_error_for_all_reqs(St, Reason) ->
10491058
inc_sent_failed(Config, Count - AlreadyInc, 1).
10501059

10511060
%% use process dictionary for upgrade without restart
1052-
maybe_log_discard(St, Increment) ->
1061+
maybe_log_discard(St, Increment, Reason) ->
10531062
Last = get_overflow_log_state(),
10541063
#{last_cnt := LastCnt, acc_cnt := AccCnt} = Last,
10551064
case LastCnt =:= AccCnt andalso Increment =:= 0 of
10561065
true -> %% no change
10571066
ok;
10581067
false ->
1059-
maybe_log_discard(St, Increment, Last)
1068+
maybe_log_discard(St, Increment, Last, Reason)
10601069
end.
10611070

10621071
maybe_log_discard(#{topic := Topic, partition := Partition},
10631072
Increment,
1064-
#{last_ts := LastTs, last_cnt := LastCnt, acc_cnt := AccCnt}) ->
1073+
#{last_ts := LastTs, last_cnt := LastCnt, acc_cnt := AccCnt},
1074+
Reason) ->
10651075
NowTs = now_ts(),
10661076
NewAccCnt = AccCnt + Increment,
10671077
DiffCnt = NewAccCnt - LastCnt,
10681078
case NowTs - LastTs > ?MIN_DISCARD_LOG_INTERVAL of
10691079
true ->
10701080
log_warn(Topic, Partition,
1071-
"replayq_overflow_dropped_produce_calls",
1072-
#{count => DiffCnt}),
1081+
"dropped_produce_requests",
1082+
#{count => DiffCnt, cause => Reason}),
10731083
put_overflow_log_state(NowTs, NewAccCnt, NewAccCnt);
10741084
false ->
10751085
put_overflow_log_state(LastTs, LastCnt, NewAccCnt)
@@ -1166,12 +1176,12 @@ set_process_label(_ClientId, _Topic, _Partition) ->
11661176
-include_lib("eunit/include/eunit.hrl").
11671177

11681178
maybe_log_discard_test_() ->
1169-
[ {"no-increment", fun() -> maybe_log_discard(undefined, 0) end}
1179+
[ {"no-increment", fun() -> maybe_log_discard(undefined, 0, ?DISCARD_OVERFLOW) end}
11701180
, {"fake-last-old",
11711181
fun() ->
11721182
Ts0 = now_ts() - ?MIN_DISCARD_LOG_INTERVAL - 1,
11731183
ok = put_overflow_log_state(Ts0, 2, 2),
1174-
ok = maybe_log_discard(#{topic => <<"a">>, partition => 0}, 1),
1184+
ok = maybe_log_discard(#{topic => <<"a">>, partition => 0}, 1, ?DISCARD_OVERFLOW),
11751185
St = get_overflow_log_state(),
11761186
?assertMatch(#{last_cnt := 3, acc_cnt := 3}, St),
11771187
?assert(maps:get(last_ts, St) - Ts0 > ?MIN_DISCARD_LOG_INTERVAL)
@@ -1180,7 +1190,7 @@ maybe_log_discard_test_() ->
11801190
fun() ->
11811191
Ts0 = now_ts(),
11821192
ok = put_overflow_log_state(Ts0, 2, 2),
1183-
ok = maybe_log_discard(#{topic => <<"a">>, partition => 0}, 2),
1193+
ok = maybe_log_discard(#{topic => <<"a">>, partition => 0}, 2, ?DISCARD_OVERFLOW),
11841194
St = get_overflow_log_state(),
11851195
?assertMatch(#{last_cnt := 2, acc_cnt := 4}, St),
11861196
?assert(maps:get(last_ts, St) - Ts0 < ?MIN_DISCARD_LOG_INTERVAL)

0 commit comments

Comments
 (0)