Distributed Matching Engine Architecture with Redis and Valkey

Putting order books in a Redis cluster does not make them safe for several writers. It spreads the race over more machines. What worked for us at Bitsten is a matching cell: one writer owns a fixed group of markets, commits book state and events in one transaction, and scales by adding more cells. This is how a cell is built, how a command moves through it, and what the two-cell benchmark showed.

By Oleksii Vasylenko, Technical Lead · Published · Updated · 15 min read

Where this comes from. Re-engineering the Bitsten matching path around fixed market partitions, same-slot atomic state and outbox commits, permanent order identities, replicated AOF durability, and failure tests for broker redelivery and conductor takeover. The design is accepted for prototype. Production rollout is waiting on soak and failure testing.

Bitsten is a cryptocurrency exchange. Its matching engine takes orders for a trading pair such as BTC-USD and decides which resting orders they trade against. The first version of that engine kept its order books in Redis, and when throughput became a question the obvious proposal was to run more engine replicas against a Redis cluster. This article is about why that does not work for a distributed matching engine, and what we built instead on Redis-compatible storage (Valkey for authoritative state, Dragonfly for state we can rebuild).

Redis and Valkey give you fast data structures, MULTI/EXEC transactions, replication, append-only persistence, Sentinel discovery, and cluster sharding. None of those features decides which order trades first. Two engine processes reading and mutating the same book at the same time will both see the same best ask and both fill it. Moving the keys into a cluster spreads that race across more primaries. It does not remove it.

The pattern we settled on is a matching cell. A fixed set of markets maps to one command queue, one active engine process (we call it the conductor), one primary/replica storage pair, and one outbox worker that publishes committed events. You add capacity by adding cells for independent markets. The datastore is the cell's atomic state and event boundary and nothing more.

Here is the cell as it runs in infra/dockerfiles/matching-cells/docker-compose.valkey.yml, the local topology we benchmark against. Two partitions, p0 and p1, each with its own Valkey 8.1 primary and replica. Both run --appendonly yes --appendfsync always, so every write reaches the AOF file before the command returns. The replica is started with --replicaof matching-p0-primary 6379 and announces its own hostname so Sentinel can reach it. Three Sentinels watch both groups with quorum two, a three-second down-after timeout, and a ten-second failover timeout.

sentinel monitor matching-p0 matching-p0-primary 6379 2
sentinel monitor matching-p1 matching-p1-primary 6379 2
sentinel down-after-milliseconds matching-p0 3000
sentinel failover-timeout matching-p0 10000
sentinel parallel-syncs matching-p0 1
infra/dockerfiles/matching-cells/sentinel-a.conf. Each partition is a named master group. The conductor asks Sentinel for the group by name, never for a fixed host.

In front of the storage sits RabbitMQ. Every partition has one command queue named matching:commands.pN, declared as a quorum queue with the single-active-consumer argument. Whoever produces a command (the order API, the market-maker bots, the cancel path) computes pairId % MATCHING_PARTITION_COUNT and publishes to that queue.

Two Node.js services run per partition. The conductor (apps/conductor) is a NestJS process that consumes the command queue and is the only thing allowed to write to that partition's keys. The outbox (apps/conductor-outbox) reads the events the conductor committed to a Redis stream and publishes them to RabbitMQ for settlement, market data, and the user-facing services. Each process learns which partition it owns from environment variables:

MATCHING_PARTITION_COUNT=16
MATCHING_PARTITION=0
REDIS_SENTINEL_HOSTS=sentinel-a:26379,sentinel-b:26379,sentinel-c:26379
REDIS_SENTINEL_NAME=matching-p0
MATCHING_DURABILITY=aof-replicated
MATCHING_DURABILITY_TIMEOUT_MS=2000
Per-instance configuration from the matching cells ADR. MATCHING_PARTITION_COUNT is shared by every producer and consumer. The rest is set per conductor and per outbox.

matchingAssignedPartition() in libs/shared-lib/src/matching/contracts.ts reads those at module load and throws if the partition is outside 0..count-1, so a mistyped value stops the process before it declares a queue. Then OrdersStorage.onModuleInit runs a capability check: in aof-replicated mode it issues WAITAOF 0 0 1 once and refuses to boot if the datastore returns an error, which is what Dragonfly does. That one call is what stops a Dragonfly-backed cell from starting with a durability policy it cannot honour.

  • One deterministic route per pair, computed from the pair id and the partition count.
  • One quorum queue and one active conductor per partition.
  • One Valkey primary/replica pair per partition, discovered through Sentinel by group name.
  • One outbox worker per partition, holding a lease so a second copy stays idle.
  • One durability policy, checked at boot and enforced on every command.

