advanced 2 min answer

Alibaba reported a peak of 583,000 order creations per second during Singles' Day 2020. Take that as the sizing input for a log-based order-event pipeline with per-buyer ordering. Roughly how many partitions does the order topic need, and which assumption dominates the error?

alibabapartitionscapacitypeak-event-readinessordering
Show the full answer Hide the answer

The assumptions, stated

  • 583,000 events per second at peak, keyed by buyer id so that one buyer's create, pay and cancel stay ordered.
  • Event size 600 bytes serialised with a registry id rather than an inline schema.
  • Replication factor 3, producers batching with compression at roughly 4:1.
  • A consumer instance processes 3,000 events per second, which assumes one validation pass and one indexed write per event.
  • Partitions are a one-way door: raising the count re-maps hash(key) % N, so the number chosen now has to survive the next few years.

The arithmetic

Raw ingest: 583,000 × 600 B ≈ 350 MB/s before replication, about 90 MB/s on the wire after 4:1 compression, and 270 MB/s of broker write across three replicas.

Partitions needed for bytes: a partition on modern NVMe sustains roughly 10 MB/s without the tail getting ugly, so bytes alone want about 35 partitions. That is not the constraint.

Partitions needed for consumers: 583,000 ÷ 3,000 ≈ 195 consumer instances, and a consumer group cannot have more useful members than partitions, so ≥195 partitions.

Headroom for skew and growth: buyer-id hashing is close to uniform, but partition load still varies by roughly ±25% at this count, and a partition must be able to absorb its own backlog after a consumer restart. Take 2×: about 400 partitions, round to 512.

The number and its range

512 partitions, with a defensible range of 256 to 1,024. That is 1,536 replicas, which a cluster of nine to twelve brokers carries comfortably: well under the roughly 4,000 partitions per broker that was the practical ceiling on ZooKeeper-backed clusters, and far under what KRaft metadata handles.

Which assumption dominates the error

The per-consumer processing rate, by a wide margin. Halve it to 1,500 events per second and the answer doubles. Double the event size and the answer does not move at all, because bytes were never the binding constraint. Measure one consumer instance against real records before choosing the number, and measure it with the downstream write included, since that is usually what sets the rate.

The second-largest error source is the peak ratio. 583,000 per second is a seconds-long spike, and the steady state is orders of magnitude lower. Partitions are cheap to hold and expensive to add, so sizing for the spike and running it year-round is the correct call here, which is unusual among capacity decisions.

When this is the wrong number to size on

If per-buyer ordering is not actually required, drop the key and the whole exercise changes. Round-robin partitioning lets consumers scale independently of partition count semantics, and 35 partitions suffice. Ordering per buyer matters only when a later event can invalidate an earlier one, as a cancel invalidates a create. If the downstream store is an upsert keyed by order id with a version number, out-of-order delivery resolves itself and you have bought 195 partitions of coupling for nothing.