Packages
brod
2.3.4
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_group_subscriber.erl
%%%
%%% Copyright (c) 2016 Klarna AB
%%%
%%% 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
%%% A group subscriber is a gen_server which subscribes to partition consumers
%%% (poller) and calls the user-defined callback functions for message
%%% processing.
%%%
%%% An overview of what it does behind the scene:
%%% 1. Start a consumer group coordinator to manage the consumer group states,
%%% @see brod_group_coordinator:start_link/4.
%%% 2. Start (if not already started) topic-consumers (pollers) and subscribe
%%% to the partition workers when group assignment is received from the
%% group leader, @see brod:start_consumer/3.
%%% 3. Call CallbackModule:handle_message/4 when messages are received from
%%% the partition consumers.
%%% 4. Send acknowledged offsets to group coordinator which will be committed
%%% to kafka periodically.
%%% @copyright 2016 Klarna AB
%%% @end
%%%=============================================================================
-module(brod_group_subscriber).
-behaviour(gen_server).
-behaviour(brod_group_member).
-export([ ack/4
, commit/1
, start_link/7
, stop/1
]).
%% callbacks for brod_group_coordinator
-export([ get_committed_offsets/2
, assignments_received/4
, assignments_revoked/1
, assign_partitions/3
]).
-export([ code_change/3
, handle_call/3
, handle_cast/2
, handle_info/2
, init/1
, terminate/2
]).
-include("brod_int.hrl").
-type cb_state() :: term().
%% Initialize the callback module s state.
-callback init(group_id(), term()) -> {ok, cb_state()}.
%% Handle a message. Return one of:
%%
%% {ok, NewCallbackState}:
%% The subscriber has received the message for processing async-ly.
%% It should call brod_group_subscriber:ack/4 to acknowledge later.
%%
%% {ok, ack, NewCallbackState}
%% The subscriber has completed processing the message.
%%
%% While this callback function is being evaluated, the fetch-ahead
%% partition-consumers are fetching more messages behind the scene
%% unless prefetch_count is set to 0 in consumer config.
%%
-callback handle_message(topic(), partition(), #kafka_message{}, cb_state()) ->
{ok, cb_state()} | {ok, ack, cb_state()}.
%% This callback is called only when subscriber is to commit offsets locally
%% instead of kafka.
%% Return {ok, Offsets, cb_state()} where Offsets can be [],
%% or only the ones that are found in e.g. local storage or database.
%% For the topic-partitions which have no committed offset found,
%% the consumer will take 'begin_offset' in consumer config as the start point
%% of data stream. If 'begin_offset' is not found in consumer config, the
%% default value -1 (latest) is used.
%
% commented out as it's an optional callback
%-callback get_committed_offsets(group_id(), [{topic(), partition()}],
% cb_state()) ->
% {ok, [{{topic(), partition()}, offset()}], cb_state()}.
%% This function is called only when 'partition_assignment_strategy' is
%% 'callback_implemented' in group config.
%
% commented out as it's an optional callback
%-callback assign_partitions([kafka_group_member()],
% [{topic(), partition()}]
% cb_state()) -> [{kafka_group_member_id(),
% [brod_partition_assignment()]}].
-define(DOWN(Reason), {down, brod_utils:os_time_utc_str(), Reason}).
-record(consumer,
{ topic_partition :: {topic(), partition()}
, consumer_pid :: ?undef %% initial state
| pid() %% normal state
| {down, string(), any()} %% consumer restarting
, consumer_mref :: reference()
, begin_offset :: offset()
, acked_offset :: offset()
}).
-type ack_ref() :: {topic(), partition(), offset()}.
-record(state,
{ client :: client()
, client_mref :: reference()
, groupId :: group_id()
, memberId :: member_id()
, generationId :: integer()
, coordinator :: pid()
, consumers = [] :: [#consumer{}]
, consumer_config :: consumer_config()
, is_blocked = false :: boolean()
, cb_module :: module()
, cb_state :: cb_state()
}).
%% delay 2 seconds retry the failed subscription to partiton consumer process
-define(RESUBSCRIBE_DELAY, 2000).
-define(LO_CMD_SUBSCRIBE_PARTITIONS, '$subscribe_partitions').
%%%_* APIs =====================================================================
%% @doc Start (link) a group subscriber.
%% Client:
%% Client ID (or pid, but not recommended) of the brod client.
%% GroupId:
%% Consumer group ID which should be unique per kafka cluster
%% Topics:
%% Predefined set of topic names to join the group.
%% NOTE: The group leader member will collect topics from all members and
%% assign all collected topic-partitions to members in the group.
%% i.e. members can join with arbitrary set of topics.
%% GroupConfig:
%% For group coordinator, @see brod_group_coordinator:start_link/5
%% ConsumerConfig:
%% For partition consumer, @see brod_consumer:start_link/4
%% CbModule:
%% Callback module which should have the callback functions
%% implemented for message processing.
%% CbInitArg:
%% The term() that is going to be passed to CbModule:init/1 when
%% initializing the subscriger.
%% @end
-spec start_link(client(), group_id(), [topic()],
group_config(), consumer_config(), module(), term()) ->
{ok, pid()} | {error, any()}.
start_link(Client, GroupId, Topics, GroupConfig,
ConsumerConfig, CbModule, CbInitArg) ->
Args = {Client, GroupId, Topics, GroupConfig,
ConsumerConfig, CbModule, CbInitArg},
gen_server:start_link(?MODULE, Args, []).
-spec stop(pid()) -> ok.
stop(Pid) ->
Mref = erlang:monitor(process, Pid),
ok = gen_server:cast(Pid, stop),
receive
{'DOWN', Mref, process, Pid, _Reason} ->
ok
end.
%% @doc Acknowledge an offset.
%% The subscriber may ack a later (greater) offset which will be considered
%% as multi-acking the earlier (smaller) offsets. This also means that
%% disordered acks may overwrite offset commits and lead to unnecessary
%% message re-delivery in case of restart.
%% @end
-spec ack(pid(), topic(), partition(), offset()) -> ok.
ack(Pid, Topic, Partition, Offset) ->
gen_server:cast(Pid, {ack, Topic, Partition, Offset}).
%% @doc Commit all acked offsets. NOTE: This is an async call.
-spec commit(pid()) -> ok.
commit(Pid) ->
gen_server:cast(Pid, commit_offsets).
%%%_* APIs for group coordinator ===============================================
%% @doc Called by group coordinator when there is new assignemnt received.
-spec assignments_received(pid(), member_id(), integer(),
brod_received_assignments()) -> ok.
assignments_received(Pid, MemberId, GenerationId, TopicAssignments) ->
gen_server:cast(Pid, {new_assignments, MemberId,
GenerationId, TopicAssignments}).
%% @doc Called by group coordinator before re-joinning the consumer group.
-spec assignments_revoked(pid()) -> ok.
assignments_revoked(Pid) ->
gen_server:call(Pid, unsubscribe_all_partitions, infinity).
-spec assign_partitions(pid(), [kpro_GroupMemberMetadata()],
[{topic(), partition()}]) ->
[{kafka_group_member_id(),
[brod_partition_assignment()]}].
assign_partitions(Pid, MemberMetadataList, TopicPartitionList) ->
Call = {assign_partitions, MemberMetadataList, TopicPartitionList},
gen_server:call(Pid, Call, infinity).
%% @doc Called by group coordinator when initializing the assignments
%% for subscriber.
%% NOTE: this function is called only when it is DISABLED to commit offsets
%% to kafka.
%% @end
-spec get_committed_offsets(pid(), [{topic(), partition()}]) ->
{ok, [{{topic(), partition()}, offset()}]}.
get_committed_offsets(Pid, TopicPartitions) ->
gen_server:call(Pid, {get_committed_offsets, TopicPartitions}, infinity).
%%%_* gen_server callbacks =====================================================
init({Client, GroupId, Topics, GroupConfig,
ConsumerConfig, CbModule, CbInitArg}) ->
ok = brod_utils:assert_client(Client),
ok = brod_utils:assert_group_id(GroupId),
ok = brod_utils:assert_topics(Topics),
{ok, CbState} = CbModule:init(GroupId, CbInitArg),
{ok, Pid} = brod_group_coordinator:start_link(Client, GroupId, Topics,
GroupConfig, ?MODULE, self()),
State = #state{ client = Client
, client_mref = erlang:monitor(process, Client)
, groupId = GroupId
, coordinator = Pid
, consumer_config = ConsumerConfig
, cb_module = CbModule
, cb_state = CbState
},
{ok, State}.
handle_info({_ConsumerPid,
#kafka_message_set{ topic = Topic
, partition = Partition
, messages = Messages
}}, State) ->
NewState = handle_messages(Topic, Partition, Messages, State),
{noreply, NewState};
handle_info({'DOWN', Mref, process, _Pid, _Reason},
#state{client_mref = Mref} = State) ->
%% restart, my supervisor should restart me
%% brod_client DOWN reason is discarded as it should have logged
%% in its crash log
{stop, client_down, State};
handle_info({'DOWN', _Mref, process, Pid, Reason},
#state{consumers = Consumers} = State) ->
case lists:keyfind(Pid, #consumer.consumer_pid, Consumers) of
#consumer{topic_partition = TP} = Consumer ->
NewConsumer = Consumer#consumer{ consumer_pid = ?DOWN(Reason)
, consumer_mref = ?undef
},
NewConsumers = lists:keyreplace(TP, #consumer.topic_partition,
Consumers, NewConsumer),
NewState = State#state{consumers = NewConsumers},
{noreply, NewState};
false ->
{noreply, State}
end;
handle_info(?LO_CMD_SUBSCRIBE_PARTITIONS, State) ->
NewState =
case State#state.is_blocked of
true ->
State;
false ->
{ok, #state{} = NewState_} = subscribe_partitions(State),
NewState_
end,
_ = send_lo_cmd(?LO_CMD_SUBSCRIBE_PARTITIONS, ?RESUBSCRIBE_DELAY),
{noreply, NewState};
handle_info(Info, State) ->
log(State, info, "discarded message:~p", [Info]),
{noreply, State}.
handle_call({get_committed_offsets, TopicPartitions}, _From,
#state{ groupId = GroupId
, cb_module = CbModule
, cb_state = CbState
} = State) ->
case CbModule:get_committed_offsets(GroupId, TopicPartitions, CbState) of
{ok, Result, NewCbState} ->
NewState = State#state{cb_state = NewCbState},
{reply, Result, NewState};
Unknown ->
erlang:error({bad_return_value,
{CbModule, get_committed_offsets, Unknown}})
end;
handle_call({assign_partitions, MemberMetadataList, TopicPartitions}, _From,
#state{ cb_module = CbModule
, cb_state = CbState
} = State) ->
Result = CbModule:assign_partitions(MemberMetadataList,
TopicPartitions, CbState),
{reply, Result, State};
handle_call(unsubscribe_all_partitions, _From,
#state{ consumers = Consumers
} = State) ->
lists:foreach(
fun(#consumer{ consumer_pid = ConsumerPid
, consumer_mref = ConsumerMref
}) ->
case is_pid(ConsumerPid) of
true ->
_ = brod:unsubscribe(ConsumerPid, self()),
_ = erlang:demonitor(ConsumerMref, [flush]);
false ->
ok
end
end, Consumers),
{reply, ok, State#state{ consumers = []
, is_blocked = true
}};
handle_call(Call, _From, State) ->
{reply, {error, {unknown_call, Call}}, State}.
handle_cast({ack, Topic, Partition, Offset}, State) ->
AckRef = {Topic, Partition, Offset},
NewState = handle_ack(AckRef, State),
{noreply, NewState};
handle_cast(commit_offsets, State) ->
ok = brod_group_coordinator:commit_offsets(State#state.coordinator),
{noreply, State};
handle_cast({new_assignments, MemberId, GenerationId, Assignments},
#state{ client = Client
, consumer_config = ConsumerConfig
} = State) ->
AllTopics =
lists:map(fun(#brod_received_assignment{topic = Topic}) ->
Topic
end, Assignments),
lists:foreach(
fun(Topic) ->
ok = brod:start_consumer(Client, Topic, ConsumerConfig)
end, lists:usort(AllTopics)),
Consumers =
[ #consumer{ topic_partition = {Topic, Partition}
, consumer_pid = ?undef
, begin_offset = BeginOffset
, acked_offset = ?undef
}
|| #brod_received_assignment{ topic = Topic
, partition = Partition
, begin_offset = BeginOffset
} <- Assignments
],
_ = send_lo_cmd(?LO_CMD_SUBSCRIBE_PARTITIONS),
NewState = State#state{ consumers = Consumers
, is_blocked = false
, memberId = MemberId
, generationId = GenerationId
},
{noreply, NewState};
handle_cast(stop, State) ->
{stop, normal, State};
handle_cast(_Cast, State) ->
{noreply, State}.
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
terminate(_Reason, #state{}) ->
ok.
%%%_* Internal Functions =======================================================
handle_messages(_Topic, _Partition, [], State) ->
State;
handle_messages(Topic, Partition, [Msg | Rest], State) ->
#kafka_message{offset = Offset} = Msg,
#state{cb_module = CbModule, cb_state = CbState} = State,
AckRef = {Topic, Partition, Offset},
{AckNow, NewCbState} =
case CbModule:handle_message(Topic, Partition, Msg, CbState) of
{ok, NewCbState_} ->
{false, NewCbState_};
{ok, ack, NewCbState_} ->
{true, NewCbState_};
Unknown ->
erlang:error({bad_return_value,
{CbModule, handle_message, Unknown}})
end,
State1 = State#state{cb_state = NewCbState},
NewState =
case AckNow of
true -> handle_ack(AckRef, State1);
false -> State1
end,
handle_messages(Topic, Partition, Rest, NewState).
-spec handle_ack(ack_ref(), #state{}) -> #state{}.
handle_ack(AckRef, #state{ generationId = GenerationId
, consumers = Consumers
, coordinator = Coordinator
} = State) ->
{Topic, Partition, Offset} = AckRef,
case lists:keyfind({Topic, Partition},
#consumer.topic_partition, Consumers) of
#consumer{consumer_pid = ConsumerPid} = Consumer ->
ok = consume_ack(ConsumerPid, Offset),
ok = commit_ack(Coordinator, GenerationId, Topic, Partition, Offset),
NewConsumer = Consumer#consumer{acked_offset = Offset},
NewConsumers = lists:keyreplace({Topic, Partition},
#consumer.topic_partition,
Consumers, NewConsumer),
State#state{consumers = NewConsumers};
false ->
%% stale ack, ignore.
State
end.
%% @private Tell consumer process to fetch more (if pre-fetch count allows).
consume_ack(Pid, Offset) when is_pid(Pid) ->
ok = brod:consume_ack(Pid, Offset);
consume_ack(_Down, _Offset) ->
%% consumer is down, should be restarted by its supervisor
ok.
%% @private Send an async message to group coordinator for offset commit.
commit_ack(Pid, GenerationId, Topic, Partition, Offset) ->
ok = brod_group_coordinator:ack(Pid, GenerationId, Topic, Partition, Offset).
send_lo_cmd(CMD) -> send_lo_cmd(CMD, 0).
send_lo_cmd(CMD, 0) -> self() ! CMD;
send_lo_cmd(CMD, DelayMS) -> erlang:send_after(DelayMS, self(), CMD).
subscribe_partitions(#state{ client = Client
, consumers = Consumers0
} = State) ->
Consumers =
lists:map(fun(C) -> subscribe_partition(Client, C) end, Consumers0),
{ok, State#state{consumers = Consumers}}.
subscribe_partition(Client, Consumer) ->
#consumer{ topic_partition = {Topic, Partition}
, consumer_pid = Pid
, begin_offset = BeginOffset0
, acked_offset = AckedOffset
} = Consumer,
case brod_utils:is_pid_alive(Pid) of
true ->
Consumer;
false ->
%% fetch from the last committed offset + 1
%% otherwise fetch from the begin offset
BeginOffset = case AckedOffset of
?undef -> BeginOffset0;
N when N >= 0 -> N + 1
end,
Options =
case BeginOffset =:= ?undef of
true -> []; %% fetch from 'begin_offset' in consumer config
false -> [{begin_offset, BeginOffset}]
end,
case brod:subscribe(Client, self(), Topic, Partition, Options) of
{ok, ConsumerPid} ->
Mref = erlang:monitor(process, ConsumerPid),
Consumer#consumer{ consumer_pid = ConsumerPid
, consumer_mref = Mref
};
{error, Reason} ->
Consumer#consumer{ consumer_pid = ?DOWN(Reason)
, consumer_mref = ?undef
}
end
end.
log(#state{ groupId = GroupId
, memberId = MemberId
, generationId = GenerationId
}, Level, Fmt, Args) ->
brod_utils:log(
Level,
"group subscriber (groupId=~s,memberId=~s,generation=~p,pid=~p):\n" ++ Fmt,
[GroupId, MemberId, GenerationId, self() | Args]).
%%%_* Emacs ====================================================================
%%% Local Variables:
%%% allout-layout: t
%%% erlang-indent-level: 2
%%% End: