AI System Design
← All chapters
Chapter 17 6 min read

Design Nearby Friends

Nearby friends shows a user which of their opted-in friends are currently within a few miles, updating as everyone moves. Unlike a proximity service, every data point moves constantly, so the design is a real-time fan-out system: WebSocket connections, a TTL location cache, and a Redis pub/sub channel per user that pushes each location update to that user's friends.

Architecture at a glance
  1. Mobile client
  2. Load balancer
  3. WebSocket servers
  4. Redis location cache (TTL)
  5. Location history DB (Cassandra)
  6. Redis pub/sub (channel per user)
  7. Subscribed friends' WebSocket handlers

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.

Figure 1Stateful sockets, stateless REST
Stateful sockets, stateless RESTWS +HTTPpersistentRESTfriend listpub/ subSET + TTLappend pointMobile clientssend location every 30 sLoad balancerWebSocket serversstateful · 1 conn per userRedis pub/subone channel per userREST API serversfriends, profilesLocation cacheRedis · latest pos + TTLLocation history DBCassandra, by user_idUser DBprofiles + friendships
Location traffic rides long-lived WebSocket connections, while friend and profile changes use ordinary stateless REST. Redis plays two roles: a TTL cache of each active user's latest position, and a pub/sub bus with one channel per user.

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.
Figure 2One update, filtered at the receiver
One update, filtered at the receiverAliceWS server ALocation cacheRedis pub/subWS server BBob1. location (lat, lng, t)2. append point to history DB3. SET alice → pos, refresh TTL4. PUBLISH channel:alice5. deliver to Bob's handler6. Bob 1.2 mi away: within radius7. push Alice's location + timestamp8. deliver to Dave's handler9. Dave 40 mi away: drop update
Alice's server publishes once and never needs to know who her friends are near; each subscriber's handler does the distance check against its own user's last position. Bob is close and gets the update, Dave is far away and it is silently dropped.

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.

Figure 3Sharding channels across Redis
Sharding channels across Redischannel:alicechannel:bobchannel:carolchannel:davepubsub-1pubsub-2pubsub-3pubsub-4keys go to the nextserver clockwise ↻
Channels are placed on pub/sub servers by consistent hashing, with the ring kept in etcd or ZooKeeper so every WebSocket server routes the same way. Adding a fifth server would move only the channels on one arc, but each moved channel's subscribers must re-subscribe.

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

Concurrent users
10 million
Location update rate
~334K/s (10M / 30 s)
Fan-out pushes
~13–14M/s
Pub/sub memory
~200 GB
Pub/sub servers (CPU-bound)
~140

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

Now practise it