Take a limit bid on BTC-USD. The API builds a PlaceOrderCommand with a random commandId, the pair, the computed partition, and the order payload, then publishes it to matching:commands.p3 (pair 19 with 16 partitions, say). RabbitMQ delivers it to the one active consumer on that queue. That consumer is the @RabbitRPC handler in conductor.controller.ts, which hands the payload to ConductorService.processCommand.

The service chains every command onto a promise tail so two messages can never be in flight at once inside the process, even though the broker already prevents that. Then processCommandSerially does the work in a fixed order:

  1. Parse and validate the command with parseMatchingCommand. This recomputes the partition from the pair id and rejects a mismatch.
  2. Refuse the command if command.partition is not this conductor's assigned partition.
  3. Look up processed-command:<commandId>. If it exists, run the durability wait and return the stored result. This is broker redelivery, and it is counted as a duplicate rather than logged as an error.
  4. Open a MULTI. Check the accepted-order:<orderId> marker. If present, the order was accepted before under a different command id and the place is skipped.
  5. Page through candidate makers on the opposite side with ZRANGEBYSCORE or ZREVRANGEBYSCORE, 64 at a time, and run calculateDeal from matching.domain.ts against each one.
  6. For each fill, stage the maker update or removal, the taker update, and a deal.executed.v1 event into the same MULTI. Events go in as XADD to the partition stream.
  7. If the taker still has volume and is a limit order, stage it as a resting maker: HSET the order record and ZADD it into the book index. A market order with leftover volume is rejected with insufficient-liquidity instead.
  8. Stage processed-command:<commandId> with the result. EXEC. Then WAITAOF 1 1 with the configured timeout.
  9. Return the result to RabbitMQ, which is the acknowledgement.
The full command path through one matching cellPartition routing, a single-active-consumer quorum queue, one conductor, an atomic same-hash-tag transaction covering state and outbox, replicated AOF acknowledgement before the broker ack, then bounded publisher confirms.WAITAOF 1 1pairId → partitionRabbitMQ quorum queuesingle active consumerconductor — one writer, in orderatomic same-hash-tagMULTI/EXECstate + book + dedup + outboxValkey primary — AOF → diskValkey replica — AOF → diskonly now: ACK the commandpartition outboxbounded publisher confirmsat-least-once event+ consumer inbox
The path of one command. Everything between opening the MULTI and EXEC is staged rather than written, and the broker is acknowledged only after the durability wait returns.

Nothing is written to Redis until EXEC. The candidate reads and the two dedup lookups happen before the transaction is opened, and the matching loop only stages commands. If the process dies anywhere in steps four through eight, Redis has nothing from this command and RabbitMQ redelivers it to the next consumer. If it dies after EXEC but before the acknowledgement, the redelivery hits the processed-command record in step three and returns the same result. That is the whole recovery story for the conductor. It is short because the transaction boundary is in one place.

A MULTI/EXEC only guarantees atomicity when every key lives on the same node. Valkey Cluster divides the keyspace into 16,384 hash slots and, when a key contains braces, hashes only the part inside them. So every key the conductor touches carries the partition tag {matching-pN}. Under Sentinel this is a naming convention. Under a real cluster it is what keeps one command's transaction on one primary.

matching:{matching-p3}:order:<orderId>              # HASH   order record
matching:{matching-p3}:book:<pairId>:ask:rates     # ZSET   price index, member = priority:id
matching:{matching-p3}:book:<pairId>:bid:rates     # ZSET   price index
matching:{matching-p3}:processed-command:<cmdId>   # STRING command result, never expires
matching:{matching-p3}:accepted-order:<orderId>    # STRING permanent marker
matching:{matching-p3}:events                      # STREAM outbox
matching:{matching-p3}:events:dead-letter          # STREAM quarantine
matching:{matching-p3}:outbox-lease                # STRING which outbox process is active
matching:{matching-p3}:outbox-attempts             # HASH   retry counts per eventId
apps/conductor/src/orders.storage.ts. Every key builder prefixes the same tag, so a partition can later move into a cluster slot without renaming anything.

The transaction has to include every fact whose separation would leave an impossible state. Closing an order removes the hash and the sorted-set member together (removeOrder stages a ZREM and a DEL). A partial fill writes the new remaining volume and appends its event in the same commit. The processed-command record and the accepted-order marker go in as well. Without them, a crash between "I matched" and "I remembered that I matched" lets the same order match a second time.

