The broker incentive platform we built takes a stream of executed trades and turns it into money owed to introducing brokers — the partners who referred the clients doing the trading. Rebates accrue per deal, and a partner's rate improves as their cumulative volume crosses tier boundaries.
At a few thousand deals a day, processing them one at a time in arrival order is fine and nobody thinks about it. At hundreds of thousands a day we had to process them in parallel, and the moment we did, a problem appeared that had been invisible before.
Listener
Consumes canonical deals off the queue into a raw staging table.Enrichment
Joins symbols, trading accounts and user groups; merges open and closed legs.Reward worker
Calculates rebates and applies tier changes. This is where order starts to matter.The problem
A tier change caused by one deal changes the rate applied to the next one. So if two deals straddle a tier boundary, the order in which they are processed determines which one gets the old rate and which gets the new one — and the totals differ.
A partner earns $5 per lot up to 100 cumulative lots, and $7 per lot above it. The rate applied to a deal is the partner's tier at the moment that deal is processed. The partner is sitting at 98 lots. Two deals arrive at almost the same instant: Deal A, 3 lots and Deal B, 4 lots.
$43
- A: partner at 98 lots, tier 1 → 3 × $5 = $15
- Cumulative 101 → tier 2
- B: tier 2 → 4 × $7 = $28
$41
- B: partner at 98 lots, tier 1 → 4 × $5 = $20
- Cumulative 102 → tier 2
- A: tier 2 → 3 × $7 = $21
The obvious fix, and why it fails
The obvious move was to preserve ordering: partition the stream so that all deals for a given partner land on the same worker, and process each partition strictly in sequence. Correct, and it would have cost most of our parallelism on the hottest partners — exactly the ones generating the volume.
But there was a better reason not to reach for it, and it is the detail that reframed the whole problem:
The upstream trading platforms do not guarantee a perfectly sequential stream of trades in the first place.
So strict in-order processing would not actually have delivered the guarantee it appears to. We would have paid the full cost of serialisation and still been reconstructing an order that the source never promised. A correctness mechanism that depends on an assumption the upstream does not honour is not a correctness mechanism. It is a performance cost with a comforting name.
Once that was clear, the question stopped being "how do we preserve order?" and became "how much does order actually matter, and where?"
The insight
Almost nowhere.
Rebate accrual is accumulative. Within a tier, ten deals in any sequence produce the same total — addition does not care about order. The calculation is commutative across the overwhelming majority of the stream.
It stops being commutative at exactly one event: a tier change. That is the only moment where which-deal-came-first changes the answer, and tier changes are rare relative to deal volume. A busy partner might cross a boundary a handful of times a month against tens of thousands of trades.
So the expensive guarantee is being bought for the whole stream to protect a fraction of a percent of it.
The design
We chose to process deals as they arrive, load-balanced freely across workers with no ordering constraint. Then, when a tier change occurs, a reconciliation runs over a ±10-minute window around it and verifies that the deals near the boundary were applied in the correct sequence — correcting them if not.
The window exists because near-simultaneous deals are the only ones that can race. A deal that arrived twenty minutes before the boundary was not in contention with one that arrived after it; the ordering question is real only for deals close enough in time that their processing could have interleaved. So the reconciliation scope is bounded by the physics of the race rather than by the size of the ledger.
What this buys: full parallelism on the hot path, correctness where correctness is actually at stake, and a reconciliation cost proportional to tier changes rather than to deals.
No partitioning constraint. Full fan-out across workers. Reconcile in a bounded window when a tier boundary is crossed. Parallelism is unconstrained and the correction cost scales with a rare event.
Deals from one tenant processed in execution order; different tenants in parallel. Distributes load well and is simple to reason about — but reintroduces serialisation inside the largest tenants, which is where the volume actually is.
Group by the introducing broker whose referred client traded — finer-grained parallelism than by tenant. Viable, and it depends on whether the upstream queue already partitions on the right key. Deprioritised: tenant-level grouping is sufficient for correctness and this adds coupling to an upstream partitioning scheme we do not control.
Worth noting that the grouping key is not static even in the alternatives. Some rebate programmes attribute to a master partner rather than the direct one, so the correct partition key depends on programme configuration — another reason to avoid making parallelism depend on it.
A second decision, for the same reason
We also made the internal queues between these stages carry only a handle to the deal, not the deal itself. The payload stays in the database and workers fetch it by reference — the claim check pattern.
Message brokers are optimised for high throughput of small messages, typically a few kilobytes. Pushing enriched trade records through them works until volume climbs, at which point the queue becomes the bottleneck and the failure is unpleasant to diagnose because nothing is logically wrong. The alternative — data-rich messages carrying the full payload — is simpler to build and simpler to debug, and it is the right choice at lower volume.
Both decisions came from the same place: at hundreds of thousands of events a day, what moves matters more than what is computed.
The generalisable part
The reusable move here is not the ten-minute window. It is the sequence of questions.
When a calculation appears to require ordering, ask whether it requires ordering everywhere or only at identifiable boundaries. Accumulation, aggregation, counting and summing are commutative; the order-dependence usually lives in a small number of state transitions layered on top — tier changes, limit breaches, threshold crossings, status flips.
If you can name the transitions, you can usually parallelise everything else and pay for correctness only at them. Then check the assumption underneath: does the upstream actually guarantee the order you are working to preserve? Often it does not, and discovering that is what makes the cheaper design obviously right rather than merely faster.
Serialising a whole stream to protect a rare event is a common and expensive default. It is worth the ten minutes to work out how rare.