On this page
Tracks

System Design — Distributed Email Service

Last reviewed 11 Sept 2026

Part of the system design series. See the framework and building blocks first if you haven’t.

1. Requirements

Functional

  • Send email: compose, attach files, deliver to internal or external (other providers’) mailboxes.
  • Receive email: accept inbound SMTP, run spam/virus checks, land it in the right mailbox/folder.
  • Read/search: list a mailbox, thread conversations, full-text search across a user’s mail.
  • Standard mailbox operations: labels/folders, read/unread, delete, spam/not-spam feedback.

Non-functional

  • Massive scale: billions of users, tens of billions of messages/day at Gmail scale.
  • Durability of stored mail is non-negotiable — losing someone’s inbox is a catastrophic failure, not a degraded experience.
  • Delivery must be reliable across an adversarial, decentralized network (other mail servers, spammers, misconfigured senders) via SMTP, which long predates modern reliability expectations.
  • Search must return results in well under a second across a mailbox that can hold years of history.
  • Sender reputation must be actively protected — a shared sending infrastructure (many users’ outbound mail) is one bad actor away from getting the whole IP range blacklisted.

2. High-level architecture

flowchart TD
subgraph Outbound
  U1[User composes] --> SUB[Submission API]
  SUB --> OQ[[Outbound queue]]
  OQ --> SEND[Sending MTA<br/>DKIM/SPF sign]
  SEND -->|SMTP| EXT[External mail server]
  SEND --> BOUNCE[Bounce/complaint handler]
end
subgraph Inbound
  EXT2[External sender] -->|SMTP| MX[Inbound MTA / MX servers]
  MX --> SPAM[Spam + virus filter]
  SPAM --> ROUTE[Mailbox router]
  ROUTE --> STORE[(Mail store,<br/>sharded by user)]
  STORE --> IDX[(Search index)]
end
U2[User reads mail] --> API[Mailbox API]
API --> STORE
API --> IDX
Inbound and outbound paths are largely independent; both funnel through spam/reputation checks
  • Submission vs. sending MTA split: the API that accepts a user’s outgoing message is separate from the MTA (mail transfer agent) that actually does SMTP delivery to the recipient’s server — this lets the submission path return fast (durably queue, respond “sent”) while delivery, retries, and bounce handling happen asynchronously and can take minutes to fail over multiple retries.
  • Inbound MX servers are the public-facing SMTP endpoint (what a DNS MX record points to). They accept mail from any sender on the internet, which is why spam/virus filtering has to sit immediately behind them, before anything touches a user’s mailbox.
  • Storage sharded by user: since data access patterns are almost entirely per-user (nobody queries across mailboxes except spam-model training), horizontal partitioning by user_id is natural and lets the system scale by adding shards, with no cross-shard transactions needed for the common case.

3. Storage & search trade-offs

ConcernOptionTrade-off
Mailbox storageRelational DB (per-shard)Simple, transactional, doesn’t scale past a point without heavy sharding work
Mailbox storageWide-column store (Cassandra/Bigtable-style)Horizontally scales cleanly, matches per-user partition key, weaker cross-row transactions (rarely needed here)
AttachmentsInline in mail storeSimple, but bloats the hot row/table with large binary blobs
AttachmentsBlob store (S3-style) + reference in mail storeKeeps the mail store small and fast; attachment fetch is a separate, cacheable path
Full-text searchSeparate index (Elasticsearch)Fast to stand up, good for small-to-mid scale, but now you have two systems to keep consistent
Full-text searchNative index embedded with storageWhat Gmail/Outlook-scale systems build custom — avoids the dual-write consistency problem and the cost of syncing a whole side index for billions of mailboxes

4. Deep dive — reliable delivery and idempotency under SMTP’s constraints

SMTP itself gives you almost no reliability guarantees — it’s a decades-old, best-effort, store-and-forward protocol, so all the reliability engineering has to happen at your layer, on top of it:

  • Submission is durable-queue-then-ack: the moment a user hits send, the message is written to a durable outbound queue and the API returns success — the user doesn’t wait on actual SMTP delivery to a possibly-slow external server, which can legitimately take many seconds or fail and need retries.
  • Delivery retries with backoff: a receiving server being temporarily unavailable (4xx SMTP response) triggers retry with exponential backoff over a bounded window (classically up to ~4–5 days before giving up, per RFC convention) — after which the sender gets a bounce notification. A permanent failure (5xx, e.g. “no such mailbox”) bounces immediately, no retry.
  • Idempotent inbound processing: the same inbound message can legitimately be retried/resent by an upstream relay. Dedup on the Message-ID header (or a hash of it) before it lands twice in a mailbox — same idempotency principle as any at-least-once delivery system.
  • DKIM/SPF/DMARC on the outbound path: every outbound message is cryptographically signed (DKIM) and sent from IPs the domain has authorized (SPF), with a DMARC policy telling receivers what to do if those checks fail. This isn’t optional at scale — receiving providers (Gmail, Outlook) increasingly reject or spam-bucket unauthenticated mail outright.