The dual-write gap and the outbox that closes itPublishing after committing lets a crash lose the event entirely. Appending the event inside the state transaction makes delivery at-least-once instead, which a consumer inbox deduplicates.Append event inside the commitstate + event in ONE transactioncrashoutbox replays from the streamevent delivered, possibly twiceconsumer inbox dedupes byeventIdPublish after commit — the gapcommit state to Rediscrashevent never publishedtrade in the book,settlement never hears
Writing state to Redis and then publishing to RabbitMQ leaves a gap where the trade exists but the event was never sent. Appending the event to a stream inside the same MULTI closes it. The remaining risk is a duplicate publish, which consumers can detect.

Broker redelivery and client retries look alike and need different defences. The processed-command record handles the same commandId arriving twice after a crash or a lost acknowledgement. The conductor finds the record, runs ensureDurability again so the caller still gets the durability guarantee, and returns the stored result without touching the book.

That record does nothing for a client that retries the same logical order with a freshly generated command id. If the first attempt filled and the order hash was deleted on close, there is no live record left for command dedup to find. So stagePlaceOrder checks a second key first, keyed by order id, and writes it inside the same MULTI as the first state transition. It is never deleted when the order closes.

const acceptedByCommand = await this.ordersStorage.getAcceptedOrderCommand(
  input.id,
  partition,
);
if (acceptedByCommand) {
  this.duplicateOrders += 1;
  this.logger.warn('Skipping duplicate matching order identity', {
    orderId: input.id, pairId: input.pairId, commandId, acceptedByCommand,
  });
  return;
}
this.ordersStorage.markOrderAccepted(input.id, commandId, partition, transaction);
apps/conductor/src/conductor.service.ts, stagePlaceOrder. A re-placed order becomes a counted no-op instead of a second fill.
Two different duplicate problems needing two different defencesA repeated command identity is handled by a processed-command record. A client retry under a fresh command identity needs a permanent accepted-order marker that outlives the closed order.yesnoyesa repeated request arrivessame commandId?broker redeliveryor lost responseprocessed-command recordreturns the prior resultclient retried under aNEW command identityoriginal alreadyfilled and deleted?command dedup sees nothing→ would match twicepermanent accepted-ordermarkerkeyed by orderId, never removed
Command dedup catches redelivery of the same command. The accepted-order marker catches a retry under a new command id, which only shows up after the original order has already closed.

Two operational caveats. Turning the marker on for an existing dataset means backfilling accepted-order:<orderId> for every historical and live order. A partial backfill leaves old closed orders replayable. And the marker is partition-local on purpose. The orders database still has to enforce globally unique order ids, because making the matching layer reject reuse across partitions would need the global coordination the partitioning exists to avoid.

Outbox delivery stays at-least-once whatever the conductor does. Every financial consumer keeps an inbox keyed by eventId, committed in the same database transaction as the projection it drives.

Asynchronous replication can lose an acknowledged write if the primary dies before the replica receives it. The Valkey Cluster specification says so directly: write safety is best-effort. A service that owns financial state has to define what "success" means instead of inheriting a default, so the conductor waits for the AOF fsync on the primary and on one replica before it acknowledges anything.

async ensureDurability(): Promise<void> {
  const { localAofFsyncs, replicaAofFsyncs, timeoutMs } = this.durability;
  if (localAofFsyncs === 0 && replicaAofFsyncs === 0) return;   // memory mode
  const result = await this.redis.call(
    'WAITAOF', localAofFsyncs, replicaAofFsyncs, timeoutMs,
  );
  if (
    !Array.isArray(result) ||
    Number(result[0]) < localAofFsyncs ||
    Number(result[1]) < replicaAofFsyncs
  ) {
    throw new Error(
      `Matching durability acknowledgement failed: ${JSON.stringify(result)}`,
    );
  }
}
apps/conductor/src/orders.storage.ts. Called after every EXEC and on every duplicate-command return. A short reply means the command fails.

matchingDurabilityPolicy() in config.ts turns MATCHING_DURABILITY=aof-replicated into { localAofFsyncs: 1, replicaAofFsyncs: 1, timeoutMs: 2000 }. Any other value except memory is a startup error. When the datastore cannot satisfy the wait within two seconds, the command throws, RabbitMQ sees a rejected delivery, and nothing has been acknowledged with weaker semantics than the policy says.

This costs latency. The multi-cell benchmark below puts the number at roughly eight times fewer transactions per second than memory mode on the same laptop. Memory mode is still a valid setting, but only as an explicit statement that recovery comes from a separate durable log. The boot check enforces the difference: a datastore without WAITAOF cannot start in replicated-AOF mode.

