Email 101
Email predates the web, and its protocols reflect that. SMTP moves mail from one server to another. POP lets a client download messages, traditionally deleting them from the server, which makes multi-device access awkward. IMAP keeps mail on the server and lets clients fetch headers and bodies on demand, which suits phones and multiple devices. Modern webmail usually ignores these at the client edge and speaks plain HTTPS to its own API, using SMTP only between providers.
To deliver to another domain, a sending server looks up that domain's MX records in DNS; each record names a mail server and a priority, and the sender tries the lowest priority number first. Attachments travel inside the message as MIME parts, typically base64-encoded, which inflates binary data by about a third. A traditional mail server stored each mailbox as files on local disk, for example in Maildir format. That works for thousands of users but not billions: disk I/O becomes a bottleneck, a single server is a single point of failure, and features like search, filtering and multi-device sync demand more than a directory of files.
Requirements and estimation
Target 1 billion users who send and receive mail, fetch their inbox, read and filter by folder or label, search by sender, subject or body, and get protection from spam and viruses. The system must be highly available, durable (lost email is unacceptable), reasonably fast for the user's own view, and flexible enough to add features.
Assume each user sends 10 emails per day: 1e9 x 10 / 1e5 seconds ≈ 100,000 send QPS. Assume each user receives 40 emails per day with about 50 KB of metadata each: 1e9 x 40 x 365 x 50 KB ≈ 730 PB of metadata per year. Double-check: 4e10 emails/day x 365 = 1.46e13 emails/year, x 5e4 bytes = 7.3e17 bytes, which is 730 PB. If 20% of emails carry an attachment averaging 500 KB, attachments add 1.46e13 x 0.2 x 500 KB ≈ 1,460 PB per year. These numbers rule out any single-machine or single-database design and tell us attachments belong in cheap object storage.
High-level distributed architecture
Users reach webmail over HTTPS. Stateless web servers handle login, sending, folder listing and fetching messages, so they scale horizontally behind a load balancer. Real-time servers keep persistent WebSocket connections to push new mail to online clients; they are stateful, so they need connection-aware routing, and long polling is a fallback for clients without WebSocket support.
Storage splits by access pattern. A metadata database holds subject, sender, recipients, folder, read flags and body for each email. An attachment store such as Amazon S3 holds large binary files, referenced by key from the metadata. A distributed cache like Redis keeps the most recent emails, because users overwhelmingly read what just arrived. A separate search store, essentially an inverted index, supports full-text queries. Separating these lets each one be scaled and tuned for its own workload instead of forcing one database to be good at everything.
Sending and receiving flows
When a user sends, the web server validates the request (size limits, recipient format) and, if the recipient is on the same domain, may run spam and virus checks and deliver internally. Otherwise the message goes onto an outgoing queue and is copied into the sender's Sent folder. SMTP outgoing workers pull from the queue, check for spam and viruses, and deliver to the recipient's MX server. Failures that may be transient, such as a greylisting response, go to retry; persistent failures land in an error queue for inspection. Decoupling with a queue means a slow remote provider never blocks the user's click, and the queue size becomes a useful health metric.
Incoming mail arrives at SMTP servers that apply acceptance policies early, rejecting invalid recipients or obvious spam before wasting resources. Accepted messages go to an incoming queue, which absorbs bursts and lets processing scale independently. Mail processing workers filter spam, scan for viruses, apply user rules, then write metadata, attachments and the search index. If the recipient is online, the message is pushed via the real-time servers; if not, it simply waits in storage until the next fetch.
Deep dive: the metadata database
Email metadata has a distinctive shape. Headers are small and read constantly; bodies are larger and read once or twice; nearly all reads concern a single user's own mail; most reads target recent messages; and durability matters enormously. Relational databases struggle at this size, blob stores cannot answer folder and unread queries efficiently, and general document stores were not built for these patterns. Large providers use custom systems, but a wide-column NoSQL store like Bigtable or Cassandra partitioned by user_id is a reasonable approximation.
Partitioning by user keeps one user's mail on one partition, so most queries are single-partition. Tables are denormalised around queries: one table for a user's folders, one for emails in a folder keyed by (user_id, folder_id) and clustered by a time-ordered email_id such as a TIMEUUID so newest-first reads are cheap, and one per-email table for body and attachments. NoSQL stores are poor at filtering on non-key columns like is_read, so unread mail is modelled as separate read_emails and unread_emails tables. Marking a message read becomes a delete from one and insert into the other, a write cost paid to make the common read fast.
Threading, consistency and deliverability
Conversation threading uses standard headers: each message has a Message-Id, a reply sets In-Reply-To to its parent, and References lists the chain of ancestors. A client or server can reconstruct a thread tree from these without guessing by subject line. On consistency, a mailbox is assigned a single primary at a time; during failover the mailbox may be briefly unavailable for sync rather than serving divergent copies, because users notice missing or duplicated mail far more than a short delay. Email is one domain where correctness beats availability for the owning user.
Sending mail is only useful if it lands in the inbox, not the spam folder. Providers use dedicated sending IPs, separate IPs for marketing versus transactional mail, and IP warm-up that ramps volume over weeks so receivers build trust. Spammers must be banned quickly so they do not poison shared reputation. Feedback loops with major ISPs report complaints and bounces. Authentication standards matter: SPF lists authorised sending IPs, DKIM signs messages cryptographically, and DMARC tells receivers what to do when those checks fail.
Search and scalability
Email search differs from web search: each query covers only one user's mail, results are sorted by time and filters, and the index is updated with every incoming message, so the workload is write-heavy. One option is Elasticsearch, partitioning documents by user_id and feeding it asynchronously from Kafka; it is mature but means running and syncing a second storage system. The alternative is a custom index embedded in the primary datastore. Because writes dominate, a log-structured merge tree is attractive: writes go to memory and are flushed sequentially to disk, and compaction merges files later, avoiding random disk writes.
For scale and availability, users are pinned to a home data center and their data replicated to others. If one data center fails, users are served from a replica. Stateless tiers scale freely; queues absorb spikes; metadata, cache and search all partition by user so capacity grows by adding nodes.
- Elasticsearch: easy to adopt, but a second system to keep consistent.
- Native LSM-based index: one store, tuned for heavy writes, but costly to build.
Failure handling and wrap-up
Queues are the main safety valve: if outgoing workers or a remote provider are slow, mail waits rather than disappearing, and growing queue depth triggers scaling or alerting. Retries with backoff handle transient SMTP errors, while an error queue captures messages that need attention. Durable replication across data centers protects stored mail. In summary: separate stateless and stateful tiers, partition everything by user, keep attachments in object storage, and treat deliverability as an ongoing operational discipline rather than a one-time configuration.
Key numbers
Key terms
- SMTP
- The protocol used to transfer email between mail servers.
- IMAP
- A retrieval protocol that keeps mail on the server and lets clients fetch parts on demand across many devices.
- MX record
- A DNS record naming the mail servers that accept email for a domain, ordered by priority.
- MIME
- A standard for encoding non-text content such as attachments inside an email message.
- Message-Id / In-Reply-To / References
- Headers that uniquely identify a message and link it to its parent and ancestors, enabling threading.
- IP warm-up
- Gradually increasing send volume from a new IP so receiving providers build trust in it.
- SPF / DKIM / DMARC
- Sender authentication standards covering authorised IPs, cryptographic signatures and the policy for failures.
- LSM tree
- A storage structure that buffers writes in memory and flushes them sequentially, ideal for write-heavy indexes.
Common mistakes
- Storing attachments inline in the metadata database instead of in an object store.
- Partitioning metadata by email ID, scattering a single user's inbox across every node.
- Filtering unread mail with a non-key predicate on a wide-column store instead of modelling separate tables.
- Sending from a brand-new IP at full volume and landing in spam folders.
- Delivering synchronously from the web tier so a slow remote MX server blocks user requests.
Further study
- Bigtable: A Distributed Storage System for Structured Data (2006)
- Cassandra: A Decentralized Structured Storage System (2009)
- The Log-Structured Merge-Tree (O'Neil et al., 1996)
- RFC 5321 (SMTP) and RFC 3501 (IMAP)
- Elasticsearch