A distributed system is a group of computers that work together over a network and look like one system to the people using it. It lets an app grow past what one server can handle, but it also swaps one simple kind of failure for many strange ones, and most of the work is in handling those.
From one server to many
Most apps start on a single server. The web app, the business logic and often the database all live in one place, and a function call from 'place order' to 'take payment' either happens or it doesn't.
That stops working when traffic grows. You can buy a bigger machine for a while (scaling up), but eventually you split the work across several machines (scaling out). A common split is by job: one service handles orders, another handles payments, another sends notifications. Each can be scaled, deployed and restarted on its own, and a bug in the notification code no longer takes down checkout.
From the outside it still looks like one shop. Behind it, a single order now involves several processes on several machines, talking to each other over the network. That is what makes it a distributed system, and it is where the trouble starts.
Why the network changes everything
Inside one process, calling a function is reliable: it runs, and you get a result or an exception. A call to another service over the network is a message, and messages can:
- arrive late, because a link is congested or a server is busy;
- arrive twice, because something along the way retried;
- never arrive, because a connection dropped or a machine restarted.
On top of that, servers crash halfway through their work, and a service may answer from a copy of the data that is a few seconds out of date (stale data).
The hardest part is partial failure. Some pieces work while others don't, and the caller often can't tell which. If you send a request and hear nothing back, there are three very different possibilities:
- The request never reached the other service.
- The service did the work, but the reply was lost.
- The service is still working and is just slow.
From the caller's side, all three look identical: a timeout. This is why 'did the payment go through?' becomes a genuinely hard question. The CAP theorem is one formal take on the same problem: when the network splits, you have to choose between answering with possibly stale data and refusing to answer.
The double-charge problem
Here is how a careless retry turns a lost reply into a real cost:
How a careless retry charges a customer twice
Step 1 of 5: The order service asks the payment service to charge the customer.
Retrying is not the mistake. In a distributed system you have to retry, because many failures are brief. The mistake is retrying an operation that is not safe to repeat.
Idempotency keys
An operation is idempotent when doing it twice has the same effect as doing it once. 'Set the order status to paid' is idempotent. 'Charge £12' is not, unless you make it so.
The usual fix is an idempotency key: a unique id the caller creates once for the operation and sends with every attempt. The receiving service remembers which keys it has already handled and, for a repeat, returns the stored result instead of doing the work again.
On the payment side, that looks like this:
def charge(amount, idempotency_key):
done = results.get(idempotency_key)
if done:
return done # a repeat: no new charge
result = card.charge(amount)
results.save(idempotency_key, result)
return resultAnd the caller retries with the same key every time:
def charge_with_retry(order, attempts=3):
for _ in range(attempts):
try:
return payments.charge(
amount=order.total,
idempotency_key=order.payment_key,
)
except TimeoutError:
continue
raise PaymentUnknown(order.id)Two details matter. The key belongs to the order, created once and stored, not generated fresh for each attempt; a new key per retry defeats the point. And in a real system the check and the save must be atomic (a unique constraint on the key column, for example), or two attempts arriving at the same moment could both charge. Payment providers such as Stripe accept an idempotency key on requests for exactly this reason.
Workflows and durable execution
Idempotency makes each step safe to retry. It doesn't answer the bigger question: if the server running the order process crashes after charging the customer but before shipping, who picks up where it stopped?
Writing that logic yourself means storing the progress of every order after every step, scanning for stuck orders, and deciding what to resume. It is easy to get subtly wrong. This is the problem a workflow engine such as Temporal solves.
In Temporal you write the whole business process as ordinary code, called a workflow. Each step that talks to the outside world (a database, a payment API, an email service) is an activity:
from datetime import timedelta
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from shop.activities import (
create_order, charge, ship, notify,
)
T = timedelta(seconds=30)
@workflow.defn
class OrderWorkflow:
@workflow.run
async def run(self, order_id: str) -> None:
await workflow.execute_activity(
create_order, order_id,
start_to_close_timeout=T)
await workflow.execute_activity(
charge, order_id,
start_to_close_timeout=T)
await workflow.execute_activity(
ship, order_id,
start_to_close_timeout=T)
await workflow.execute_activity(
notify, order_id,
start_to_close_timeout=T)It reads like a plain script, but it runs differently. The Temporal service records an event history for each workflow: which activities started, and what each one returned. If an activity fails, Temporal retries it according to a retry policy, with backoff between attempts, so a brief outage in the payment service doesn't fail the order. If the worker running the workflow crashes, another worker picks it up, replays the workflow code against the recorded history, skips the steps that already finished (their results come from the history), and carries on from the next one.
That is durable execution: the program's progress survives crashes, restarts and deploys, as if the process had never stopped.
Temporal doesn't remove the need for idempotency. An activity can still run more than once, for example if a worker crashes after charging the card but before reporting success. So the activities themselves should be idempotent, and the charge activity passes a key built from the order id:
from temporalio import activity
@activity.defn
async def charge(order_id: str) -> str:
order = await orders.get(order_id)
return await payments.charge(
amount=order.total,
idempotency_key=f"charge-{order_id}",
)Trade-offs and when to use it
Distributing a system is a cost you pay for scale and independence, not a goal in itself. A single well-built service with one database is easier to build, test and debug, and many apps never outgrow one. Split things up when a real limit forces you to: load one machine can't carry, teams blocking each other, or parts that need to scale or fail separately.
A workflow engine is worth it when a process spans several services, runs for more than a moment, and must finish correctly: orders, payments, sign-up flows, data pipelines. It adds a service to run (or a hosted one to pay for) and a new programming model to learn. For a single, quick call, a retry with an idempotency key is often enough. Other common tools for the same problem are message queues with retries, the outbox pattern, and sagas, where each step has a compensating action (a refund for a charge) to undo it if a later step fails.
Common mistakes
- Retrying non-idempotent calls. Every retried operation that changes something needs an idempotency key or a natural way to detect a repeat.
- Treating a timeout as a failure. A timeout means 'unknown', not 'didn't happen'. Check, or retry safely.
- No timeouts at all. A call that waits forever ties up a worker and hides the failure.
- Retrying immediately and forever. Back off between attempts and cap them, or a struggling service gets buried under retries.
- Non-deterministic workflow code. Temporal replays workflow code, so it must make the same decisions every time. Reading the clock, generating random numbers or calling an API directly belongs in activities, or in the SDK's replay-safe helpers such as
workflow.now().
Key takeaways
- A distributed system is several machines working together over a network, appearing as one system.
- Networks delay, duplicate and lose messages, and servers crash midway, so partial failure is normal.
- A timeout means the outcome is unknown; retries are necessary but must be safe to repeat.
- Idempotency keys make repeated requests harmless, such as a payment retried after a lost reply.
- Durable execution, as in Temporal, records each step of a workflow so it resumes after a crash instead of starting over.