While we are on the topic I am really curious to know how they solve it actually.
While we are on the topic I am really curious to know how they solve it actually.
The traditional mailbox architectural model can be improved via a few methods though to increase scale.
1. The entire tweet doesn’t need to be duplicated for each follower, only a lookup reference (“ID”) to the tweet.
2. Unlike in a traditional mailbox sense, each user’s timeline is a bounded collection. So rather than maintaining every tweet for every followed user in the timeline, it’s capped to the X most recent.
Denormalizing data refers to moving away from a model where you store one and only one copy of each ‘tweet’ in something like a relational database, to a model where you might actually store a separate copy of each tweet for each follower. That is kind of an extreme example of denormalization, but it’s a good way to illustrate how you could make it near-instant for any user to load their twitter homepage. If you stored the interesting tweets of every person I follow in a table just for me, you would make ‘reads’ (loading my homepage) incredibly cheap, but writes (someone with a lot of followers tweeting) very expensive.
Those trade offs exist everywhere in a system like this, which is (2), you get to decide when you do work. If celebrities with a million followers tweet an average of once a second, but people load their feeds a hundred thousand times a second, it is entirely acceptable to do 10000x more work for a celebrity tweet posting. This is actually the same concept behind using indexes in relational databases, but done more explicitly.
My personal answer to this question would probably start somewhat space-innefficient. I would take each tweet and put it into an event processing queue which writes a reference to it into each followers feed. This sounds nasty, but it scales with the number of writes, not reads, and we have less writes, and it scales linearly with the follower count of the writer. I would then improve efficiency by thinking about dormant accounts, caching, and maybe doing a bit more work on read.
And for something like Twitter, you'd probably Publish but then also log to some kind of "Notifications" store, so if a user did care but was not actively watching, on their next subscription they'd receive the messages they'd missed.
They've got lots of engineering blog content, there may be a better answer here https://blog.twitter.com/engineering/en_us/topics/infrastruc...
What I'd do is have a shared highfanout queue. Once you have more than say 10k followers your tweets go to the high fanout cluster. You'd have hundreds or low thousands of machines, each of which serves a slice of consumers. When the tweet is sent, you write it to this queue. Each consumer is pulling from one of the shards. If you have a thousand workers that means each worker only needs to send a thousand messages in three seconds, which is very doable.
Only about 180k twitter users have more than 20k followers. If each users tweets every 200 seconds, which seems like a high estimate, then that implies a load of about 1kwps for this system, which seems doable, especially if you have a small intermediate layer of distributors which consolidates the write load.
That's just my sketch.
Seems to me you’d need to do the latter right? That way you ensure you process each follower at least once, but each worker doesn’t need to be aware of others it’s just pulling messages from the queue.
If you were going to write, you'd write into the cache, rather than persisting anything to disk. That's my guess at least.
If you only worked in single Wordpress servers, I could see it a big leap.
Me too. Unfortunately, I have no control over recruiters or the questionable candidates they sometimes choose to send in. I just have to do my best to ascertain whether the person is a good potential teammate.
https://blog.twitter.com/engineering/en_us/topics/infrastruc...