Current section
Files
Jump to
Current section
Files
src/franz@topic_subscriber.erl
-module(franz@topic_subscriber).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/franz/topic_subscriber.gleam").
-export([ack/1, new/6, with_config/2, with_commited_offset/3, start/1]).
-export_type([builder/1, ack/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
-opaque builder(FLH) :: {builder,
franz:client(),
binary(),
franz@partitions:partitions(),
list({integer(), integer()}),
franz@message_type:message_type(),
fun((integer(), franz:kafka_message(), FLH) -> ack()),
FLH,
list(franz@consumer_config:consumer_config())}.
-type ack() :: any().
-file("src/franz/topic_subscriber.gleam", 24).
?DOC(" Acknowledge the processing of the message.\n").
-spec ack(any()) -> ack().
ack(Cb_state) ->
franz_ffi:ack(Cb_state).
-file("src/franz/topic_subscriber.gleam", 39).
?DOC(" Create a new topic subscriber builder.\n").
-spec new(
franz:client(),
binary(),
franz@partitions:partitions(),
franz@message_type:message_type(),
fun((integer(), franz:kafka_message(), FLO) -> ack()),
FLO
) -> builder(FLO).
new(Client, Topic, Partitions, Message_type, Callback, Init_callback_state) ->
{builder,
Client,
Topic,
Partitions,
[],
Message_type,
Callback,
Init_callback_state,
[]}.
-file("src/franz/topic_subscriber.gleam", 60).
?DOC(" Add a consumer configuration to the topic subscriber builder.\n").
-spec with_config(builder(FLQ), franz@consumer_config:consumer_config()) -> builder(FLQ).
with_config(Builder, Consumer_config) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder),
[Consumer_config | erlang:element(9, Builder)]}.
-file("src/franz/topic_subscriber.gleam", 72).
?DOC(
" Add a commited offset to the topic subscriber builder.\n"
" CommittedOffsets are the offsets for the messages that have been successfully processed (acknowledged), not the begin-offset to start fetching from\n"
).
-spec with_commited_offset(builder(FLT), integer(), integer()) -> builder(FLT).
with_commited_offset(Builder, Partition, Offset) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
[{Partition, Offset} | erlang:element(5, Builder)],
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder),
erlang:element(9, Builder)}.
-file("src/franz/topic_subscriber.gleam", 84).
?DOC(" Start a new topic subscriber.\n").
-spec start(builder(any())) -> {ok, gleam@erlang@process:pid_()} |
{error, franz:franz_error()}.
start(Builder) ->
franz_ffi:start_topic_subscriber(
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(9, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder)
).