We had changed this once without noticing. Replacing KeyDB's every-second AOF with Dragonfly's snapshot cron moved the worst-case unreplicated loss window from one second to five minutes in the earlier deployment (the current Dragonfly compose file snapshots every minute). It arrived as part of a datastore migration, not as a durability decision, and it is the reason the policy is now a named setting.

Sentinel and Cluster solve different problems. Sentinel discovers and promotes a primary inside one replicated group. Cluster shards a keyspace across many primaries and redirects clients by hash slot. Because the application already partitions markets, a cell can run on Sentinel per partition: failure domains stay obvious and multi-key transactions need no cluster awareness. A larger deployment could place several partition tags in a real cluster, provided the client is a cluster client and every transaction stays inside one tag.

The old configuration listed REDIS_CLUSTER_NODES for a Nest module that constructs an ioredis.Redis, not an ioredis.Cluster. The list was ignored. Everything worked in development because the single host was reachable, and the "cluster" existed only in the environment file. That setting now throws at boot:

export function matchingRedisOptions(config: RedisConfig): RedisOptions {
  if (config.REDIS_CLUSTER_NODES) {
    throw new Error(
      'REDIS_CLUSTER_NODES requires a real ioredis Cluster client; ' +
      'use partitioned matching cells or wire ClusterModule explicitly',
    );
  }
  // ...
  if (config.REDIS_SENTINEL_HOSTS || config.REDIS_SENTINEL_NAME) {
    if (!config.REDIS_SENTINEL_HOSTS || !config.REDIS_SENTINEL_NAME) {
      throw new Error('REDIS_SENTINEL_HOSTS and REDIS_SENTINEL_NAME must be configured together');
    }
    return {
      ...common,
      sentinels: parseSentinelAddresses(config.REDIS_SENTINEL_HOSTS),
      name: config.REDIS_SENTINEL_NAME,
      role: 'master',
    };
  }
  return { ...common, host: config.REDIS_HOST, port: config.REDIS_PORT };
}
libs/shared-lib/src/matching/config.ts, matchingRedisOptions. Sentinel hosts and the master name must be set together. A cluster node list is rejected outright.

The same options set reconnectOnError for READONLY replies, which is what a client sees when it is still talking to a demoted primary after a Sentinel failover. Failover testing for a cell means stopping the primary container, watching Sentinel promote the replica, confirming the conductor reconnected through the new master, and comparing book and event state before letting traffic back in. We have run that on the local compose stack. We have not run it under sustained load for 24 hours yet, which is the gate before production.

Publishing to RabbitMQ from inside the conductor, after the Redis commit, would reopen the dual-write gap. Instead the conductor only appends to matching:{matching-pN}:events, and a separate process drains it. ConductorOutboxService opens two Redis connections (the module names them data and control), creates the consumer group matching-outbox-v1 with MKSTREAM if it is missing, and runs a loop:

  1. Take the partition lease: SET outbox-lease <consumerId> PX 10000 NX on the control connection. Without the lease the loop sleeps a second and tries again, and /ready returns 503. A timer renews the lease every 3.3 seconds with a compare-and-expire Lua script.
  2. Reclaim entries another consumer left pending for more than 30 seconds with XAUTOCLAIM, and publish those first.
  3. If this consumer still has pending entries of its own, wait. Do not read new ones ahead of them.
  4. Otherwise XREADGROUP ... COUNT 100 BLOCK 1000 for new entries.
  5. Decode each entry. An entry whose payload is not a valid envelope goes to the dead-letter stream in one MULTI (XADD to dead-letter, XACK, XDEL) and the loop moves on.
  6. Publish the batch and, only when every confirm came back, acknowledge the batch in Redis.

The first version published and awaited one event at a time, which preserved order and capped a worker at about 517 events per second on scripts/outbox-drain-performance.js. The current version keeps a bounded window of confirms open. Sends still go out on one channel in stream order, so RabbitMQ receives them in order. The worker just stops paying a full round trip per event.

const confirmations = deliveries.map(({ entry, event }) =>
  publishEvent(rabbit, entry, event),
);
const confirmed = await Promise.allSettled(confirmations);
if (confirmed.some((r) => r.status === 'rejected' || !r.value)) {
  throw new Error('RabbitMQ did not confirm publication batch');
}

const transaction = dataClient.multi();
transaction.xack(streamKey, consumerGroup, ...streamIds);
transaction.xdel(streamKey, ...streamIds);
transaction.hdel(attemptsKey, ...eventIds);
await transaction.exec();
apps/conductor-outbox/src/outbox-delivery.ts, publishBatchAndAcknowledge. allSettled rather than all, so a failed confirm does not leave callbacks firing against a channel that was already torn down.

