Packages
brod
2.2.15
4.5.7
4.5.6
4.5.5
4.5.4
4.5.3
4.5.2
4.5.1
4.5.0
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.3
4.3.2
4.3.1
4.3.0
4.2.0
4.1.1
4.1.0
4.0.0
3.19.1
3.19.0
3.18.0
3.17.1
3.17.0
3.16.5
3.16.4
3.16.3
3.16.2
3.16.1
3.16.0
3.15.6
3.15.5
3.15.4
3.15.3
3.15.1
3.15.0
3.14.0
3.13.0
3.12.0
3.11.0
3.10.0
3.9.5
3.9.3
3.9.2
3.9.1
3.9.0
3.8.1
3.8.0
3.7.11
3.7.10
3.7.9
3.7.8
3.7.7
3.7.6
3.7.5
3.7.4
3.7.3
3.7.2
3.7.1
3.7.0
3.6.2
3.6.1
3.6.0
3.5.2
3.5.1
3.5.0
3.4.0
3.3.5
3.3.4
3.3.3
3.3.2
3.3.1
3.3.0
3.2.0
3.0.0
2.5.0
2.4.1
2.4.0
2.3.7
2.3.6
2.3.5
2.3.4
2.3.3
2.3.1
2.2.16
2.2.15
2.2.14
2.2.12
2.2.11
2.2.10
2.2.9
2.2.8
2.2.7
2.2.6
2.2.5
2.2.4
2.2.3
2.2.2
2.2.1
2.2.0
2.1.12
2.1.11
2.1.10
2.1.8
2.1.7
2.1.4
2.1.2
2.0.0
Apache Kafka Erlang client library
Current section
113 Versions
Jump to
Current section
113 Versions
Compare versions
5
files changed
+10
additions
-7
deletions
| @@ -1,6 +1,6 @@ | |
| 1 1 | PROJECT = brod |
| 2 2 | PROJECT_DESCRIPTION = Kafka client library in Erlang |
| 3 | - PROJECT_VERSION = 2.2.14 |
| 3 | + PROJECT_VERSION = 2.2.15 |
| 4 4 | |
| 5 5 | DEPS = supervisor3 kafka_protocol |
| @@ -1,5 +1,5 @@ | |
| 1 1 | {<<"name">>,<<"brod">>}. |
| 2 | - {<<"version">>,<<"2.2.14">>}. |
| 2 | + {<<"version">>,<<"2.2.15">>}. |
| 3 3 | {<<"requirements">>, |
| 4 4 | #{<<"kafka_protocol">> => #{<<"app">> => <<"kafka_protocol">>, |
| 5 5 | <<"optional">> => false, |
| @@ -1,6 +1,6 @@ | |
| 1 1 | {application,brod, |
| 2 2 | [{description,"Apache Kafka Erlang client library"}, |
| 3 | - {vsn,"2.2.14"}, |
| 3 | + {vsn,"2.2.15"}, |
| 4 4 | {registered,[]}, |
| 5 5 | {applications,[kernel,stdlib,ssl,kafka_protocol,supervisor3]}, |
| 6 6 | {env,[]}, |
| @@ -392,7 +392,7 @@ handle_add_offset(#pending_acks{ offsets_queue = Queue | |
| 392 392 | _ -> |
| 393 393 | handle_add_offset(PendingAcks, Offsets) |
| 394 394 | end; |
| 395 | - handle_add_offset(PendingAcks, []) -> |
| 395 | + handle_add_offset(#pending_acks{} = PendingAcks, []) -> |
| 396 396 | PendingAcks. |
| 397 397 | |
| 398 398 | maybe_shrink_max_bytes(#state{ prefetch_count = PrefetchCount |
| @@ -48,7 +48,7 @@ | |
| 48 48 | ]). |
| 49 49 | |
| 50 50 | -define(DEFAULT_CONNECT_TIMEOUT, timer:seconds(5)). |
| 51 | - -define(DEFAULT_REQUEST_TIMEOUT, timer:seconds(120)). |
| 51 | + -define(DEFAULT_REQUEST_TIMEOUT, timer:minutes(4)). |
| 52 52 | |
| 53 53 | %%%_* Includes ================================================================= |
| 54 54 | -include("brod_int.hrl"). |
| @@ -303,7 +303,7 @@ handle_msg({From, stop}, #state{mod = Mod, sock = Sock}, _Debug) -> | |
| 303 303 | Mod:close(Sock), |
| 304 304 | _ = reply(From, ok), |
| 305 305 | ok; |
| 306 | - handle_msg(Msg, State, Debug) -> |
| 306 | + handle_msg(Msg, #state{} = State, Debug) -> |
| 307 307 | error_logger:warning_msg("[~p] ~p got unrecognized message: ~p", |
| 308 308 | [?MODULE, self(), Msg]), |
| 309 309 | ?MODULE:loop(State, Debug). |
| @@ -387,7 +387,10 @@ assert_max_req_age(Requests, Timeout) -> | |
| 387 387 | %% @end |
| 388 388 | -spec send_assert_max_req_age(pid(), timeout()) -> ok. |
| 389 389 | send_assert_max_req_age(Pid, Timeout) when Timeout >= 1000 -> |
| 390 | - _ = erlang:send_after(Timeout div 2, Pid, assert_max_req_age), |
| 390 | + %% Check every 1 minute |
| 391 | + %% or every half of the timeout value if it's less than 2 minute |
| 392 | + SendAfter = erlang:min(Timeout div 2, timer:minutes(1)), |
| 393 | + _ = erlang:send_after(SendAfter, Pid, assert_max_req_age), |
| 391 394 | ok. |
| 392 395 | |
| 393 396 | %%%_* Emacs ==================================================================== |