Problem & requirements
Users who opt in should see a list of friends within a radius, say 5 miles, with the distance and when the location was last updated. The list refreshes every few seconds, while each client reports its own location about every 30 seconds, a cadence chosen because people walking move only a few dozen metres in that time. Friends who have been inactive for roughly ten minutes should disappear from the list. The system also keeps a location history for analytics and machine learning.
Non-functionally we want low latency so the list feels live, but a few lost updates are fine because the next one arrives soon. Eventual consistency for the location store is acceptable; a short delay across replicas does not change what the user sees in any meaningful way.
Back-of-the-envelope
Suppose 1 billion users, 10% of whom use the feature, giving 100 million daily users, and that 10% of those are online at once: 10 million concurrent users. Each reports every 30 seconds, so location updates arrive at 10M / 30, about 334,000 updates per second.
Fan-out is the scary number. If a user has up to 400 friends and about 10% of them are online and subscribed, each update is pushed to roughly 40 channels' subscribers. 334K x 40 is about 13 million pushes per second. A sanity check: 300K x 40 = 12M, and 334K is slightly more, so 13M is right. That figure, not the raw update rate, will size the pub/sub layer.
High-level design
A peer-to-peer design where every phone keeps connections to every nearby friend is hopeless on mobile networks, so a backend broker sits in between. RESTful API servers handle ordinary stateless work: adding and removing friends, profile updates. WebSocket servers hold a persistent bidirectional connection per online client, receive location updates and push friend updates down. A Redis location cache stores each active user's latest position with a TTL; when the TTL expires the user is treated as inactive. A location history database, typically Cassandra because the workload is append-heavy and partitions cleanly by user, keeps the trail.
The heart of the design is Redis pub/sub with one channel per user. Each user's WebSocket handler subscribes to the channels of all their friends; when a user moves, the server publishes the new location to that user's own channel, and Redis delivers it to every subscribed handler. Channels are extremely cheap, so creating one for every user is fine.
The periodic update and fan-out flow
When a client sends a location, the load balancer routes it over the existing WebSocket connection. The handler writes the point to the history database, refreshes the user's entry and TTL in the location cache, stores it in the connection's own state for later distance checks, and publishes it to the user's pub/sub channel. These steps can run in parallel.
Redis broadcasts the message to every subscriber, which is the WebSocket handler of each online friend. That handler computes the distance between the publisher's new location and its own user's last known location. Only if the friend is within the radius does it forward the update down the socket, together with the timestamp. This design deliberately does the distance filter on the receiving side, which keeps publishing trivially cheap and avoids any central query over all users.
- On connect: load friends, fetch their cached locations, push those within range, subscribe to their channels.
- On disconnect: unsubscribe from friend channels and let the TTL expire the location entry.
Deep dive: scaling each tier
API servers scale trivially because they are stateless. WebSocket servers are stateful: a node can only be removed after its connections drain, so mark it as draining at the load balancer, stop routing new connections to it, and wait for existing ones to close before terminating. Autoscaling must respect this, which makes scale-in slower than scale-out.
The location cache holds only one entry per active user, about 10M keys of under 100 bytes, which would fit in one box, but 334K writes per second is too much for a single Redis node. Because each user's data is independent, shard by user ID across several nodes and replicate each shard for availability. The location cache and the history database both partition naturally on user ID, which is why this problem shards so cleanly.
Deep dive: scaling Redis pub/sub
Memory first: 100M channels, each tracking perhaps 100 subscribers at about 20 bytes per pointer, is 100M x 100 x 20 B = 200 GB, which is two servers at 100 GB each. CPU is the real constraint. At 13–14 million pushes per second and a conservative 100,000 pushes per second per Redis server, we need around 140 servers. So pub/sub must be a distributed cluster.
Shard channels by user ID using consistent hashing, and store the hash ring in a service-discovery system such as etcd or ZooKeeper so WebSocket servers can locate the node for any channel and get notified when the ring changes. Treat this cluster as stateful: resizing moves channels, every subscriber of a moved channel must resubscribe, and the resulting burst of resubscriptions can cause dropped updates and CPU spikes. Over-provision, resize rarely and do it at low-traffic times.
Edge cases and alternatives
Adding a friend triggers a callback to the user's WebSocket handler to subscribe to the new friend's channel; removing one unsubscribes. Users with very many friends are bounded by the platform's friend cap, often a few thousand, and their subscriptions are spread across many pub/sub nodes, so the load is absorbed. A 'nearby random people' feature uses a pub/sub channel per geohash cell: a user subscribes to their own cell and its neighbours and publishes to their cell, mixing the proximity idea with real-time fan-out.
Redis pub/sub is not the only option. Erlang/OTP or another actor model runtime can represent each user as a lightweight process that receives location messages and forwards them to friends' processes, collapsing the WebSocket and pub/sub tiers into one. It is elegant and very efficient, but it requires engineers with that skill set and a platform the team must operate.
Failure handling & wrap-up
Most failures are tolerable because the data is ephemeral. A lost update is replaced within 30 seconds. A failed WebSocket server drops its clients, which reconnect elsewhere, reload friends from the cache and resubscribe. A failed location cache shard loses only the latest points, which are repopulated by the next round of updates. A failed pub/sub node is replaced via service discovery, after which affected subscribers resubscribe.
The core lesson is to recognise a fan-out problem and size it by pushes, not requests. Per-user channels plus receiver-side filtering keep the publisher simple, and every stateful tier needs a plan for draining or resharding.
Key numbers
Key terms
- WebSocket
- A persistent, full-duplex connection that lets a server push data to a client without polling.
- Redis pub/sub
- A lightweight message bus where publishers send to named channels and Redis forwards each message to all current subscribers without storing it.
- Location cache TTL
- An expiry on each user's cached location so inactive users disappear automatically.
- Fan-out
- Delivering one event to many recipients; here, one location update to every online friend.
- Connection draining
- Stopping new connections to a server and waiting for existing ones to finish before removing it.
- Service discovery
- A registry such as etcd or ZooKeeper that tracks live servers and shard ownership and notifies clients of changes.
- Actor model
- A concurrency model where independent lightweight processes communicate only by message passing.
Common mistakes
- Sizing the system on update QPS and ignoring the far larger fan-out push rate.
- Treating WebSocket servers as stateless and killing them without draining.
- Querying all users for proximity on each update instead of filtering per friend on the receiving side.
- Resizing the pub/sub cluster during peak hours and triggering a resubscription storm.
- Over-engineering consistency for data that is replaced every 30 seconds.
Further study
- Redis Pub/Sub documentation
- Facebook Nearby Friends launch (2014)
- Erlang/OTP and the actor model
- Consistent Hashing and Random Trees (Karger et al., 1997)
- Apache Cassandra