Broadway producer that subscribes to a Commanded event store and yields
events as %Scriba.Event{}-wrapped Broadway.Messages.
Usage
source: {Scriba.Source.Commanded, application: MyApp.CommandedApp}Options
:application(required) — the user'sCommanded.Applicationmodule.:subscription_name(default"scriba") — name passed to Commanded for durable subscription tracking.:start_from(default:origin) — where to start reading::origin,:current, or a specific event number.:buffer_size,:concurrency_limit,:partition_by— forwarded verbatim to the event store adapter's subscription. Unset means the adapter's own default applies. See "Subscription buffer and throughput".
Subscription buffer and throughput
:buffer_size is how many events the store will send before it requires an
acknowledgement. EventStore's default is 1, and that default, not
:parallelism, is what bounds a projection's catch-up rate: Scriba
acknowledges after the batch commits, so a batcher waiting on a single
in-flight event waits out its full :batch_timeout before acking and
releasing the next one.
Measured against a real EventStore, 5,000 events over 100 streams:
:buffer_size | Throughput | 10M events |
|---|---|---|
| unset (adapter default, 1) | 9.1 events/sec | 12.7 days |
| 500 | 6,002 events/sec | 28 minutes |
Short runs read lower — 500 events reach roughly 2,200/sec, because startup is a larger share of the measurement.
Scriba sets no default of its own — the adapter's applies unless configured.
Raising it trades memory and redelivered-work-after-a-crash for throughput:
up to :buffer_size events are held in flight, and an unclean restart
replays whatever had not been acknowledged.
source: {Scriba.Source.Commanded,
application: MyApp.CommandedApp,
buffer_size: 500}Acknowledgement watermark
Two properties of this producer make a larger buffer safe.
Acknowledgement is issued by this process, not by the Broadway batch
processor that committed the batch. An event store may resolve the acking
subscriber from self() and silently ignore an ack from anywhere else,
which stalls the subscription for good once its buffer fills.
And only the longest gapless run of committed events is acknowledged.
Acks are prefix acks — acking event 7 acks everything up to 7 — while
batches commit out of source order whenever :parallelism exceeds 1. A
handler still working on event 5 must therefore hold the watermark at 4,
however many later events have committed; otherwise a crash in that window
loses event 5 with no dead letter, no cursor anomaly and no log line,
because the store believes it was delivered and nothing redelivers it.
Optional dependency (§11)
This module is the only place in Scriba that depends on :commanded. To
honor optional: true, every Commanded reference is dynamic:
- Module atoms are built from string literals (
:"Elixir.Commanded.X"), so the compiler does not record them as static module dependencies. - Calls go through
apply/3.
Result: Scriba (and this module) compile cleanly even when :commanded is
absent. start_link/1 raises a clear error in that case; child_spec/1 is
always safe.
Pause/resume memory caveat (v0.1)
pause/1 sets a paused: true flag — handle_demand/2 returns no
messages while paused, accumulating it in the state's :demand. The
Commanded subscription, however, keeps pushing events into the
source's pending :queue regardless of pause state (we don't
unsubscribe).
Memory grows during pause, bounded by how many events the upstream EventStore delivers in the pause window. For operator-driven pauses (seconds to minutes on low-volume projections), this is fine. For long pauses on high-throughput projections, this is a documented sharp edge.
The cleaner alternative — unsubscribe on pause, re-subscribe on resume from the current cursor — is v0.5 hardening territory and changes EventStore subscription state in non-trivial ways. Out of scope for v0.1.