Packages
brod
3.3.5
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_kafka_request.erl
%%%
%%% Copyright (c) 2017 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 Help functions to build request messages.
-module(brod_kafka_request).
-export([ fetch_request/7
, metadata_request/2
, offsets_request/4
, produce_request/7
, offset_fetch_request/3
]).
-include("brod_int.hrl").
-type api() :: brod_kafka_apis:api().
-type vsn() :: brod_kafka_apis:vsn().
-type topic() :: brod:topic().
-type partition() :: brod:partition().
-type offset() :: brod:offset().
%% @doc Make a produce request, If the first arg is a `brod_sock' pid, call
%% `brod_kafka_apis:pick_version/2' to resolve version.
%%
%% NOTE: `pick_version' is essentially a ets lookup, for intensive callers
%% like `brod_producer', we should pick version before hand
%% and re-use it for each produce request.
%% @end
-spec produce_request(pid() | vsn(), topic(), partition(),
brod:kv_list(), integer(), integer(),
brod:compression()) -> kpro:req().
produce_request(MaybePid, Topic, Partition, KvList,
RequiredAcks, AckTimeout, Compression) ->
Vsn = pick_version(produce_request, MaybePid),
kpro:produce_request(Vsn, Topic, Partition, KvList,
RequiredAcks, AckTimeout, Compression).
%% @doc Make a fetch request, If the first arg is a `brod_sock' pid, call
%% `brod_kafka_apis:pick_version/2' to resolve version.
%%
%% NOTE: `pick_version' is essentially a ets lookup, for intensive callers
%% like `brod_producer', we should pick version beforehand
%% and re-use it for each produce request.
%% @end
-spec fetch_request(pid(), topic(), partition(), offset(),
kpro:wait(), kpro:count(), kpro:count()) -> kpro:req().
fetch_request(Pid, Topic, Partition, Offset,
WaitTime, MinBytes, MaxBytes) ->
Vsn = pick_version(fetch_request, Pid),
kpro:fetch_request(Vsn, Topic, Partition, Offset,
WaitTime, MinBytes, MaxBytes).
%% @doc Make a 'offsets_request' message for offset resolution.
%% In kafka protocol, -2 and -1 are semantic 'time' to request for
%% 'earliest' and 'latest' offsets.
%% In brod implementation, -2, -1, 'earliest' and 'latest'
%% are semantic 'offset', this is why often a variable named
%% Offset is used as the Time argument.
%% @end
-spec offsets_request(pid(), topic(), partition(), brod:offset_time()) ->
kpro:req().
offsets_request(SockPid, Topic, Partition, TimeOrSemanticOffset) ->
Time = ensure_integer_offset_time(TimeOrSemanticOffset),
Vsn = pick_version(offsets_request, SockPid),
kpro:offsets_request(Vsn, Topic, Partition, Time).
%% @doc Make a metadata request.
-spec metadata_request(pid(), [topic()]) -> kpro:req().
metadata_request(SockPid, Topics) ->
Vsn = pick_version(metadata_request, SockPid),
TopicsForEncoding =
case Vsn of
0 -> Topics;
_ when Topics =:= [] -> ?kpro_null;
_ -> Topics
end,
kpro:req(metadata_request, Vsn, [{topics, TopicsForEncoding}]).
%% @doc Make a offset fetch request.
%% NOTE: empty topics list only works for kafka 0.10.2.0 or later
%% @end
-spec offset_fetch_request(pid(), brod:group_id(), Topics) -> kpro:req()
when Topics :: [{topic(), [partition()]}].
offset_fetch_request(SockPid, GroupId, Topics0) ->
Topics =
lists:map(
fun({Topic, Partitions}) ->
[ {topic, Topic}
, {partitions, [[{partition, P}] || P <- Partitions]}
]
end, Topics0),
Body = [ {group_id, GroupId}
, {topics, case Topics of
[] -> ?kpro_null;
_ -> Topics
end}
],
Vsn = pick_version(offset_fetch_request, SockPid),
kpro:req(offset_fetch_request, Vsn, Body).
%% @private
-spec pick_version(api(), pid()) -> vsn().
pick_version(_API, Vsn) when is_integer(Vsn) -> Vsn;
pick_version(API, SockPid) when is_pid(SockPid) ->
brod_kafka_apis:pick_version(SockPid, API);
pick_version(API, _) ->
brod_kafka_apis:default_version(API).
%% @private
-spec ensure_integer_offset_time(brod:offset_time()) -> integer().
ensure_integer_offset_time(?OFFSET_EARLIEST) -> -2;
ensure_integer_offset_time(?OFFSET_LATEST) -> -1;
ensure_integer_offset_time(T) when is_integer(T) -> T.
%%%_* Emacs ====================================================================
%%% Local Variables:
%%% allout-layout: t
%%% erlang-indent-level: 2
%%% End: