StatifierRouter.Broadway (StatifierRouter v0.9.2)

Copy Markdown View Source

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:

OptionValueDefault
:namethe pipeline's name, an atomrequired
:producer{module, opts}: any Broadway producer and its optionsrequired
:routera %StatifierRouter.Config{}required
:normalizea fun of one %Broadway.Message{} returning the event StatifierRouter.route/3 takesnormalize/1
:processorsBroadway'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

normalize(message)

@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"}
}

partition(message, router, normalize)

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.

start_link(opts)

@spec start_link(keyword()) :: Broadway.on_start()

Starts the pipeline, linked to the calling process. The module documentation lists the options.