Macro for creating Pulsar consumer callback modules that support internal state.
This module provides a use macro that sets up your module as a consumer callback
with default implementations for optional callbacks.
Usage
Use this module in your consumer callback:
defmodule MyApp.MessageCounter do
use Pulsar.Consumer.Callback
# Implement required callbacks
def handle_message(message, state) do
# Process message
{:ok, state}
end
endWhen you use Pulsar.Consumer.Callback, your module automatically:
- Gets the
@behaviour Pulsar.Consumer.Callbackannotation - Receives default implementations for all optional callbacks
- Can override any callback by simply defining it
Required Callbacks
handle_message/2- Handle incoming messages with current state
Optional Callbacks
All optional callbacks have sensible defaults that you can override:
init/2- Initialize the callback module state, given the consumer's resolved topic and subscription (default:{:ok, nil})terminate/2- Cleanup when consumer terminates (default::ok)handle_call/3- Handle synchronous calls (default:{:reply, {:error, :not_implemented}, state})handle_cast/2- Handle asynchronous casts (default:{:noreply, state})handle_info/2- Handle other messages (default:{:noreply, state})became_active/1- Called when this consumer becomes the active consumer of a:failoversubscription (default:{:noreply, state})became_passive/1- Called when this consumer becomes a passive (standby) consumer of a:failoversubscription (default:{:noreply, state})handle_invalid_message/2- Called instead ofhandle_message/2for a message whose bytes could not be trusted (default: log a warning and acknowledge it, so it is not redelivered). Override it to record or divert such messages; seePulsar.Message.valid?/1for what they contain. Because they are routed here,handle_message/2only ever receives messages that arrived intact.
Message Format
handle_message/2 receives a Pulsar.Message struct containing all message information.
See Pulsar.Message for detailed field documentation.
You can pattern match on the struct:
def handle_message(%Pulsar.Message{payload: payload}, state) do
# Process the payload
{:ok, state}
endOr access fields directly:
def handle_message(%Pulsar.Message{} = message, state) do
payload = message.payload
producer = Pulsar.Message.producer_name(message)
{:ok, state}
endExample Implementation
defmodule MyApp.MessageCounter do
use Pulsar.Consumer.Callback
def init(opts, _context) do
max_messages = Keyword.get(opts, :max_messages, 1000)
{:ok, %{count: 0, max_messages: max_messages, messages: []}}
end
def handle_message(%Pulsar.Message{payload: payload}, callback_state) do
new_state = %{
callback_state
| count: callback_state.count + 1,
messages: [payload | callback_state.messages]
}
# Stop processing if we've reached the limit
if new_state.count >= new_state.max_messages do
{:stop, new_state}
else
{:ok, new_state}
end
end
def terminate(_reason, callback_state) do
IO.puts("Processed #{callback_state.count} messages")
:ok
end
endMultiple Consumers
For shared or key-shared subscriptions, :consumer_count runs several consumer processes
on one subscription to increase throughput:
{Pulsar.Client,
host: "pulsar://localhost:6650",
consumers: [
[topic: topic,
subscription_name: subscription,
callback_module: MyCallback,
subscription_type: :key_shared,
consumer_count: 3,
name: :orders]
]}Each process has independent callback state. The consumer facade intentionally hides those short-lived worker processes, so state that must be shared, queried, or aggregated across consumers should live in an application process with its own public API.
Return Values
init/2 (Optional)
If not implemented, defaults to {:ok, nil}.
{:ok, state}- Successful initialization with initial state{:error, reason}- Initialization failed
Its second argument is the consumer's resolved context; see ## Consumer Context below.
handle_message/2
{:ok, new_state}- Message processed successfully, acknowledge message automatically{:error, reason, new_state}- Processing failed, track for redelivery{:noreply, new_state}- Message processed, but don't automatically ACK/NACK. UsePulsar.Consumer.ack/2orPulsar.Consumer.nack/2for manual acknowledgment{:stop, new_state}- Message processed successfully, acknowledge, then stop consumer
terminate/2 (Optional)
If not implemented, defaults to :ok.
:ok- Cleanup completed successfully- Any other value is ignored
Consumer Context
A consumer on a partitioned topic subscribes to one partition of it, so several callback
processes share a configured topic while each handles a different partition. init/2
receives which one this process ended up on:
%{
topic: "persistent://public/default/orders-partition-2",
base_topic: "persistent://public/default/orders",
partition: 2,
subscription_name: "order-service",
subscription_type: :shared,
consumer_name: "orders-order-service-partition-2-1"
}:topic is what this consumer subscribed to and :base_topic what it was configured with;
they are equal, and :partition is nil, when the topic is not partitioned. Use :topic
as a per-partition key for metrics or flow accounting, and :base_topic where business
logic cares about the logical topic.
A consumer configured with an already concrete partition topic such as
"orders-partition-2" reports partition: nil, because that is what the broker reports:
metadata for a concrete partition gives a partition count of zero, so the consumer is not a
partitioned one. The index is not recovered from the name, since a topic legitimately named
"events-partition-3" would otherwise be given a partition it does not have. Such a consumer
still gets its topic in :topic; only the index is unavailable.
Keep whatever you need in your state rather than looking it up per message:
defmodule MyApp.PerPartitionMetrics do
use Pulsar.Consumer.Callback
def init(_init_args, context) do
{:ok, %{topic: context.topic, partition: context.partition, count: 0}}
end
def handle_message(%Pulsar.Message{}, state) do
:telemetry.execute([:my_app, :message], %{count: 1}, %{topic: state.topic})
{:ok, %{state | count: state.count + 1}}
end
endFailover Consumer Events
For :failover subscriptions, became_active/1 and became_passive/1
notify a consumer when the broker promotes it to (or demotes it from) the
active role. Useful for metrics/alerting or warming up state before a
takeover; not meaningful for other subscription types. A passive consumer
never receives messages, so there is nothing to pause or resume.
Manual Acknowledgment
When you return {:noreply, state} from handle_message/2, the message will NOT be automatically
acknowledged or negatively acknowledged. This gives you full control over when and how to ACK/NACK messages,
which is useful for:
- Broadway pipelines that handle acknowledgment in batches
- Async processing where acknowledgment happens after the callback returns
- Custom acknowledgment logic based on downstream processing
Example with manual ACK:
def handle_message(%Pulsar.Message{payload: payload, message_id: message_id}, state) do
consumer = self()
# Send to async processor
Task.start(fn ->
case process_async(payload) do
:ok -> Pulsar.Consumer.ack(consumer, message_id)
{:error, _} -> Pulsar.Consumer.nack(consumer, message_id)
end
end)
{:noreply, state}
end
Summary
Types
@type context() :: %{ topic: String.t(), base_topic: String.t(), partition: non_neg_integer() | nil, subscription_name: String.t(), subscription_type: atom(), consumer_name: String.t() | nil }
The consumer's resolved identity, passed to init/2.
:topic- the topic this consumer subscribed to, which for a partitioned topic is one concrete partition:base_topic- the topic the consumer was configured with; equals:topicunless the topic is partitioned:partition- the partition index, ornilwhen the topic is not partitioned:subscription_name- the subscription this consumer belongs to:subscription_type- how that subscription is shared:consumer_name- this worker's broker-visible name, which is its group's name suffixed with the worker's position in it, not the:namethe consumer was configured with
@type init_arg() :: term()
@type message_args() :: Pulsar.Message.t()
@type reason() :: term()
@type state() :: term()
Callbacks
@callback handle_call(term(), GenServer.from(), state()) :: {:reply, term(), state()} | {:reply, term(), state(), timeout() | :hibernate | {:continue, term()}} | {:noreply, state()} | {:noreply, state(), timeout() | :hibernate | {:continue, term()}} | {:stop, term(), term(), state()} | {:stop, term(), state()}