The architecture seems rather overengineered considering a single raspberry pi could do the job, even after 100x scaling!
The architecture seems rather overengineered considering a single raspberry pi could do the job, even after 100x scaling!
WSE back then had a limit of 10k messages per second (as in one message every 1/10000th of a second). Messages came by two separate network operators so that was 20k messages to be deduplicated to 10k operations. Incomming messages came in compressed so they required uncompressing.
Responding to the message required complex processing and then may have resulted in an order to market which, again, had to be constructed and validated.
All this worked on a single server (regular two-socket Xeon-based server).
There was 10us (microseconds..) time budget to send response to the market and it had to work every time even under maximum possible load (10k messages per second).
Routing and forwarding 300 messages per second doesn't seem like something to brag about...
That said, it really doesn't sound too difficult with straightforward architecting (says the armchair critic).
To start, there's nothing here about what machine this architecture runs on, it could be running on a Raspberry Pi for all we know.
Then, this ignores the cost of database lookups for the keys. That data is probably small enough to be on the one machine, but then you have to have service support for (reliably, in real time) syncing that data to the service. A separate database is therefore probably the right solution here, which means you're doing networking in each message send, which makes it unlikely that a Pi could do this.
Next up you've got the issues of reliability. The message queue separation gives you better reliability in the face of issues such as upstream APIs going down or erroring, or issues for a specific user. All the business logic around handling this, the message queue handling persistence and ACID semantics (or parts of it), this all takes additional resources, not to mention potentially a fair bit of disk space (for a Pi) to queue up undelivered messages should an upstream API slow down or stop accepting new messages.
Then you have hardware failure, at this scale you don't want a single machine failure to wipe out your primary communication method with millions of customers. You'd therefore want to have a distributed system, even if that's only for reliability rather than performance.
Lastly, 21 million clock cycles might sound like a lot, and might go a long way with C/C++/Rust, but as you move up to more dynamic languages that will reduce significantly. It happens that they are using Go here, and that's likely to get pretty good performance out of the hardware, but writing this service in Python/Ruby would be a very valid choice for developer productivity, or based on existing skills they have in the team. That might be 1/10th the performance, but since you need a distributed system for reliability here anyway, adding a few more machines to the pool might be a better choice than introducing a lower level language that takes longer to develop and exposes you to memory safety or threading issues.
There may well be other factors I haven't considered here, but I think for the use case of delivering that scale of messages, the reliability options you get with a system like this are well worth the additional hardware requirements and architecture overhead.
Edit: lmilcin makes a good point about trading systems, but there are several differences – that system still has a single point of failure, it was probably written in a low level language with input from experts on performance, and it was running on a much faster machine. The single point of failure of a server-grade machine like that is probably an acceptable risk if you own the hardware, but in a cloud environment (which brings other benefits) hardware is less reliable so probably not an acceptable risk there. I don't think it's an apples-to-apples comparison, although it is interesting.
> Then you have hardware failure, at this scale
It's really, really, not scale!
> Lastly, 21 million clock cycles might sound like a lot, and might go a long way with C/C++/Rust, but as you move up to more dynamic languages that will reduce significantly
Then you're prob using the wrong lang.
I don't really accept that level of engineering is necessary all round, unless the business case requires it, and then I'd speak to whoever put those business requirements together and ask hard questions.
Talking to 1 million different people an hour is _business_ scale. Regardless of the tech required to do that, it not working would likely be a significant business impact.
> Then you're prob using the wrong lang.
There's so much more here than performance. There's developer productivity, there's tooling availability, all sorts. They happen to be using Go which is probably the best trade-offs for this particular system, but if you were doing complex machine learning you'd probably want to use Python due to all the excellent tooling available, and that means having a "slow" language for parts that aren't optimised for you.
If the difference is 1 engineer, ~1k lines of code, 3 machines, vs 3 engineers, ~10k lines of code and 1 machine, the former is likely to be the right trade off for most companies. It's cheaper to build, and since number of bugs typically correlates to lines of code, it will likely be much more reliable.
As for the ACID semantics, you're right you wouldn't need them all here, redelivery is probably fine within some bounds, and the upstream APIs might even have idempotency tokens to prevent this, but the downtime of losing a machine for a few hours, the time to regenerate all those messages that were lost on the downed machine, etc, could equate to quite a bit of downtime and poor UX for users.
You're not wrong from a performance perspective, but taking into account reliability, business impact, user experience, and developer productivity/costs, I think the solution in the blog post is a better set of trade-offs than you're suggesting.
Replace raspberry pi with any local server configuration.
One million per hour? I forgot how to count that low.