Watermill made event-driven Go feel almost as ergonomic as an HTTP router. That comfort is a gift. It is also how I talked myself into shipping consistency bugs that did not show up until a quiet Sunday night, when a DNS blip and a partial deploy collided.
I already wrote about getting started elsewhere. This essay is about the physics that remain after the boilerplate disappears. Dual writes. Ordering surprises inside handler groups. Duplicate delivery that everyone "knows about" until money is involved. Middleware that looks complete until one poison payload pins a consumer forever.
The personal stake is simple. Event-driven architecture does not fail loudly the way a 500 does. It fails by drifting: the database believes something happened, and the rest of the world never heard.
The outage that is not a bug
A pod eviction. A rebalance. A process stopped between two steps that both needed to happen. The program did not "have a bug" in the usual sense. The execution simply ended. Durable design accepts interruptions and recovers to a correct state anyway. Watermill will not invent that durability for you. It will make durable patterns easier to express once you decide they matter.
Publish inside the same transaction, or admit you did not
The classic race looks innocent.
- Insert the order in Postgres
- Publish
OrderCreated - Crash between 1 and 2
The database believes the order exists. Downstream services never saw the fact. Support tickets arrive days later, and the on-call engineer stares at two systems that are both "correct" in isolation.
Watermill's Forwarder (outbox) pattern exists for this. Write domain state and the outbound message in one database transaction. A Forwarder process later moves outbox rows to Kafka, Redis Streams, or whatever broker you chose.
Command handler
└─ BEGIN
├─ mutate domain tables
└─ insert outbox row
COMMIT
Forwarder
└─ read outbox → publish → mark forwarded
If you skip the outbox, document the publish as best-effort and design compensations on purpose. Silent dual-write is how event-driven architecture earns a bad name inside a company that then refuses to try again.
Handler groups change the meaning of "this message"
By default, CQRS-style processors often create one subscriber per handler. That is fine when each topic carries one event type. When many event types share a topic and order matters, handler groups put those handlers behind one subscriber so the stream is processed sequentially.
The sharp edge caught me later than it should have. If handler A succeeds and handler B fails inside the same group, redelivery runs A again. So keep one handler per event type per group when you can, make every handler idempotent anyway, and treat AckOnUnknownEvent as a loaded gun. Silence is not the same as correctness.
Exactly-once is a wish; idempotency is a design
Most brokers Watermill wraps deliver at-least-once. Rebalances, retries, and partial failures will replay messages. Plan for it on day one, not after the first double charge.
Practical options that have held up for me:
- Store
message.UUIDor a business key in a processed-events table with a unique index - Prefer naturally idempotent writes (
UPSERT, conditional updates) - Use a short TTL cache only for soft dedupe, never for money
func (h *MarkPaid) Handle(ctx context.Context, e *OrderPaid) error {
ok, err := h.dedupe.TryOnce(ctx, e.EventID)
if err != nil {
return err
}
if !ok {
return nil
}
return h.orders.MarkPaid(ctx, e.OrderID)
}
Retries without a poison path are optional self-harm
Retry middleware is mandatory. Infinite retries without a dead-letter path are how one bad payload pins a consumer forever. Wire three layers: backoff retries for transient errors, poison or dead-letter after a capped budget, and a circuit breaker in front of fragile dependencies.
Separate handler bugs from dependency outages. A nil pointer should not share the same budget as a 503 from payments. Alert on dead-letter depth the moment the first message lands. A quiet DLQ is a forgotten one.
Request-reply is not your default hop
Watermill supports request-reply. Used everywhere, it reintroduces the coupling events were meant to remove. Prefer async facts for things that already happened, and plain RPC when you truly need an immediate answer. Keep request-reply for the rare case where a client is waiting and the command must be processed now.
Topics that mean everything age badly
One domain-events topic with forty payload types feels convenient in week one. By month three it becomes an ordering, ACL, and schema-evolution mess. Prefer one aggregate type or bounded context per topic family. Version payloads deliberately. Pass correlation and causation ids in metadata so a cascade is greppable when someone asks "what started this?"
Prove the hard parts
Before calling a service done, I want a short checklist that is not theater:
- State and event publish are atomic (outbox) or explicitly best-effort
- Handlers survive duplicate delivery
- Retry plus dead-letter are configured and alerted
- Ordering needs map to subscriber topology
- Lag, error class, and DLQ depth are on a dashboard
A surprising number of "Watermill bugs" are schema bugs wearing a messaging costume. Add fields carefully. Keep consumers tolerant of unknowns. Contract-test critical events in CI.
Closing
Watermill removes boilerplate. It does not remove distributed-systems physics. Treat the outbox as default for anything money-adjacent, design for at-least-once, and be intentional about handler groups. The library will happily help you build a resilient backbone, or a very fast way to spread inconsistency. The difference is almost never the import path. It is whether you respected the gaps between steps.
Top comments (1)
The forwarder deserves the same scrutiny as the handlers. If it crashes between publish and marking the outbox row forwarded, that event goes out twice, so your dedupe key has to cover forwarder replays too, not just consumer retries. And ordering is a forwarder problem as well: two forwarder workers polling without ORDER BY id will scramble per-aggregate order before the broker ever sees the events. One worker per partition key, or SELECT FOR UPDATE SKIP LOCKED with strict ordering, keeps the outbox sequence intact end to end.