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.

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.

  1. A producer builds a command (order.place, order.cancel, or orders.replace-batch), computes the partition from the pair, and publishes it to that partition's RabbitMQ queue, matching:commands.pN.
  2. The queue delivers it to one consumer. The conductor for partition N receives it in conductor.controller.ts, which hands it to ConductorService.processCommand.
  3. The service appends the command to a promise chain, so it cannot start until the previous command has settled.
  4. 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.
  5. 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.
  6. For each candidate it calls calculateDeal in matching.domain.ts, a pure function over decimal.js values that returns either a deal or a reason not to match.
  7. 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.
  8. EXEC runs the MULTI. Then WAITAOF 1 1 blocks until the write is fsynced on the primary and one replica. Only then does the RPC return and the broker message get acknowledged.
  9. 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 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);
}
apps/conductor/src/orders.storage.ts, December 2024. Acquire with SET NX PX, release with DEL.

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');
}
apps/conductor/src/conductor.service.ts, December 2024. The wrapper around every add, cancel, and batch.

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.

Here is the sequence. Each step is something the code above will do without complaint.

  1. Writer A acquires the lock for pair 12 with a 10-second TTL and starts matching an order.
  2. 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.
  3. The lock expires. Redis does not tell anyone. Expiry is silent by design.
  4. Writer B acquires the lock for pair 12 cleanly and starts matching the same book.
  5. A finishes, reaches the finally, and calls releasePairLock. That DEL removes B's lock while B is still in the middle of its match.
  6. Writer C acquires the now-free lock. B and C are both mutating one order book.
How an expiring lock with an unconditional release allows two writersWriter A's lock expires mid-match, writer B acquires it, then A's unconditional DEL removes B's lock. Writer C then acquires the free key while B is still mutating the same order book.Writer CWriter BRedis lock keyWriter Adeep sweep / GC pause /Redis stall — takes > 5sTTL expires silentlyunconditional — deletesB's lock, not A'sB and C now both mutatethe same order bookSET lock:BTC-USD PX 5000 NX1OK2SET lock:BTC-USD PX 5000 NX3OK — B starts matching4DEL lock:BTC-USD5SET lock:BTC-USD PX 5000 NX6OK7
Steps 5 and 6. A's release deletes B's lock, so C can acquire a key that B still believes it holds. From that point two writers are mutating one 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.

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);
}
apps/conductor/src/conductor.controller.ts. One queue per partition, one active consumer, ordered delivery.
Deterministic routing to one active writer per partitionProducers route each pair to a partition queue. Each quorum queue permits a single active consumer, so exactly one conductor mutates each partition's state.producersrouting: pairId % partitionCountquorum queuematching:commands.p0x-single-active-consumerquorum queuematching:commands.p1quorum queuematching:commands.pNconductor p0ONE writerconductor p1ONE writerconductor pNONE writerstate cell p0state cell p1state cell pN
Commands for a partition queue up in order. The broker delivers them to one conductor at a time, so ordering no longer depends on how long any step takes.

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.

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;
}
apps/conductor/src/conductor.service.ts. Serialisation inside the process, with no mutex.

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.

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.

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}:events
libs/shared-lib/src/matching/contracts.ts. The routing function. Every producer, conductor, and outbox has to agree on the count.
What parallelises around a serial matching decisionAuthentication, validation, and risk checks run in parallel before the ordered decision. Settlement, market data, surveillance, and analytics run in parallel after it. Only the choice of which resting order trades next is serial.Parallel — consumes sequenced factssettlementmarket datasurveillanceanalyticsSerial — the ONE ordered decisionwhich resting ordertrades nextParallel — no priority decided hereauthvalidaterisk checks
Everything that does not decide execution priority runs in parallel. The one state transition that does is serialised per partition.

Partition 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.

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.

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/.

BoundaryWhat it includesScriptThroughput
Pure domaincalculateDeal only, no I/Omatching-domain-performance.js107,000 to 135,000 calc/s
ConductorQueue, validation, Redis MULTI, WAITAOF, event appendconductor-performance.js292 to 382 commands/s
Full pipelineConductor plus outbox publishing with broker confirmsconductor-performance.js with BENCH_WAIT_OUTBOX=1106 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.

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 1 before the broker acknowledgement, and the process refuses to boot in that mode against a datastore that does not support the command.

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.

ConcernSingle writer per marketParallel writers on one book
PriorityOne explicit command orderNeeds coordination and conflict rules
Hot pathNo book-state locks at allLocks, atomics, retries, or a deterministic merge
ReplayDeterministic by constructionMust reproduce scheduling-independent outcomes
ScalingAcross markets and partitionsPotentially within one hot market
Failure analysisOne owner, one sequenceMore interleavings and partial outcomes
Default choiceContinuous price-time booksSpecialised algorithms with proven benefit

We use plenty of parallel infrastructure. Parallel authority over one continuous order book is the part that needs exceptional evidence.

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 →