Packages
kafka_ex_tc
0.12.1-11-debug
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
Current section
Files
lib/kafka_ex/protocol/metadata.ex
defmodule KafkaEx.Protocol.Metadata do
alias KafkaEx.Protocol
import KafkaEx.Protocol.Common
@supported_versions_range {0, 1}
@default_api_version 0
@moduledoc """
Implementation of the Kafka Hearbeat request and response APIs
"""
defmodule Request do
@moduledoc false
defstruct topic: nil
@type t :: %Request{topic: binary}
end
defmodule Broker do
@moduledoc false
alias KafkaEx.Socket
defstruct node_id: -1, host: "", port: 0, socket: nil, is_controller: nil
@type t :: %__MODULE__{
node_id: integer,
host: binary,
port: integer,
socket: KafkaEx.Socket.t(),
is_controller: boolean
}
def connected?(%Broker{} = broker) do
broker.socket != nil && Socket.open?(broker.socket)
end
end
defmodule Response do
@moduledoc false
alias KafkaEx.Protocol.Metadata.Broker
alias KafkaEx.Protocol.Metadata.TopicMetadata
defstruct brokers: [], topic_metadatas: [], controller_id: nil
@type t :: %Response{
brokers: [Broker.t()],
topic_metadatas: [TopicMetadata.t()],
controller_id: integer
}
def broker_for_topic(metadata, brokers, topic, partition) do
case Enum.find(metadata.topic_metadatas, &(topic == &1.topic)) do
nil ->
nil
topic_metadata ->
find_lead_broker(metadata.brokers, topic_metadata, brokers, partition)
end
end
def partitions_for_topic(
metadata,
topic,
partitioner \\ :default_partitioner
) do
case Enum.find(metadata.topic_metadatas, &(&1.topic == topic)) do
nil ->
# topic doesn't exist yet, no partitions
[]
topic_metadata ->
case partitioner do
:default_partitioner ->
Enum.map(topic_metadata.partition_metadatas, & &1.partition_id)
:tt_partitioner ->
leader_available_partition_ids(topic_metadata.partition_metadatas)
end
end
end
defp find_lead_broker(metadata_brokers, topic_metadata, brokers, partition) do
case Enum.find(
topic_metadata.partition_metadatas,
&(partition == &1.partition_id)
) do
nil -> nil
lead_broker -> find_broker(lead_broker, metadata_brokers, brokers)
end
end
defp find_broker(lead_broker, metadata_brokers, brokers) do
case Enum.find(metadata_brokers, &(lead_broker.leader == &1.node_id)) do
nil ->
nil
broker ->
Enum.find(brokers, &broker_for_host?(&1, broker.host, broker.port))
end
end
defp broker_for_host?(broker, host, port) do
broker.host == host && broker.port == port && Broker.connected?(broker)
end
defp leader_available_partition_ids(partition_metadatas) do
Enum.reduce(
partition_metadatas,
[],
fn partition_metadata, acc ->
case partition_metadata.error_code do
:leader_not_available -> acc
_ -> [partition_metadata.partition_id | acc]
end
end
)
end
end
defmodule TopicMetadata do
@moduledoc false
alias KafkaEx.Protocol.Metadata.PartitionMetadata
defstruct error_code: 0,
topic: nil,
is_internal: nil,
partition_metadatas: []
@type t :: %TopicMetadata{
error_code: integer | :no_error,
topic: nil | binary,
is_internal: nil | boolean,
partition_metadatas: [PartitionMetadata.t()]
}
end
defmodule PartitionMetadata do
@moduledoc false
defstruct error_code: 0,
partition_id: nil,
leader: -1,
replicas: [],
isrs: []
@type t :: %PartitionMetadata{
error_code: integer,
partition_id: nil | integer,
leader: integer,
replicas: [integer],
isrs: [integer]
}
end
def api_version(api_versions) do
case KafkaEx.ApiVersions.find_api_version(
api_versions,
:metadata,
@supported_versions_range
) do
{:ok, version} ->
version
# those three should never happen since :metadata is part of the protocol since the beginning.
# they are left here as this will server as reference implementation
# :unknown_message_for_server ->
# :unknown_message_for_client ->
# :no_version_supported ->
_ ->
@default_api_version
end
end
def create_request(
correlation_id,
client_id,
topics,
api_version \\ @default_api_version
)
def create_request(correlation_id, client_id, nil, api_version) do
create_request(correlation_id, client_id, "", api_version)
end
def create_request(correlation_id, client_id, "", api_version) do
topic_count = if 0 == api_version, do: 0, else: -1
KafkaEx.Protocol.create_request(
:metadata,
correlation_id,
client_id,
api_version
) <> <<topic_count::32-signed>>
end
def create_request(correlation_id, client_id, topic, api_version)
when is_binary(topic) do
create_request(correlation_id, client_id, [topic], api_version)
end
def create_request(correlation_id, client_id, topics, api_version)
when is_list(topics) do
KafkaEx.Protocol.create_request(
:metadata,
correlation_id,
client_id,
api_version
) <> <<length(topics)::32-signed, topic_data(topics)::binary>>
end
def parse_response(data), do: parse_response(data, @default_api_version)
def parse_response(data, nil), do: parse_response(data, @default_api_version)
def parse_response(
<<_correlation_id::32-signed, brokers_size::32-signed, rest::binary>>,
api_version
) do
case api_version do
1 ->
{brokers, rest} = parse_brokers(1, brokers_size, rest, [])
<<controller_id::32-signed, rest::binary>> = rest
<<topic_metadatas_size::32-signed, rest::binary>> = rest
%Response{
brokers: brokers,
controller_id: controller_id,
topic_metadatas: parse_topic_metadatas_v1(topic_metadatas_size, rest)
}
0 ->
{brokers, rest} = parse_brokers(0, brokers_size, rest, [])
<<topic_metadatas_size::32-signed, rest::binary>> = rest
%Response{
brokers: brokers,
topic_metadatas: parse_topic_metadatas(topic_metadatas_size, rest)
}
end
end
defp parse_brokers(_api_version, 0, rest, brokers), do: {brokers, rest}
defp parse_brokers(
0,
brokers_size,
<<node_id::32-signed, host_len::16-signed, host::size(host_len)-binary,
port::32-signed, rest::binary>>,
brokers
) do
parse_brokers(0, brokers_size - 1, rest, [
%Broker{node_id: node_id, host: host, port: port} | brokers
])
end
# defp parse_brokers_v1(0, rest, brokers), do: {brokers, rest}
defp parse_brokers(
1,
brokers_size,
<<
node_id::32-signed,
host_len::16-signed,
host::size(host_len)-binary,
port::32-signed,
# rack is nullable
-1::16-signed,
rest::binary
>>,
brokers
) do
parse_brokers(1, brokers_size - 1, rest, [
%Broker{node_id: node_id, host: host, port: port} | brokers
])
end
defp parse_brokers(
1,
brokers_size,
<<
node_id::32-signed,
host_len::16-signed,
host::size(host_len)-binary,
port::32-signed,
rack_len::16-signed,
_rack::size(rack_len)-binary,
rest::binary
>>,
brokers
) do
parse_brokers(1, brokers_size - 1, rest, [
%Broker{node_id: node_id, host: host, port: port} | brokers
])
end
defp parse_topic_metadatas(0, _), do: []
defp parse_topic_metadatas(
topic_metadatas_size,
<<error_code::16-signed, topic_len::16-signed,
topic::size(topic_len)-binary, partition_metadatas_size::32-signed,
rest::binary>>
) do
{partition_metadatas, rest} =
parse_partition_metadatas(partition_metadatas_size, [], rest)
[
%TopicMetadata{
error_code: Protocol.error(error_code),
topic: topic,
partition_metadatas: partition_metadatas
}
| parse_topic_metadatas(topic_metadatas_size - 1, rest)
]
end
defp parse_topic_metadatas_v1(0, _), do: []
defp parse_topic_metadatas_v1(
topic_metadatas_size,
<<
error_code::16-signed,
topic_len::16-signed,
topic::size(topic_len)-binary,
# booleans are actually 8-signed
is_internal::8-signed,
partition_metadatas_size::32-signed,
rest::binary
>>
) do
{partition_metadatas, rest} =
parse_partition_metadatas(partition_metadatas_size, [], rest)
[
%TopicMetadata{
error_code: Protocol.error(error_code),
topic: topic,
partition_metadatas: partition_metadatas,
is_internal: is_internal == 1
}
| parse_topic_metadatas_v1(topic_metadatas_size - 1, rest)
]
end
defp parse_partition_metadatas(0, partition_metadatas, rest),
do: {partition_metadatas, rest}
defp parse_partition_metadatas(
partition_metadatas_size,
partition_metadatas,
<<error_code::16-signed, partition_id::32-signed, leader::32-signed,
rest::binary>>
) do
{replicas, rest} = parse_replicas(rest)
{isrs, rest} = parse_isrs(rest)
parse_partition_metadatas(
partition_metadatas_size - 1,
[
%PartitionMetadata{
error_code: Protocol.error(error_code),
partition_id: partition_id,
leader: leader,
replicas: replicas,
isrs: isrs
}
| partition_metadatas
],
rest
)
end
defp parse_replicas(<<num_replicas::32-signed, rest::binary>>) do
parse_int32_array(num_replicas, rest)
end
defp parse_isrs(<<num_isrs::32-signed, rest::binary>>) do
parse_int32_array([], num_isrs, rest)
end
defp parse_int32_array(array \\ [], num, data)
defp parse_int32_array(array, 0, rest) do
{Enum.reverse(array), rest}
end
defp parse_int32_array(array, num, <<value::32-signed, rest::binary>>) do
parse_int32_array([value | array], num - 1, rest)
end
end