forked from kafka4beam/wolff
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathwolff.erl
More file actions
235 lines (202 loc) · 10.3 KB
/
Copy pathwolff.erl
File metadata and controls
235 lines (202 loc) · 10.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
%% Copyright (c) 2018 EMQ Technologies Co., Ltd. All Rights Reserved.
%%
%% 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.
-module(wolff).
-include("wolff.hrl").
%% Supervised client management APIs
-export([ensure_supervised_client/3,
stop_and_delete_supervised_client/1
]).
%% Primitive producer worker management APIs
-export([start_producers/3,
stop_producers/1
]).
%% Supervised producer management APIs
-export([ensure_supervised_producers/3,
ensure_supervised_dynamic_producers/2,
stop_and_delete_supervised_producers/1
]).
%% Messaging APIs
-export([send/3,
send_sync/3,
cast/3
]).
%% Messaging APIs of dynamic producer.
-export([send2/4,
cast2/4,
send_sync2/4,
add_topic/2,
remove_topic/2
]).
-export([check_connectivity/1,
check_connectivity/2,
check_if_topic_exists/2,
check_if_topic_exists/3]).
%% for test
-export([get_producer/2]).
-export_type([client_id/0, host/0, producers/0, msg/0, ack_fun/0, partitioner/0,
name/0, offset_reply/0, topic/0, gname/0]).
-deprecated({check_if_topic_exists, 3}).
-type gname() :: wolff_producers:gname().
-type client_id() :: binary().
-type host() :: kpro:endpoint().
-type topic() :: kpro:topic().
-type partition() :: kpro:partition().
-type name() :: atom() | binary().
-type offset() :: kpro:offset().
%% Kafka offset for successfully produced messages (`-1' when `required_acks'
%% is `none'), or the reason atom when the message is dropped instead:
%% - `buffer_overflow_discarded': pushed out of the buffer by newer messages
%% (replayq overflow, or high memory pressure with `drop_if_highmem').
%% - `partition_lost': the partition is gone (e.g. topic deleted or shrunk).
%% - `message_expired': all messages in the batch are older than `max_batch_age'.
%% - `max_retry_exceeded': dropped after `max_retry' failed retries.
%% - `message_too_large': a single-call batch is too large for the topic.
-type offset_reply() :: offset()
| buffer_overflow_discarded
| partition_lost
| message_expired
| max_retry_exceeded
| message_too_large.
-type producers_cfg() :: wolff_producers:config().
-type producers() :: wolff_producers:producers().
-type partitioner() :: wolff_producers:partitioner().
-type msg() :: #{key := binary(),
value := binary(),
ts => pos_integer(),
headers => [{binary(), binary()}]
}.
-type ack_fun() :: fun((partition(), offset_reply()) -> ok)
| {fun(), [term()]}. %% apply(F, [Partition, Offset | Args])
%% @doc Start supervised client process.
-spec ensure_supervised_client(client_id(), [host()], wolff_client:config()) ->
{ok, pid()} | {error, any()}.
ensure_supervised_client(ClientId, Hosts, Config) ->
wolff_client_sup:ensure_present(ClientId, Hosts, Config).
%% @doc Stop and delete client under supervisor.
-spec stop_and_delete_supervised_client(client_id()) -> ok.
stop_and_delete_supervised_client(ClientId) ->
wolff_client_sup:ensure_absence(ClientId).
%% @doc Start producers with the per-partition workers linked to caller.
-spec start_producers(pid(), topic(), producers_cfg()) -> {ok, producers()} | {error, any()}.
start_producers(Client, Topic, ProducerCfg) when is_pid(Client) ->
wolff_producers:start_linked_producers(Client, Topic, ProducerCfg).
%% @doc Stop linked producers.
-spec stop_producers(#{workers := map(), _ => _}) -> ok.
stop_producers(Producers) ->
wolff_producers:stop_linked(Producers).
%% @doc Ensure supervised producers are started.
-spec ensure_supervised_producers(client_id(), topic(), producers_cfg()) ->
{ok, producers()} | {error, any()}.
ensure_supervised_producers(ClientId, Topic, ProducerCfg) ->
wolff_producers:start_supervised(ClientId, Topic, ProducerCfg).
%% @doc Ensure supervised dynamic-producers are started.
-spec ensure_supervised_dynamic_producers(client_id(), producers_cfg()) ->
{ok, producers()} | {error, any()}.
ensure_supervised_dynamic_producers(ClientId, ProducerCfg) ->
wolff_producers:start_supervised_dynamic(ClientId, ProducerCfg).
%% @doc Ensure supervised producers are stopped then deleted.
-spec stop_and_delete_supervised_producers(wolff_producers:producers()) -> ok.
stop_and_delete_supervised_producers(Producers) ->
wolff_producers:stop_supervised(Producers).
%% @doc Pick a partition producer and send a batch asynchronously.
%% The callback function is evaluated by producer process when ack is received from kafka.
%% In case `required_acks' is configured to `none', the callback is evaluated immediately after send.
%% The partition number and the per-partition worker pid are returned in a tuple to caller,
%% so it may use them to correlate the future `AckFun' evaluation.
%% NOTE: This API is blocked until the batch is enqueued to the producer buffer, otherwise no backpressure.
%% High produce rate may cause excessive ram and disk usage.
%% NOTE: In case producers are configured with `required_acks = none',
%% the second arg for callback function will always be `?UNKNOWN_OFFSET' (`-1').
-spec send(producers(), [msg()], ack_fun()) -> {partition(), pid()}.
send(Producers, Batch, AckFun) ->
{Partition, ProducerPid} = wolff_producers:pick_producer(Producers, Batch),
ok = wolff_producer:send(ProducerPid, Batch, AckFun),
{Partition, ProducerPid}.
%% @doc Topic as argument for dynamic producers, otherwise equivalent to `send/3'.
-spec send2(producers(), topic(), [msg()], ack_fun()) -> {partition(), pid()}.
send2(Producers, Topic, Batch, AckFun) ->
{Partition, ProducerPid} = wolff_producers:pick_producer2(Producers, Topic, Batch),
ok = wolff_producer:send(ProducerPid, Batch, AckFun),
{Partition, ProducerPid}.
%% @doc Cast a batch to a partition producer.
%% Even less backpressure than `send/3'.
%% It does not wait for the batch to be enqueued to the producer buffer.
-spec cast(producers(), [msg()], ack_fun()) -> {partition(), pid()}.
cast(Producers, Batch, AckFun) ->
{Partition, ProducerPid} = wolff_producers:pick_producer(Producers, Batch),
ok = wolff_producer:send(ProducerPid, Batch, AckFun, no_wait_for_queued),
{Partition, ProducerPid}.
%% @doc Topic as argument for dynamic producers, otherwise equivalent to `cast/3'.
-spec cast2(producers(), topic(), [msg()], ack_fun()) -> {partition(), pid()}.
cast2(Producers, Topic, Batch, AckFun) ->
{Partition, ProducerPid} = wolff_producers:pick_producer2(Producers, Topic, Batch),
ok = wolff_producer:send(ProducerPid, Batch, AckFun, no_wait_for_queued),
{Partition, ProducerPid}.
%% @doc Pick a partition producer and send a batch synchronously.
%% Raise error exception in case produce pid is down or when timed out.
%% NOTE: In case producers are configured with `required_acks => none',
%% the returned offset will always be `?UNKNOWN_OFFSET' (`-1').
%% In case the batch is discarded due to buffer overflow, the offset
%% is `buffer_overflow_discarded'.
%% In case a single message is too large (Kafka topic config max.message.bytes)
%% the offset is `message_too_large'.
-spec send_sync(producers(), [msg()], timeout()) -> {partition(), offset_reply()}.
send_sync(Producers, Batch, Timeout) ->
{_Partition, ProducerPid} = wolff_producers:pick_producer(Producers, Batch),
wolff_producer:send_sync(ProducerPid, Batch, Timeout).
%% @doc Topic as argument for dynamic producers, otherwise equivalent to `send_sync/3'.
-spec send_sync2(producers(), topic(), [msg()], timeout()) -> {partition(), offset_reply()}.
send_sync2(Producers, Topic, Batch, Timeout) ->
{_Partition, ProducerPid} = wolff_producers:pick_producer2(Producers, Topic, Batch),
wolff_producer:send_sync(ProducerPid, Batch, Timeout).
%% @doc Add a topic to dynamic producer.
%% Returns `ok' if the topic is already addded.
-spec add_topic(producers(), topic()) -> ok | {error, any()}.
add_topic(Producers, Topic) ->
wolff_producers:add_topic(Producers, Topic).
%% @doc Remove a topic from dynamic producer.
%% Returns `ok' if the topic is already removed.
-spec remove_topic(producers(), topic()) -> ok.
remove_topic(Producers, Topic) ->
wolff_producers:remove_topic(Producers, Topic).
%% @hidden For test only.
get_producer(Producers, Partition) ->
wolff_producers:lookup_producer(Producers, Partition).
%% @doc Check if the client is connected to the cluster.
-spec check_connectivity(client_id()) ->
ok | {error, [{FormatedHostPort :: binary(), any()}]}.
check_connectivity(ClientId) ->
case wolff_client_sup:find_client(ClientId) of
{ok, Pid} -> wolff_client:check_connectivity(Pid);
{error, Error} -> {error, Error}
end.
%% @doc Check if the cluster is reachable.
-spec check_connectivity([host()], wolff_client:config()) ->
ok | {error, [{FormatedHostPort :: binary(), any()}]}.
check_connectivity(Hosts, ConnConfig) ->
wolff_client:check_connectivity(Hosts, ConnConfig).
%% @hidden Deprecated. Check if the cluster is reachable and the topic is created.
-spec check_if_topic_exists([host()], wolff_client:config(), topic()) ->
ok | {error, unknown_topic_or_partition | [#{host := binary(), reason := term()}] | any()}.
check_if_topic_exists(Hosts, ConnConfig, Topic) ->
wolff_client:check_if_topic_exists(Hosts, ConnConfig, Topic).
%% @doc Check if a topic exists using a supervised client or a client porcess.
-spec check_if_topic_exists(client_id() | pid(), topic()) -> ok | {error, unknown_topic_or_partition | any()}.
check_if_topic_exists(ClientId, Topic) when is_binary(ClientId) ->
case wolff_client_sup:find_client(ClientId) of
{ok, Pid} -> check_if_topic_exists(Pid, Topic);
{error, Error} -> {error, Error}
end;
check_if_topic_exists(ClientPid, Topic) when is_pid(ClientPid) ->
wolff_client:check_topic_exists_with_client_pid(ClientPid, Topic).