Packages

kafka_ex_tc

0.12.1-21
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 protocol sync_group.ex
Raw

lib/kafka_ex/protocol/sync_group.ex

defmodule KafkaEx.Protocol.SyncGroup do
@moduledoc """
Implementation of the Kafka SyncGroup request and response APIs
"""
@member_assignment_version 0
defmodule Request do
@moduledoc false
defstruct member_id: nil,
group_name: nil,
generation_id: nil,
assignments: []
@type t :: %Request{
member_id: binary,
group_name: binary,
generation_id: integer,
assignments: [
{member :: binary, [{topic :: binary, partitions :: [integer]}]}
]
}
end
defmodule Assignment do
@moduledoc false
defstruct topic: nil, partitions: []
@type t :: %Assignment{topic: binary, partitions: [integer]}
end
defmodule Response do
@moduledoc false
defstruct error_code: nil, assignments: []
@type t :: %Response{
error_code: atom | integer,
assignments: [Assignment.t()]
}
end
@spec create_request(integer, binary, Request.t()) :: binary
def create_request(correlation_id, client_id, %Request{} = request) do
KafkaEx.Protocol.create_request(:sync_group, correlation_id, client_id) <>
<<byte_size(request.group_name)::16-signed, request.group_name::binary,
request.generation_id::32-signed,
byte_size(request.member_id)::16-signed, request.member_id::binary,
length(request.assignments)::32-signed,
group_assignment_data(request.assignments, "")::binary>>
end
@spec parse_response(binary) :: Response.t()
def parse_response(
<<_correlation_id::32-signed, error_code::16-signed,
member_assignment_len::32-signed,
member_assignment::size(member_assignment_len)-binary>>
) do
%Response{
error_code: KafkaEx.Protocol.error(error_code),
assignments: parse_member_assignment(member_assignment)
}
end
# Helper functions to create assignment data structure
defp group_assignment_data([], acc), do: acc
defp group_assignment_data([h | t], acc),
do: group_assignment_data(t, acc <> member_assignment_data(h))
defp member_assignment_data({member_id, member_assignment}) do
assignment_bytes_for_member = <<
@member_assignment_version::16-signed,
length(member_assignment)::32-signed,
topic_assignment_data(member_assignment, "")::binary,
# UserData
0::32-signed
>>
<<byte_size(member_id)::16-signed, member_id::binary,
byte_size(assignment_bytes_for_member)::32-signed,
assignment_bytes_for_member::binary>>
end
defp topic_assignment_data([], acc), do: acc
defp topic_assignment_data([h | t], acc),
do: topic_assignment_data(t, acc <> partition_assignment_data(h))
defp partition_assignment_data({topic_name, partition_ids}) do
<<byte_size(topic_name)::16-signed, topic_name::binary,
length(partition_ids)::32-signed,
partition_id_data(partition_ids, "")::binary>>
end
defp partition_id_data([], acc), do: acc
defp partition_id_data([h | t], acc),
do: partition_id_data(t, acc <> <<h::32-signed>>)
# Helper functions to parse assignments
defp parse_member_assignment(<<>>), do: []
defp parse_member_assignment(
<<@member_assignment_version::16-signed, assignments_size::32-signed,
rest::binary>>
) do
parse_assignments(assignments_size, rest, [])
end
defp parse_assignments(0, _rest, assignments), do: assignments
defp parse_assignments(
size,
<<topic_len::16-signed, topic::size(topic_len)-binary,
partition_len::32-signed, rest::binary>>,
assignments
) do
{partitions, rest} = parse_partitions(partition_len, rest, [])
parse_assignments(size - 1, rest, [{topic, partitions} | assignments])
end
defp parse_partitions(0, rest, partitions), do: {partitions, rest}
defp parse_partitions(
size,
<<partition::32-signed, rest::binary>>,
partitions
) do
parse_partitions(size - 1, rest, [partition | partitions])
end
end