Unique IDs Across Shards

Every new order needs an ID. With one database, we can use a counter to assign the next number: 91, 92, 93. With shards, no single database receives every new order. If each shard starts its own counter at 1, two shards can create orders with the same ID.

For Shopend, we want order IDs to be unique across shards, including when we move a shop’s data. A shop can still use its own counter for a display number such as “Order #1001.” That number only needs to be unique within the shop. The order ID is different: it must be unique across all of Shopend.

This is not just about orders. It is about any document or entity that needs a unique identifier across the entire system.

We will explore some strategies for generating unique IDs across shards. We will continue to use orders as the example.

A central counter

We could have one dedicated counter server assign the next number to every new order. This gives us short, increasing IDs. However, every request to create an order also requires a request to that counter server. It adds overhead to the system and latency to each order.

Also, it creates a single point of failure. If that server is down, the shards cannot get new IDs.

A range per shard

Another option is to give each shard a large, non-overlapping range of 64-bit integers. Shard A could start at 1 and stop before 1,000,000,000,000. Shard B could start at 1,000,000,000,000 and use the next range. Each shard then assigns IDs without contacting another server for each order.

We have to assign each new shard an unused range and give a shard another range if it uses up its current one. When we move a shop to another shard, its existing orders keep their IDs. New orders use the destination shard’s range. IDs increase within each shard, but do not sort by creation time across shards.

We could also interleave the numbers: A uses 1, 4, 7; B uses 2, 5, 8; and C uses 3, 6, 9. This pattern works only with exactly three shards. If we add a fourth shard, we have to change the allocation, and we cannot reuse numbers that were already assigned.

UUIDs

A UUID (universally unique identifier) is a 128-bit value, usually written in this form:

f47ac10b-58cc-4372-a567-0e02b2c3d479

There are several versions. A version 4 UUID uses 122 random bits; the remaining bits identify its version and format. Any server can generate one without requesting a number or an assigned range from another server. Two UUIDs could match, but if the random bits are generated properly, the chance is small enough to ignore. Even at a billion UUIDs per second, it would take about 86 years to reach a 50 percent chance of at least one collision.

A UUID takes twice as much space as a 64-bit integer. Random IDs also cause a problem for an index, because the index keeps its entries in sorted order. Each new ID can land in a different part of the index, so the index pages an insert needs are less likely to already be in memory.

Version 7 UUIDs start with a timestamp, so new IDs tend to be inserted near the end of the index. Both versions are specified in RFC 9562, published in May 2024.

Snowflake IDs

Twitter introduced Snowflake IDs in 2010. With this approach, each Shopend API instance can create order IDs itself, without contacting another server for each order.

Structure of a Snowflake ID

A Snowflake ID is a 64-bit integer with four fields, from left to right:

  • 1 bit that is always 0, so the ID is a positive number.
  • 41 bits for a timestamp: the number of milliseconds since a chosen start date. 41 bits cover about 69 years.
  • 10 bits for the worker ID. A worker is a process that creates IDs. For Shopend, each API instance is a worker. We give each worker a different worker ID when it starts. 10 bits allow 1,024 workers. Twitter split these bits into a data center number and a machine number.
  • 12 bits for a sequence number. It counts the IDs a worker creates within the same millisecond, so one worker can create up to 4,096 IDs per millisecond.

To see how the fields work, suppose we write them as decimal groups for readability. At time 5000, API instance 07 could produce 5000 07 001, then 5000 07 002. Instance 12 could produce 5000 12 001 in that same millisecond. In the next millisecond, instance 07 could produce 5001 07 001. The real ID packs these fields into bits.

The timestamp comes before the worker ID and sequence number, so IDs sort roughly by creation time. Clocks on different API instances may differ, and orders may reach the database in a different order. So new IDs tend to be inserted near the end of a sorted index, but not always at the very end.

An instance must never use the same timestamp and sequence number twice with its worker ID. If it uses all 4,096 sequence values in one millisecond, it waits for the next millisecond. If its clock moves backward, it stops creating IDs until the clock passes the last timestamp it used.

MongoDB’s ObjectId uses a related design. It combines a 4-byte timestamp, a 5-byte random value per process, and a 3-byte counter.

For Shopend, we use Snowflake IDs for orders. MongoDB’s default ObjectId would also work, because each process creates it without contacting another server. The examples in this and earlier chapters still use short IDs such as "92" so they are easy to read.