Packages
brod
4.3.1
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
4
files changed
+19
additions
-6
deletions
| @@ -24,4 +24,4 @@ | |
| 24 24 | [{<<"app">>,<<"kafka_protocol">>}, |
| 25 25 | {<<"optional">>,false}, |
| 26 26 | {<<"requirement">>,<<"4.1.9">>}]}]}. |
| 27 | - {<<"version">>,<<"4.3.0">>}. |
| 27 | + {<<"version">>,<<"4.3.1">>}. |
| @@ -1,6 +1,6 @@ | |
| 1 1 | {application,brod, |
| 2 2 | [{description,"Apache Kafka Erlang client library"}, |
| 3 | - {vsn,"4.3.0"}, |
| 3 | + {vsn,"4.3.1"}, |
| 4 4 | {registered,[]}, |
| 5 5 | {applications,[kernel,stdlib,kafka_protocol]}, |
| 6 6 | {env,[]}, |
| @@ -388,8 +388,8 @@ handle_call({stop_producer, Topic}, _From, State) -> | |
| 388 388 | ok = brod_producers_sup:stop_producer(State#state.producers_sup, Topic), |
| 389 389 | {reply, ok, State}; |
| 390 390 | handle_call({stop_consumer, Topic}, _From, State) -> |
| 391 | - ok = brod_consumers_sup:stop_consumer(State#state.consumers_sup, Topic), |
| 392 | - {reply, ok, State}; |
| 391 | + Reply = brod_consumers_sup:stop_consumer(State#state.consumers_sup, Topic), |
| 392 | + {reply, Reply, State}; |
| 393 393 | handle_call({get_leader_connection, Topic, Partition}, _From, State) -> |
| 394 394 | {Result, NewState} = do_get_leader_connection(State, Topic, Partition), |
| 395 395 | {reply, Result, NewState}; |
| @@ -395,12 +395,15 @@ handle_info(_Info, State) -> | |
| 395 395 | %%-------------------------------------------------------------------- |
| 396 396 | -spec terminate(Reason :: normal | shutdown | {shutdown, term()} | term(), |
| 397 397 | State :: term()) -> any(). |
| 398 | - terminate(_Reason, #state{workers = Workers, |
| 398 | + terminate(_Reason, #state{config = Config, |
| 399 | + workers = Workers, |
| 399 400 | coordinator = Coordinator, |
| 400 401 | group_id = GroupId |
| 401 402 | }) -> |
| 402 403 | ok = terminate_all_workers(Workers), |
| 403 | - ok = flush_offset_commits(GroupId, Coordinator). |
| 404 | + ok = flush_offset_commits(GroupId, Coordinator), |
| 405 | + ok = stop_consumers(Config), |
| 406 | + ok. |
| 404 407 | |
| 405 408 | %%%=================================================================== |
| 406 409 | %%% Internal functions |
| @@ -529,6 +532,16 @@ do_ack(Topic, Partition, Offset, #state{ workers = Workers | |
| 529 532 | {error, unknown_topic_or_partition} |
| 530 533 | end. |
| 531 534 | |
| 535 | + stop_consumers(Config) -> |
| 536 | + #{ client := Client |
| 537 | + , topics := Topics |
| 538 | + } = Config, |
| 539 | + lists:foreach( |
| 540 | + fun(Topic) -> |
| 541 | + _ = brod_client:stop_consumer(Client, Topic) |
| 542 | + end, |
| 543 | + Topics). |
| 544 | + |
| 532 545 | %%%_* Emacs ==================================================================== |
| 533 546 | %%% Local Variables: |
| 534 547 | %%% allout-layout: t |