Packages

kafka_ex_tc

0.12.1-25-dev
0.13.0 0.12.1 0.12.1-51-1 0.12.1-50-3 0.12.1-50-2 0.12.1-50-1 0.12.1-49-2 0.12.1-49-1 0.12.1-49-0 0.12.1-48-2 0.12.1-48-1 0.12.1-47-1 0.12.1-46-2 0.12.1-45-2 0.12.1-45-1 0.12.1-44-1 0.12.1-43-1 0.12.1-42-8 0.12.1-42-6 0.12.1-42-5 0.12.1-42-4 0.12.1-42-1 0.12.1-41-2 0.12.1-41-1 0.12.1-40-9 0.12.1-40-8 0.12.1-40-7 0.12.1-40-6 0.12.1-40-5 0.12.1-40-4 0.12.1-40-3 0.12.1-40-2 0.12.1-40-10 0.12.1-40-1 0.12.1-39-2 0.12.1-39-1 0.12.1-36-3 0.12.1-36-2 0.12.1-36-1 0.12.1-34-b 0.12.1-34-a 0.12.1-33-1 0.12.1-32-a1 0.12.1-29-s-3 0.12.1-29-s-2 0.12.1-29-s-1 0.12.1-28-rumble-1 0.12.1-27-2 0.12.1-27-1 0.12.1-26-rumble-2 0.12.1-26-rumble-1 0.12.1-26-dev-1 0.12.1-26-dev.1 0.12.1-26-3 0.12.1-26-2 0.12.1-26-1 0.12.1-25-s-8 0.12.1-25-s-7 0.12.1-25-s-6 0.12.1-25-s-5 0.12.1-25-s-4 0.12.1-25-s-3 0.12.1-25-s-2 0.12.1-25-s-1 0.12.1-25-dev-22 0.12.1-25-dev 0.12.1-24-2 0.12.1-24-1 0.12.1-22-poc-b-7 0.12.1-20-poc-b-6 0.12.1-20-poc-b-5 0.12.1-20-poc-b-4 0.12.1-20-poc-b-3 0.12.1-19-poc-b-3 0.12.1-19-poc-b-2 0.12.1-19-poc-b-1 0.12.1-19-poc-7 0.12.1-19-poc-6 0.12.1-19-poc-4 0.12.1-19-poc-3 0.12.1-19-poc-2 0.12.1-19-poc-1 0.12.1-19-poc 0.12.1-11-debug-1 0.12.1-11-debug 0.12.1-39 0.12.1-38.a 0.12.1-38 0.12.1-37.c 0.12.1-37.b 0.12.1-37.a 0.12.1-37.1 0.12.1-37 0.12.1-35.b 0.12.1-35.a 0.12.1-34 0.12.1-33 0.12.1-32 0.12.1-31 0.12.1-30 0.12.1-29 0.12.1-28 0.12.1-27 0.12.1-26 0.12.1-25 0.12.1-24 0.12.1-23 0.12.1-22.1 0.12.1-22 0.12.1-21 0.12.1-20 0.12.1-19 0.12.1-18 0.12.1-17 0.12.1-16 0.12.1-15 0.12.1-14 0.12.1-13 0.12.1-12 0.12.1-11 0.12.1-10 0.12.1-9 0.12.1-8 0.12.1-7 0.12.1-6 0.12.1-5 0.12.1-4 0.12.1-3 0.12.1-2 0.12.1-1

Kafka client for Elixir/Erlang.

Current section

Files

Jump to
kafka_ex_tc lib kafka_ex consumer_groups_supervisor.ex
Raw

lib/kafka_ex/consumer_groups_supervisor.ex

defmodule KafkaEx.ConsumerGroupsSupervisor do
@moduledoc false
alias KafkaEx.Config
alias KafkaEx.Utils.RegisterName
import KafkaEx.Utils.Guards
def start_link() do
supervisors =
for config <- Config.consumer_group_configs() do
case config do
{config, n} -> children(config, n)
config -> child(config)
end
end
supervisors = :lists.flatten(supervisors)
supervisor_opts =
Keyword.merge(
[strategy: :one_for_one],
name: __MODULE__
)
{:ok, _} = Supervisor.start_link(supervisors, supervisor_opts)
end
defp children(config, n) when is_pos_integer(n) do
for x <- 1..n, do: child(config, x)
end
defp child(config, n \\ nil) do
import Supervisor.Spec
{consumer_group_name, suffix} =
RegisterName.update_consumer_group_name_and_suffix(config, n)
gen_consumer_impl = Keyword.get(config, :gen_consumer_impl)
topic_names = Keyword.get(config, :topic_names)
callback_fn = Keyword.get(config, :callback_fn)
callback_module = Keyword.get(config, :callback_module)
commit_method = Keyword.get(config, :commit_method)
init_offset_f = Keyword.get(config, :init_offset_f)
extra_consumer_args = %{
callback_fn: callback_fn,
callback_module: callback_module,
commit_method: commit_method,
init_offset_f: init_offset_f
}
consumer_group_opts = [
heartbeat_interval: Config.heartbeat_interval(),
commit_interval: Config.commit_interval(),
extra_consumer_args: extra_consumer_args,
suffix: suffix
]
supervisor(
KafkaEx.ConsumerGroup,
[
gen_consumer_impl,
consumer_group_name,
topic_names,
consumer_group_opts
],
id:
RegisterName.consumer_group_sup_id(
consumer_group_name,
config
)
)
end
end