Distributed Order Fulfillment
When you order a new bike, a lot happens before it reaches your door. Someone checks that a bike is on the shelf and puts your name on it. The bank takes the money. A packer boxes it up. A truck picks it up, and a driver brings it to your house. Each job is done by a different team, at a different speed, and any of them can hit a snag. Distributed order fulfillment is how we keep this relay race organized: one coordinator who knows the plan, a clear record of where every order stands, patient retries when something hiccups, and a special pile for problems that need a person to look at them.
Long-running workflows
A relay race where every runner is on a different team
Buying a bike is not one step. It is a chain: reserve the bike, take the payment, pack the box, ship it, and deliver it. Each step belongs to a different service, which is a separate program run by its own team, keeping its own records.
These steps do not all happen in one quick moment. Reserving the bike takes a blink. Payment may take a few seconds. Packing waits until the warehouse opens in the morning. Delivery can take days. So we cannot keep one phone call open from start to finish. Instead, services talk by leaving messages in queues, like notes in classroom cubbies, and each one picks up its notes when it is ready. Grown-ups call this asynchronous, which simply means nobody stands around waiting for the answer.
Messages make the system calm and sturdy. If the warehouse computer is down for ten minutes, its notes simply wait in the queue until it comes back. The price is that things now happen out of sight, sometimes out of order, and sometimes twice. So we need a clear plan, a clear record, and clear rules, which is what the rest of this page is about.
Remember
A long order is a chain of separate steps connected by messages, not one long phone call.
Orchestration and the saga pattern
One coach with a clipboard calls each play
Someone has to know the plan. An orchestrator is like a coach with a clipboard. It does not do the work itself. It sends a command, such as reserve one blue bike, into the inventory team's queue. Then it waits for an event to come back, such as bike reserved. When that arrives, it sends the next command, take the payment, and so on down the list. (Another style, called choreography, has no coach: each team just reacts to the others' events. It has fewer parts in the middle, but the whole plan is harder to see in one place.)
A multi-step plan like this, where each step saves its own work and has a matching undo, is called a saga. We cannot wrap all the teams in one giant all-or-nothing save, because they keep separate records and the steps take days. So when a later step fails, the coach runs undo steps, called compensations, for the steps that already finished, newest first. If the payment fails, the reserved bike is put back on the shelf. If the bike cannot be shipped, the money is refunded and the bike is put back.
A compensation is not a magic eraser. The payment really happened, so its undo is a new action, a refund, and it shows up on the bank statement. And while a saga is still running, other people can see the in-between states, like a bike that is reserved but not yet paid for.
The coach also keeps an eye on the clock. If an event does not come back in time, say the warehouse has not packed the bike after two days, a timeout fires. The coach can retry, try another warehouse, or escalate, which means alerting a person to take a look.
Remember
The orchestrator sends commands, listens for events, watches the clock, and runs undo steps when something fails.
Finite state machine
A board game where the order sits on exactly one square
Every order is always in exactly one named state, like a game piece that sits on exactly one square of a board game: Placed, Stock reserved, Paid, Packed, Shipped, or Delivered. There are also two side squares, Cancelled and Refunded.
The rules list which moves are allowed, like the arrows printed on the board. Placed can move to Stock reserved. Paid can move to Packed. But Shipped can never jump back to Placed. If a message asks for a move that is not on the list, the order refuses it and someone is told. This turns a big, foggy question (what on earth is going on with this order) into a short, clear list that both people and programs can check.
It also makes repeated messages harmless. Messages sometimes arrive twice. If the order is already Paid and a second payment received message shows up, the rules say that move is already done, so we ignore it instead of charging twice or packing two bikes.
Remember
One state at a time, and only the moves on the list are allowed.
Optimistic concurrency
Write the page number on your change so nobody overwrites a newer page
Each time an order changes state, we save the change in the database along with a version number that goes up by one every time, like the page numbers in a diary. Placed is version 1, Stock reserved is version 2, Paid is version 3, and so on. Saving every change also gives us a history, so we can see exactly what happened to an order and when.
Sometimes two helpers try to change the same order at the same moment. Say the warehouse helper wants to mark it Packed, while another helper, handling the customer's cancel button, wants to mark it Refunded. Both read the order when it was Paid, at version 3. The rule is: save my change only if the order is still at version 3. The first save wins and makes it version 4. The second save finds version 4, not 3, so it is refused. That helper reads the order again and decides what to do now that it is Packed.
This is called optimistic concurrency, because we hopefully assume clashes are rare and only check at the moment of saving, instead of locking the order the whole time. It is fast when clashes are rare. When clashes are common, many saves get refused and must be retried, which wastes work.
Remember
Save a state change only if the version is still the one you read.
Retries and dead-letter queues
If a note will not go through, try again, then put it on the problem pile
Some failures are just hiccups: the payment service was busy for a second, or the network blinked. For those, we retry. But we wait a little longer each time, say 1 second, then 2, then 4. This is called backoff, and it stops us from hammering a service that is already struggling. Adding a little randomness to each wait helps too, so a thousand retries do not all land at the same instant.
Some messages will never work, however many times we try. Maybe a field is missing, or the order points at a product that no longer exists. That is called a poison message. If we kept retrying it forever, it could block the queue, like one kid stuck in the doorway holding up the whole line.
So after a set number of tries, say 5, the message is moved to a dead-letter queue, a special problem pile on the side. The main line keeps moving. An alert tells a person, who finds the cause and fixes it, maybe in the code, maybe in the data. Then they redrive the message, which means putting it back into the main queue to be processed again.
A dead-letter queue must be watched. If nobody looks, it becomes a black hole where orders quietly disappear, and a customer waits for a bike that is never coming. Good teams alert as soon as anything lands there and keep track of how long messages sit in it.
Remember
Retry hiccups with growing waits, move stubborn messages to a watched problem pile, and replay them after a fix.
Idempotent handlers and the transactional outbox
Doing a job twice should not count twice, and never forget to pass the note
Queues usually promise at least once delivery: every message arrives, but now and then it arrives twice. Retries and redrives send things again too. So every handler, the code that acts on a message, must be idempotent. That big word means doing it twice has the same effect as doing it once. A simple way is to give every message a unique id, remember which ids are already finished, and skip any you have seen before. The state machine helps as well, since a Paid order ignores a second paid message.
There is one more trap. When the payment helper finishes, it must do two things: save the order as Paid in its database, and send a message saying the order was paid. If it saves and then crashes before sending, nobody hears the news and the bike is never packed. If it sends first and then crashes before saving, everyone hears news that its own records do not show.
The transactional outbox fixes this. In one single save, called a transaction, which happens completely or not at all, the helper writes both the new order state and the message into an outbox table in the same database. A separate helper, the relay, reads new rows from the outbox, sends them to the queue, and marks them as sent. If anything crashes, the message is still waiting in the outbox. It might be sent twice, which is fine, because the handlers on the other side are idempotent.
Remember
Save the change and its message together, and make every handler safe to run twice.
Quick recap
- An order is a chain of asynchronous steps across inventory, payment, warehouse, and shipping, connected by queues.
- An orchestrator sends commands, listens for events, watches for timeouts, and runs compensations when a step fails.
- A compensation is a new action that undoes the effect, like a refund, not an eraser.
- A finite state machine keeps each order in exactly one state and rejects moves that are not on the list.
- Version numbers stop two helpers from overwriting each other's changes.
- Retry hiccups with backoff. After several failures, move the message to a dead-letter queue, alert someone, fix it, and redrive it.
- Handlers must be idempotent, and the transactional outbox makes sure every state change gets announced.
Grown-up words
and what they mean in plain words
- Orchestrator
- The coordinator that knows the steps and tells each service what to do next.
- Saga
- A long, multi-step process where every step has a matching undo.
- Compensation
- The undo action for a step that already finished, like a refund.
- Finite state machine
- A fixed list of states and the allowed moves between them.
- Optimistic concurrency
- Saving only if nobody changed the record since you read it, checked with a version number.
- Backoff
- Waiting a little longer before each new retry.
- Dead-letter queue (DLQ)
- A side queue for messages that failed too many times.
- Redrive
- Putting messages from the dead-letter queue back to be processed again.
- Idempotent
- Safe to do twice: the second time changes nothing.
- Transactional outbox
- Saving an outgoing message in the same database save as the change it describes.