Current section

113 Versions

Jump to

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