Packages
brod
3.7.2
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_member.erl
%%%
%%% Copyright (c) 2016-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
%%% Implement brod_group_member behaviour callbacks to allow a process act like
%%% a group member without having to deal with kafka group protocol details.
%%% A typical work flow:
%%%
%%% 1. Spawn a group coordinator by calling
%%% @see brod_group_coordinator:start_link/6.
%%% 2. Subscribe to partitions received in the assignemts from
%%% @see assignments_received/4. callback.
%%% 3. Receive messages from subscribed partitions (delivered by the partition
%%% workers (the pollers) implemented in brod_consumer);
%%% 4. Unsubscribe from all previously subscribed partitions when
%%% @see assignments_revoked/1. is called.
%%%
%%% For group members who commit offsets to kafka, they should:
%%% 1. Call @see brod_group_coordinator:ack/4. to acknowledge sucessfull
%%% consumption of the messages. Group coordinator will commit the
%%% acknowledged offsets every configured interval.
%%% 2. Call @see brod_group_coordinator:commit_offsets/1,2.
%%% to force an immediate offset commit if necessary.
%%%
%%% For group members who manages offsets locally, they should:
%%% 1. Implement the get_committed_offsets/2 callback.
%%% This callback is evaluated everytime when new assignments are received.
%%% @end
%%%=============================================================================
-module(brod_group_member).
-include("brod_int.hrl").
%% Call the callback module to initialize assignments.
%% NOTE: This function is called only when `offset_commit_policy' is
%% `consumer_managed' in group config.
%% see brod_group_coordinator:start_link/6. for more group config details
%% NOTE: The committed offsets should be the offsets for successfully processed
%% (acknowledged) messages, not the begin-offset to start fetching from.
-callback get_committed_offsets(pid(), [{brod:topic(), brod:partition()}]) ->
{ok, [{{brod:topic(), brod:partition()}, brod:offset()}]}.
%% Called when the member is elected as the consumer group leader.
%% The first element in the group member list is ensured to be the leader.
%% NOTE: this function is called only when 'partition_assignment_strategy' is
%% 'callback_implemented' in group config.
%% see brod_group_coordinator:start_link/6. for more group config details.
-callback assign_partitions(pid(), [brod:group_member()],
[{brod:topic(), brod:partition()}]) ->
[{brod:group_member_id(),
[brod:partition_assignment()]}].
%% Called when assignments are received from group leader.
%% the member process should now call brod:subscribe/5
%% to start receiving message from kafka.
-callback assignments_received(pid(), brod:group_member_id(),
brod:group_generation_id(),
brod:received_assignments()) -> ok.
%% Called before group re-balancing, the member should call
%% brod:unsubscribe/3 to unsubscribe from all currently subscribed partitions.
-callback assignments_revoked(pid()) -> ok.
%%%_* Emacs ====================================================================
%%% Local Variables:
%%% allout-layout: t
%%% erlang-indent-level: 2
%%% End: