We Deleted the Distributed Lock From Our Matching Engine
Our matching engine serialised orders per trading pair with an expiring Redis lock. Under load, an expired lock let a second writer onto the same order book. We fixed it by removing locking from the design and letting the message broker enforce one consumer per partition. This is the story of the bug, the fix, and what we measured afterwards.
By Oleksii Vasylenko, Technical Lead · Published · Updated · 15 min read
Where this comes from. Building and re-engineering the Bitsten matching engine, from the original lock-based design in late 2024 through the partitioned single-writer rework in August 2026. Code excerpts are taken from the repository history and the current tree.
What a conductor is, and what one command does
Bitsten is a cryptocurrency exchange. The matching engine is the part that takes a new order, finds resting orders on the opposite side of the book that it can trade against, and produces trades. Ours is written in TypeScript on NestJS, keeps its order books in Redis, and talks to the rest of the platform over RabbitMQ. The service that does the matching is called the conductor. That name predates me and I never found a reason to change it.
A conductor owns one partition. A partition is a fixed group of trading pairs, chosen by pairId % partitionCount, and the conductor for that partition is the only process allowed to change those books. Every other service (the API, the liquidity bots, the calculator that settles balances) sends commands to the conductor and reads events from it. It never touches the books directly.
One command goes through the conductor like this. The parts in bold are the ones the rest of this article is about.
- A producer builds a command (
order.place,order.cancel, ororders.replace-batch), computes the partition from the pair, and publishes it to that partition's RabbitMQ queue,matching:commands.pN. - The queue delivers it to one consumer. The conductor for partition N receives it in
conductor.controller.ts, which hands it toConductorService.processCommand. - The service appends the command to a promise chain, so it cannot start until the previous command has settled.
- It checks the partition matches its own, then checks a
processed-command:<commandId>key in Redis. If the key exists, the broker redelivered something already handled, and the stored result is returned without touching the book. - For a place command it checks
accepted-order:<orderId>, marks the order accepted, and starts paging the opposite side of the book out of a Redis sorted set, 64 candidates at a time, already in price-time order. - For each candidate it calls
calculateDealinmatching.domain.ts, a pure function overdecimal.jsvalues that returns either a deal or a reason not to match. - Every book change (HSET of the order record, ZADD or ZREM on the price index, DEL of a closed order) and every event (XADD to the partition's event stream) is staged into one Redis MULTI. Nothing is written until the command is fully computed.
- EXEC runs the MULTI. Then
WAITAOF 1 1blocks until the write is fsynced on the primary and one replica. Only then does the RPC return and the broker message get acknowledged. - A separate worker, the conductor-outbox, reads the event stream and publishes each event to RabbitMQ with publisher confirms. Settlement, market data, and notifications consume from there.
Steps 2 and 3 are the serialisation mechanism. There is no lock anywhere in that list. There used to be, and this is what happened to it.
The Redis lock we had, as it was written
The original conductor did not have one command queue. It had three: a subscribe queue for new orders, a subscribe queue for cancellations, and an RPC queue for batch replace. Adds and cancels for the same pair could arrive on different queues and be processed by different handlers at the same time, so something had to serialise them. In November 2024 that something was the redlock library, taking a lock per order ID with a one-second TTL. It was fragile enough that the commit log for that week reads "Fix conductor lock", "conductor lock release", and then "frop redlock" (the typo is in the commit).
What replaced it in December 2024 was a hand-rolled lock per trading pair. This is the version that shipped and ran for over a year, from orders.storage.ts at that commit:
async acquirePairLock(pairId: number, ttl = 10000): Promise<boolean> {
const lockKey = this.getPairLockKey(pairId);
const result = await this.redis.set(
lockKey,
Date.now().toString(),
'PX',
ttl,
'NX',
);
return result === 'OK';
}
async releasePairLock(pairId: number): Promise<void> {
const lockKey = this.getPairLockKey(pairId);
await this.redis.del(lockKey);
}The value stored under the key is a timestamp, not an owner token. The release deletes whatever is at the key. Most of the time that is the lock the caller took. The caller side wrapped every operation in withPairLock, which retried acquisition three times with a small backoff, and before that even ran waitForPairLock, which polled EXISTS every 100 milliseconds for up to five seconds. So a busy pair spent a noticeable fraction of its time polling Redis to ask whether it was allowed to proceed.
private async withPairLock<T>(pairId: number, operation: () => Promise<T>, retries = 3): Promise<T> {
for (let attempt = 1; attempt <= retries; attempt++) {
let locked = false;
try {
locked = await this.ordersStorage.acquirePairLock(pairId);
if (!locked) {
if (attempt === retries) {
throw new Error(`Failed to acquire pair lock for pair ${pairId}`);
}
const backoff = Math.min(50 * 2 ** (attempt - 1), 500);
await new Promise(resolve => setTimeout(resolve, backoff));
continue;
}
return await operation();
} finally {
if (locked) {
await this.ordersStorage.releasePairLock(pairId); // unconditional
}
}
}
throw new Error('Lock acquisition failed');
}Ten seconds felt like a generous TTL when it was written. A match against a normal book takes a few milliseconds. The problem is what happens on the day it does not.
The interleaving that puts two writers on one book
Here is the sequence. Each step is something the code above will do without complaint.
- Writer A acquires the lock for pair 12 with a 10-second TTL and starts matching an order.
- A's command takes longer than 10 seconds. A deep sweep across many price levels can do it on its own. So can a garbage collection pause, a Redis stall, or an unrelated handler blocking the event loop. During the incident that made us look, it was Redis commands queueing behind a slow replica.
- The lock expires. Redis does not tell anyone. Expiry is silent by design.
- Writer B acquires the lock for pair 12 cleanly and starts matching the same book.
- A finishes, reaches the
finally, and callsreleasePairLock. That DEL removes B's lock while B is still in the middle of its match. - Writer C acquires the now-free lock. B and C are both mutating one order book.
Step 6 is where price-time priority stops being a property of the exchange and becomes a property of the scheduler. Two writers read the same best ask. Both decide to fill it. Both stage a DEL of that maker and both write a deal event. Nothing in the code detects this. We found out from the reconciliation job that compares executed volumes against balances, which reported a mismatch the next morning, and it took a day of reading logs to work back to the lock.
The textbook fix is a fencing token. Store a random owner value instead of a timestamp, and release with a Lua compare-and-delete so you only remove your own lock. That closes step 5. It does nothing for step 3. A can still be in the middle of a match with an expired lock while B legitimately holds a new one. To close that, A has to re-check ownership before every write, which means the ownership check has to be inside the Redis transaction, which means you have rebuilt a worse version of what a message queue would give you for nothing. The RFC we wrote for the rework lists this as finding P0 and the "repair the lock" option as a short-lived stabilisation step at best.
Single active consumer: the queue setting that replaced the lock
Instead of coordinating writers, we stopped having more than one. The three queues became one durable queue per partition. That queue is a RabbitMQ quorum queue with the x-single-active-consumer argument, which means the broker delivers to exactly one registered consumer at a time and fails over to another only when that consumer disconnects. This is the whole mechanism, from the current controller:
@RabbitRPC({
exchange: MATCHING_COMMAND_EXCHANGE,
routingKey: matchingCommandQueue(ASSIGNED_PARTITION),
queue: matchingCommandQueue(ASSIGNED_PARTITION),
queueOptions: {
durable: true,
arguments: {
'x-queue-type': 'quorum',
'x-single-active-consumer': true,
},
},
})
processCommand(payload: unknown): Promise<ProcessedCommandResult> {
return this.service.processCommand(payload);
}RabbitMQ's own documentation recommends single active consumer for the case where messages must be processed in arrival order. The quorum queue type replicates the queue across broker nodes, so the ordering guarantee survives a broker node failure and not only a restart. We also run with prefetch 1, so a conductor that dies mid-command strands at most one unacknowledged message, which the broker redelivers to the next consumer.
What this buys over the lock, concretely:
- No TTL. There is nothing to tune, renew, or lose a race against, because there is no expiry.
- No stale owner. A conductor that lost its consumer status stops receiving messages. It cannot delete anything belonging to the new consumer because there is nothing to delete.
- Backpressure is visible. A slow conductor shows up as queue depth and message age in the RabbitMQ management UI, where before it showed up as lock contention and polling.
- Recovery is "the broker redelivers the last unacknowledged command". The dedup key handles the case where the command already committed before the acknowledgement was lost.
- Scaling is deliberate. You add partitions, not competing consumers for the same pair.
The trade is that one hot pair is limited to one consumer. For a price-time book that is the right constraint. A pair needs one serialisation point somewhere. It is better for that point to be an explicit, observable queue than an accidental lock with a timer on it.
One writer is not the same as one thread
The property that matters is exclusive ownership of the partition state. Which OS primitive provides it does not. The owner can be a thread, an actor, a process, or an event-loop task. Inside the conductor it is a promise chain. Every command appends itself to a tail promise, so command N+1 cannot start until command N has settled, whether it succeeded or threw.
private processingTail: Promise<void> = Promise.resolve();
async processCommand(value: unknown): Promise<ProcessedCommandResult> {
const result = this.processingTail.then(() =>
this.processCommandSerially(value),
);
// Never let a rejection break the chain for the next command.
this.processingTail = result.then(() => undefined, () => undefined);
return result;
}The broker already guarantees one in-flight command per partition, so this is belt and braces. It costs nothing, and it means the invariant still holds if someone later attaches a second command source or a health check starts calling into the service directly. Both of those have happened to us in other services.
Process isolation has been easier to operate than shared memory would be. A crash takes down one partition. Memory growth is visible per conductor. A standby conductor can attach to the queue and take over without a handoff protocol we had to write. The cost is serialisation and a network hop between stages, and the numbers further down say what that costs.
Why shared-write parallelism on one book is hard
Take two aggressive buys reaching the same ask level, with two matching threads. Both read the same best ask and the same remaining quantity. To avoid overfilling the maker they need coordination. To preserve arrival order they need more. To compute fees from consistent state, publish one authoritative event sequence, and keep the cancellation index correct they need more still. A mutex around that critical section makes the matching serial again, with extra steps.
Finer locks split the book into regions. That works until an order crosses several price ranges and has to take several region locks, at which point you inherit lock ordering, retry, and starvation. Optimistic concurrency moves the cost into conflict detection. It looks good under light traffic and collapses at the best price under a burst, which is when it matters.
There is a correctness problem underneath the performance one. A retried command must not receive worse time priority than it originally earned. Any design where scheduling accidents leak into execution priority is wrong regardless of its throughput graph. And a benchmark that leaves out contention, cancellations, partial fills, and invariant checks will reward the design that is least predictable under real load.
What we parallelise instead: independent markets
BTC-USD and ETH-USD share no price-time priority, so their books can advance at the same time. We route each pair to a logical partition with a stable function and give each partition its own queue, conductor, Redis key space, event stream, and outbox.
export function matchingPartitionForPair(
pairId: number,
partitionCount = matchingPartitionCount(),
): number {
return pairId % partitionCount;
}
// pairId -> partition -> queue matching:commands.pN
// -> keys matching:{matching-pN}:*
// -> stream matching:{matching-pN}:eventsPartition count is topology metadata, not a tuning knob. If producers use 16 and consumers use 32, the same pair reaches two owners and you are back to the lock bug with more infrastructure. So parseMatchingCommand recomputes the partition from the pair and rejects a command whose declared partition disagrees, and a conductor throws on any command whose partition is not its own. Changing the count is a migration: stop routing, drain the queues, snapshot or replay state, verify, cut over.
I should say where this stands. MATCHING_PARTITION_COUNT defaults to 1, and production runs with one partition today. The partition plumbing is in place and tested, and the benchmark write-up covers what happened when we ran two cells on one laptop (they were slower, for reasons that have nothing to do with the design). Multi-partition production is the next step, not the current state.
A single very hot symbol is a different problem, and none of the usual ideas survive it. Splitting a continuous book by price range fails when an order crosses ranges. Splitting bids from asks fails because every match touches both. Splitting by account fails because price priority spans accounts. Specialised venues use batch auctions, hardware pipelines, and deterministic parallel algorithms, but those are different execution models with explicit merge rules. For an ordinary continuous limit order book, optimise the single-writer loop and its round trips first, and give a dominant symbol dedicated hardware before you consider changing its semantics.
Node.js worker threads and SharedArrayBuffer
Worker threads run JavaScript in parallel and are useful for CPU-bound work: independent books, replay verification, compression, risk models. They do not make Redis or RabbitMQ calls faster, and spawning a worker per order costs far more than the match itself. If you use them, use a fixed pool or long-lived partition workers with compact commands over bounded channels.
SharedArrayBuffer avoids copies at the price of importing a shared-memory concurrency model into financial state: atomic layouts, publication barriers, ownership transfer, crash handling, version compatibility. For a TypeScript system, separate worker ownership with message passing is far easier to audit. Reach for transferable buffers or a small native component when a profile shows serialisation is the cost. We have not needed either.
Where the time actually goes
This is the part that would have saved the most time if we had measured it before the rework. We separated the engine boundary from the orchestration boundary and ran each on its own, on an M3 Max under Docker Desktop, with a workload of alternating equal-price bids and asks on one pair so that every two commands fully match. The scripts are in the repository under scripts/.
| Boundary | What it includes | Script | Throughput |
|---|---|---|---|
| Pure domain | calculateDeal only, no I/O | matching-domain-performance.js | 107,000 to 135,000 calc/s |
| Conductor | Queue, validation, Redis MULTI, WAITAOF, event append | conductor-performance.js | 292 to 382 commands/s |
| Full pipeline | Conductor plus outbox publishing with broker confirms | conductor-performance.js with BENCH_WAIT_OUTBOX=1 | 106 commands/s |
Single pair, maximum contention, one partition, local Docker. Every run had zero command failures and zero residual book depth. The benchmark write-up has the full environment and the caveats.
Decimal arithmetic and the matching function are about three orders of magnitude away from being the limit. Redis round trips, the durability wait, and confirmed outbox delivery dominate the command path. Adding shared-write threads to the match loop would have complicated ordering while leaving the bottleneck where it was.
So the work went elsewhere: fixed partitions, one active conductor per partition, one atomic state-and-event commit per command, and an outbox that publishes with a bounded confirm window instead of one round trip per event. None of that is threading.
Rules we wrote down
A design is partly defined by what it rules out. These are in the RFC so nobody has to argue them again under deadline pressure.
- No active-active writers for one pair. Not with locks, not with optimistic retry, not with "the window is tiny".
- No cross-partition transactions. Settlement across accounts is downstream work driven by a committed trade event.
- No Redlock, and no hand-rolled Redis lock, as an order-book correctness mechanism. That is the thing we removed.
- No changing the partition count in place. It is a migration with a drain and a cutover.
- No calling a replica "durable" without an acknowledgement contract behind it. Ours is
WAITAOF 1 1before the broker acknowledgement, and the process refuses to boot in that mode against a datastore that does not support the command.
How to test one design against the other
Feed both implementations the same deterministic command trace: passive limits, orders crossing one and many levels, partial fills, cancels near and far from the touch, duplicate commands, repeated order identities, bursts concentrated at one best price. Then compare the complete event sequence and the final book, not the elapsed time. If two runs produce different makers, prices, quantities, or ordering from identical input, the faster one is not an implementation of the same market.
Measure p50 through p99.9 alongside sustained throughput, queue age, conflicts, retries, allocations, and CPU. Run long enough to see garbage collection, AOF rewrites, and backpressure. Then start killing things: a conductor mid-command, a Redis primary, a broker node during a burst. The conductor-tester app in our repository exists for this and has found more bugs than the unit suite has.
A parallel design earns its complexity if it raises sustainable end-to-end capacity on representative traffic, produces byte-equivalent outcomes, and stays recoverable. Agree on that criterion before anyone writes the second matching thread.
Single Writer vs Shared Writers on One Book
| Concern | Single writer per market | Parallel writers on one book |
|---|---|---|
| Priority | One explicit command order | Needs coordination and conflict rules |
| Hot path | No book-state locks at all | Locks, atomics, retries, or a deterministic merge |
| Replay | Deterministic by construction | Must reproduce scheduling-independent outcomes |
| Scaling | Across markets and partitions | Potentially within one hot market |
| Failure analysis | One owner, one sequence | More interleavings and partial outcomes |
| Default choice | Continuous price-time books | Specialised algorithms with proven benefit |
We use plenty of parallel infrastructure. Parallel authority over one continuous order book is the part that needs exceptional evidence.
Related Matching Engine Guides
The four boundaries behind the throughput table above, with the negative results included.
How deterministic market ownership becomes independently scalable, durable matching cells.
Price-level indexes, FIFO queues, and cancellation maps for the single-writer core.
The trading, consistency, capacity, and recovery obligations that sit above this decision.
Related reading
The pillar guide connects threading to price-time priority, order types, durability, duplicate protection, settlement, recovery, and exchange operations.
Read the complete matching engine architecture guide →