For fine-grained control or when consuming from a single partition, use the low-level partition subscription approach.
Workflow:
- Start a consumer for the topic using
:brod.start_consumer/3. - Subscribe to the specific partition using
:brod.subscribe/5. This returns a consumer_pid. - Handle incoming messages in a
GenServer via handle_info/2. Messages arrive as a kafka_message_set. - Acknowledge processed messages using
:brod.consume_ack(consumer_pid, offset). - Handle errors (like
kafka_fetch_error) in handle_info/2 to manage consumer lifecycle or crashes.
defmodule BrodSample.PartitionSubscriber do
use GenServer
import Record, only: [defrecord: 2, extract: 2]
defrecord :kafka_message, extract(:kafka_message, from_lib: "brod/include/brod.hrl")
defrecord :kafka_message_set, extract(:kafka_message_set, from_lib: "brod/include/brod.hrl")
defrecord :kafka_fetch_error, extract(:kafka_fetch_error, from_lib: "brod/include/brod.hrl")
defmodule State do
@enforce_keys [:consumer_pid]
defstruct consumer_pid: nil
end
defmodule KafkaMessage do
@enforce_keys [:offset, :key, :value, :ts]
defstruct offset: nil, key: nil, value: nil, ts: nil
end
def start_link(topic, partition) do
GenServer.start_link(__MODULE__, {topic, partition})
end
@impl true
def init({topic, partition}) do
:ok = :brod.start_consumer(:kafka_client, topic, begin_offset: :latest)
{:ok, consumer_pid} = :brod.subscribe(:kafka_client, self(), topic, partition, [])
{:ok, %State{consumer_pid: consumer_pid}}
end
@impl true
def handle_info(
{consumer_pid, kafka_message_set(messages: msgs)},
%State{consumer_pid: consumer_pid} = state
) do
for msg <- msgs do
msg = kafka_message_to_struct(msg)
IO.inspect(msg)
:brod.consume_ack(consumer_pid, msg.offset)
end
{:noreply, state}
end
def handle_info({pid, kafka_fetch_error()} = error, %State{consumer_pid: pid} = state) do
{:stop, error, state}
end
defp kafka_message_to_struct(kafka_message(offset: offset, key: key, value: value, ts: ts)) do
%KafkaMessage{
offset: offset,
key: key,
value: value,
ts: DateTime.from_unix!(ts, :millisecond)
}
end
end