Packages
brod
3.6.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
Files
Jump to
Current section
Files
src/brod_consumers_sup.erl
%%%
%%% Copyright (c) 2015-2018 Klarna Bank AB (publ)
%%%
%%% Licensed under the Apache License, Version 2.0 (the "License");
%%% you may not use this file except in compliance with the License.
%%% You may obtain a copy of the License at
%%%
%%% http://www.apache.org/licenses/LICENSE-2.0
%%%
%%% Unless required by applicable law or agreed to in writing, software
%%% distributed under the License is distributed on an "AS IS" BASIS,
%%% WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
%%% See the License for the specific language governing permissions and
%%% limitations under the License.
%%%
%%%=============================================================================
%%% @doc brod consumers supervisor
%%%=============================================================================
-module(brod_consumers_sup).
-behaviour(supervisor3).
-export([ init/1
, post_init/1
, start_link/0
, find_consumer/3
, start_consumer/4
, stop_consumer/2
]).
-include("brod_int.hrl").
-define(TOPICS_SUP, brod_consumers_sup).
-define(PARTITIONS_SUP, brod_consumers_sup2).
%% By default, restart ?PARTITIONS_SUP after a 10-seconds delay
-define(DEFAULT_PARTITIONS_SUP_RESTART_DELAY, 10).
%% By default, restart partition consumer worker process after a 2-seconds delay
-define(DEFAULT_CONSUMER_RESTART_DELAY, 2).
%%%_* APIs =====================================================================
%% @doc Start a root consumers supervisor.
-spec start_link() -> {ok, pid()}.
start_link() ->
supervisor3:start_link(?MODULE, ?TOPICS_SUP).
%% @doc Dynamically start a per-topic supervisor.
-spec start_consumer(pid(), pid(), brod:topic(), brod:consumer_config()) ->
{ok, pid()} | {error, any()}.
start_consumer(SupPid, ClientPid, TopicName, Config) ->
Spec = consumers_sup_spec(ClientPid, TopicName, Config),
supervisor3:start_child(SupPid, Spec).
%% @doc Dynamically stop a per-topic supervisor.
-spec stop_consumer(pid(), brod:topic()) -> ok | {error, any()}.
stop_consumer(SupPid, TopicName) ->
supervisor3:terminate_child(SupPid, TopicName).
%% @doc Find a brod_consumer process pid running under ?PARTITIONS_SUP
-spec find_consumer(pid(), brod:topic(), brod:partition()) ->
{ok, pid()} | {error, Reason} when
Reason :: {consumer_not_found, brod:topic()}
| {consumer_not_found, brod:topic(), brod:partition()}
| {consumer_down, noproc}.
find_consumer(SupPid, Topic, Partition) ->
case supervisor3:find_child(SupPid, Topic) of
[] ->
%% no such topic worker started,
%% check sys.config or brod:start_link_client args
{error, {consumer_not_found, Topic}};
[PartitionsSupPid] ->
try
case supervisor3:find_child(PartitionsSupPid, Partition) of
[] ->
%% no such partition?
{error, {consumer_not_found, Topic, Partition}};
[Pid] ->
{ok, Pid}
end
catch exit : {noproc, _} ->
{error, {consumer_down, noproc}}
end
end.
%% @doc supervisor3 callback.
init(?TOPICS_SUP) ->
{ok, {{one_for_one, 0, 1}, []}};
init({?PARTITIONS_SUP, _ClientPid, _Topic, _Config}) ->
post_init.
post_init({?PARTITIONS_SUP, ClientPid, Topic, Config}) ->
%% spawn consumer process for every partition
%% in a topic if partitions are not set explicitly
%% in the config
%% TODO: make it dynamic when consumer groups API is ready
case get_partitions(ClientPid, Topic, Config) of
{ok, Partitions} ->
Children = [ consumer_spec(ClientPid, Topic, Partition, Config)
|| Partition <- Partitions ],
{ok, {{one_for_one, 0, 1}, Children}};
Error ->
Error
end;
post_init(_) ->
ignore.
get_partitions(ClientPid, Topic, Config) ->
case proplists:get_value(partitions, Config, []) of
[] ->
get_all_partitions(ClientPid, Topic);
[_|_] = List ->
{ok, List}
end.
get_all_partitions(ClientPid, Topic) ->
case brod_client:get_partitions_count(ClientPid, Topic) of
{ok, PartitionsCnt} ->
{ok, lists:seq(0, PartitionsCnt - 1)};
{error, _} = Error ->
Error
end.
consumers_sup_spec(ClientPid, TopicName, Config0) ->
DelaySecs = proplists:get_value(topic_restart_delay_seconds, Config0,
?DEFAULT_PARTITIONS_SUP_RESTART_DELAY),
Config = proplists:delete(topic_restart_delay_seconds, Config0),
Args = [?MODULE, {?PARTITIONS_SUP, ClientPid, TopicName, Config}],
{ _Id = TopicName
, _Start = {supervisor3, start_link, Args}
, _Restart = {permanent, DelaySecs}
, _Shutdown = infinity
, _Type = supervisor
, _Module = [?MODULE]
}.
consumer_spec(ClientPid, Topic, Partition, Config0) ->
DelaySecs = proplists:get_value(partition_restart_delay_seconds, Config0,
?DEFAULT_CONSUMER_RESTART_DELAY),
Config = proplists:delete(partition_restart_delay_seconds, Config0),
Args = [ClientPid, Topic, Partition, Config],
{ _Id = Partition
, _Start = {brod_consumer, start_link, Args}
, _Restart = {transient, DelaySecs} %% restart only when not normal exit
, _Shutdown = 5000
, _Type = worker
, _Module = [brod_consumer]
}.
%%%_* Internal Functions =======================================================
%%%_* Emacs ====================================================================
%%% Local Variables:
%%% allout-layout: t
%%% erlang-indent-level: 2
%%% End: