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
10061008handle_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 ;
10091011handle_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
10621071maybe_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
11681178maybe_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