Each published message carries messageId = eventId, plus headers for the causation id, the partition, and the stream id as x-partition-offset. A failed batch increments a per-event attempt counter in the outbox-attempts hash and retries with jittered exponential backoff capped at 30 seconds. A crash after a broker confirm and before the XACK republishes that batch, and that is the duplicate the consumer inbox exists to absorb. We do not claim exactly-once across Redis and RabbitMQ.

scripts/redis-multicell-performance.js spreads N logical partitions across the URLs it is given and, per partition, runs the conductor's write shape in miniature: HSET an order, ZADD a book member, XADD an event, all in one MULTI, optionally followed by WAITAOF 1 1 2000. We ran it against one cell and two cells on the same Docker Desktop host, for both the Dragonfly compose file and the Valkey one.

TopologyDurabilityTransactions/svs one cell
Dragonfly, 1 cellmemory + replica25,621baseline
Dragonfly, 2 cellsmemory + replicas22,986-10%
Valkey, 1 cellAOF always + WAITAOF 1 13,254baseline
Valkey, 2 cellsAOF always + WAITAOF 1 12,956-9%

Dragonfly: 100,000 operations over 64 partitions. Valkey: 10,000 operations over 16 partitions. Both cells shared one CPU and disk budget on a MacBook.

Two cells were slower than one in both configurations. That was not the result we wanted, and it does not say partitioning fails. Adding a second cell on the same host added four more processes competing for one CPU scheduler, one memory bus, and one disk. It did not add any of the three things the architecture claims to scale with: CPU, disk bandwidth, and a separate failure domain. The test did not test the claim.

The rows also put the durability cost on the table. The only difference between 25,621 and 3,254 transactions per second is the persistence policy. The next test that would mean something needs cells on separate hosts or with enforced CPU and storage budgets, and the pass condition is at least 1.5 times the single-cell throughput with zero invariant violations during the run and during a forced failover. Until that passes, the multi-cell scale claim is a hypothesis with a plausible mechanism.

Every metric carries the partition, because a healthy aggregate hides one stalled market. The conductor exposes command counts, duplicate commands, duplicate order identities, failures, events created, and the last command duration on /metrics. The outbox exposes published, retried, and dead-lettered events, the stream length, the pending count, and whether it owns the lease. Its /ready endpoint returns 503 when it does not, so a standby copy is visibly a standby.

There is also an on-demand consistency check, because "every indexed member has a record and the record agrees with its index" is cheap to verify and expensive to get wrong without noticing. GET /reconcile?pairId=<id> on the conductor scans up to 10,000 members per side, fetches the records in one pipeline, and reports:

{
  "pairId": 12, "partition": 3,
  "indexedOrders": 8412, "scannedOrders": 8412,
  "truncated": false, "healthy": false,
  "issues": [
    "ask:7f3a…:missing-order",   // member in the index, no order hash
    "bid:91c2…:index-mismatch"   // hash exists but its pair or side disagrees
  ]
}
apps/conductor/src/orders.storage.ts, auditPair. A truncated scan is reported as unhealthy even with zero issues, so a deep book cannot pass by accident.

The matching loop tolerates the first kind at runtime. getMatchingMakerPages logs and skips a candidate whose hash is gone rather than aborting the command, because one orphaned member should not wedge a market until someone runs the audit.

Capacity planning follows the busiest pair, not total traffic divided by cell count. A hot partition gets dedicated resources before anyone considers changing its execution model, and relocating a partition is a migration: stop routing, drain the queue, snapshot or replay the keys, verify, change ownership, resume, and watch the sequence continue. Changing MATCHING_PARTITION_COUNT in place is the one thing the design forbids, because producers and consumers disagreeing on the count puts the same pair in front of two writers, which is the bug the whole cell exists to prevent.

TopologyStrengthRiskUse when
In-memory + journalLowest hot-path latencyCustom replay and snapshot toolingThe engine log is authoritative
Valkey + Sentinel cellSimple atomic partition and failoverOne primary's capacity per cellThe application owns market partitioning
Valkey ClusterMany primaries, client-side routingHash-slot and resharding complexityLarge fixed partition sets share infrastructure
Dragonfly memory cellHigh Redis-compatible throughputNo WAITAOF; cannot meet the replicated-AOF contractState is rebuildable from another durable log
Database-centricFamiliar, strong durabilityLocks and tail latency under churnModerate volume, auditability dominates

Decide when a command may be acknowledged and how the book is rebuilt after a failure first. Those two answers rule out most of the rows.

The pillar guide covers the matching policy, order lifecycle, technology choices, double-match prevention, recovery, and production evidence behind these cells.

Read the complete matching engine architecture guide →