The Broadway front: a pipeline that hands every message its producer
emits to StatifierRouter.route/3.
The host starts it in its own supervision tree, with any Broadway producer it already operates; this package ships no producer and starts no process of its own.
children = [
MyApp.Repo,
{StatifierRouter.Broadway,
name: MyApp.AdEventsRouter,
producer: {BroadwayKafka.Producer, kafka_opts},
router: router_config,
processors: [default: [concurrency: 8]]}
]Options
start_link/1 takes a keyword list:
| Option | Value | Default |
|---|---|---|
:name | the pipeline's name, an atom | required |
:producer | {module, opts}: any Broadway producer and its options | required |
:router | a %StatifierRouter.Config{} | required |
:normalize | a fun of one %Broadway.Message{} returning the event StatifierRouter.route/3 takes | normalize/1 |
:processors | Broadway's :processors option, passed through | [default: []] |
Any other key, or a value the table does not allow, raises
ArgumentError, as Broadway.start_link/2 does for its own options.
There are no batchers in this release.
Each message
handle_message/3 normalizes the message and calls
StatifierRouter.route/3, which returns once every delivery's
transaction has ended. {:ok, outcomes} returns the message unchanged,
so Broadway acknowledges it after the transactions committed.
{:error, reason} returns it failed with reason. A raise inside
route/3 is not rescued here either: Broadway fails the message it
raised on (ADR-0003, section 1).
Redelivery is the producer's contract
A failed message is not redelivered by this pipeline and not by Broadway, which "does not provide any sort of retries out of the box" and acknowledges a failed message as failed immediately. Whether the event comes back is the producer's contract:
- a queue-style producer that leaves an unacknowledged message invisible for a timeout and then hands it over again, Amazon SQS the example Broadway itself names, redelivers it;
BroadwayKafka.Producer, the producer this module's example uses, "will always ack the messages even when they fail" and advances the group's offset past the failed one, so nothing hands it over again and reprocessing is, in its own words, a strategy the host rolls.
A host that needs a failed delivery retried therefore chooses a producer that gives it back, or arranges the replay itself. Nothing in this package retries, and nothing in it holds a failed message.
Partitioning
partition/3 is the pipeline's partition_by. It keeps every message
for one address on one processor, so deliveries to one execution reach
it one after another instead of queueing on its lock, each holding a
pooled connection while it waits. That is an optimisation and never the
guarantee: two deliveries to one execution are stepped one at a time by
statifier_persistence's per-execution lock, whichever processor they
came through (ADR-0003, section 5). A binding whose order is :none
is not partitioned by its key (ADR-0003, section 10).
Summary
Functions
The default :normalize: the event StatifierRouter.route/3 takes, read
from the message. :scope, :message_id and :source are read from the
message's metadata, and the message's data is the event's data.
The partition of one message: the hash, by :erlang.phash2/1, of the
address {scope, document, key} of the first binding in router's
order that is enabled, is for the event's source, has order: :by_key,
and whose match holds and key produces a key for the event.
Starts the pipeline, linked to the calling process. The module documentation lists the options.
Functions
@spec normalize(Broadway.Message.t()) :: map()
The default :normalize: the event StatifierRouter.route/3 takes, read
from the message. :scope, :message_id and :source are read from the
message's metadata, and the message's data is the event's data.
A message that does not carry all four is not refused here: the event
this function builds carries nil for what the metadata did not hold,
and StatifierRouter.route/3 refuses it. A message with no
:message_id in its metadata is refused with
{:error, :no_message_id}, which route/3 checks before anything else;
a message whose metadata holds no :scope or no :source, or whose
data is not a map, is refused with
{:error, {:invalid_event, event}}. Either way the message fails. A
producer whose messages carry them elsewhere is paired with a
:normalize of the host's own.
iex> message = %Broadway.Message{
...> data: %{"kind" => "click", "impression_id" => "imp_7f3a"},
...> metadata: %{scope: "7c1e", message_id: "ad_events/3/1107", source: "ad_events"},
...> acknowledger: Broadway.NoopAcknowledger.init()
...> }
iex> StatifierRouter.Broadway.normalize(message)
%{
scope: "7c1e",
message_id: "ad_events/3/1107",
source: "ad_events",
data: %{"kind" => "click", "impression_id" => "imp_7f3a"}
}
@spec partition( Broadway.Message.t(), StatifierRouter.Config.t(), (Broadway.Message.t() -> map()) ) :: non_neg_integer()
The partition of one message: the hash, by :erlang.phash2/1, of the
address {scope, document, key} of the first binding in router's
order that is enabled, is for the event's source, has order: :by_key,
and whose match holds and key produces a key for the event.
Under a :bindings_resolver, the bindings are the resolver's answer for
the event's scope, asked here once per message and again by route/3 in
handle_message/3 (ADR-0001, the Amendment of 2026-09-25). An answer
that is not a list of bindings raises ArgumentError; here that raise
is rescued and the message is partitioned by its message id, so it does
not take the producer down. route/3 raises the same ArgumentError
in handle_message/3, and Broadway fails the message there. A resolver
that raises ArgumentError itself is treated the same way; any other
raise from the resolver is not rescued here.
A message no such binding addresses is partitioned by the hash of its
message id. So is a message normalize builds no routable event from,
one whose resolver answer StatifierRouter.route/3 would refuse, and
one whose resolver answer is malformed: a partitioner answers for every
message, and handle_message/3 is where such a message fails.
normalize is called here as well as in handle_message/3, and here it
runs in the producer's dispatcher, where a raise takes the producer down
with the messages it holds. It is not rescued: a :normalize must answer
for every message, and the default one does.
@spec start_link(keyword()) :: Broadway.on_start()
Starts the pipeline, linked to the calling process. The module documentation lists the options.