5. What real systems do today

  • Amazon SES enforces concrete reputation thresholds: keep bounce rate under 5% and complaint rate under 0.1%; hitting those triggers an account review, and sustained rates near 10% bounces or 0.5% complaints can get sending paused entirely. This is a directly citable, real numeric policy for “how does a real system police sender reputation.”
  • SES’s architecture for bounce/complaint handling is event-driven: a configuration set publishes Bounce and Complaint events to an SNS topic, which fans out to a Lambda or SQS consumer that updates the sender’s local suppression list — the standard pattern is to suppress future sends to an address after a hard bounce rather than keep retrying and damaging reputation further.
  • Gmail/Outlook-scale architectures (per multiple 2025–2026 engineering write-ups on the topic) horizontally partition storage by user since access patterns are independent per user, replicate across data centers, and route users to a geographically closer mail server; database layers at this scale tend to be purpose-built rather than off-the-shelf, optimized specifically for the high-IOPS, per-user-partitioned access pattern mail represents.
  • Spam filtering at scale is a layered pipeline, not a single classifier: connection-level checks (sender IP reputation, rate limiting) before the message body is even fully received, then content/ML-based classification, then user-feedback loops (spam/not-spam clicks) that retrain the model — rejecting or bucketing as early in the pipeline as possible is cheaper than processing a full message just to discard it.

6. Scaling & failure

  • Inbound MX servers get overwhelmed during a spam wave → connection-level rate limiting and greylisting (temporarily 4xx-ing unfamiliar senders, which legitimate MTAs retry but most spam bots don’t) at the edge, before spend goes into full content filtering.
  • A single user’s mailbox shard is a hotspot (a mailing-list address, a support inbox receiving huge volume) → this is a partitioning skew problem; either sub-shard the hot mailbox by time range, or dedicate more resources to that shard specifically rather than resharding everyone.
  • Search index falls behind mail store (indexing lag after a write spike) → search shows stale/missing recent mail. Fix: treat indexing as an async queue consumer with its own backlog metric and alerting, and make sure the mailbox API can still serve unindexed-but-stored mail via a direct (slower) store scan as a fallback for very recent messages.

What happens when the sending MTA / outbound path dies: user-facing submission should stay up — the durable outbound queue is exactly what absorbs this. Messages queue and simply get delivered late once the MTA layer recovers; users see their message as “sent” (submission succeeded) but real delivery is delayed, which is the correct trade-off since submission and delivery are different guarantees. If the outage is prolonged, the queue itself becomes the thing to monitor (depth, oldest-message-age) and you may need to shed load by pausing high-volume automated senders first, protecting delivery latency for interactive human-sent mail.

What happens when the mail store (durable storage) has a partial outage: this is the one place “fail open” is not acceptable — reads for the affected shard should degrade to “temporarily unavailable” rather than return incomplete or stale-without-warning data, and writes (new inbound mail) should queue durably upstream rather than be accepted and dropped. Because data loss here is irreversible and directly visible to the user (a missing email), the durability guarantees (replication factor, backup/restore, point-in-time recovery) matter more here than almost anywhere else in the system.

Interview follow-ups

  • “SMTP itself is unreliable. How do you build a reliable send experience on top of it?” — Durable queue-then-ack at submission, async delivery with bounded retry/backoff, explicit bounce notification on permanent failure — the user’s “sent” and the actual delivery are decoupled guarantees.
  • “How do you stop the same inbound message from being stored twice if an upstream relay retries it?” — Dedup on Message-ID (or its hash) before it’s written to the mail store — same idempotent-consumer pattern as any at-least-once pipeline.
  • “Why would you build a custom search index instead of just using Elasticsearch?” — Scale-dependent: at Gmail scale, a fully separate index as a second source of truth for billions of mailboxes multiplies storage and consistency-maintenance cost; below that scale, Elasticsearch alongside the store with eventual consistency is the pragmatic choice.
  • “What actually gets an outbound sending IP blacklisted, and how do you prevent it?” — Sustained bounce/complaint rates above provider thresholds (SES: 5% bounce / 0.1% complaint triggers review); prevent it with address validation before send, immediate suppression-list additions on hard bounce, and DKIM/SPF/DMARC so receivers trust the mail is legitimately from you.
  • “Where does spam filtering happen, and why does the order matter?” — Cheapest checks first (connection/IP-level rate limiting and reputation) before expensive checks (full content/ML classification) — rejecting early avoids paying processing cost on mail you’re going to discard anyway.
  • “A user’s mailbox shard is getting hammered by a mailing list. What do you do?” — Diagnose as partition skew, not a global capacity problem; sub-shard by time or dedicate resources to that specific shard rather than rebalancing the whole system.
  • “Mail store has a partial outage. Do you fail open or closed?” — Fail closed for the affected shard — return “temporarily unavailable” rather than risk silently dropping or serving stale data for something as durability-sensitive and user-visible as email.

Sources: Amazon SES Reputation metrics — AWS docs · How to improve email sender reputation with Amazon SES Email Validation — AWS Messaging Blog · How to Handle SES Bounces and Complaints with SNS — OneUptime · Understanding Gmail Architecture — DhiWise · High-Level Design for Gmail — InterviewReady, Medium · Building Scalable Email Processing at Enterprise Scale — GDSKS, Medium