A file storage and sync service — think Dropbox or Google Drive at a foundational level — looks like “upload a file, download a file” until multiple devices need to stay in sync, large files need to upload without redoing the whole thing on a dropped connection, and duplicate content across millions of users starts costing real money if left unaddressed.

Problem

We’re designing a service where users upload files, organize them into folders, share them with other users, and see changes sync automatically across every device they’re logged in on. The core design challenge is separating two very different kinds of data with very different scaling needs: the actual file bytes (large, immutable once written, rarely read after the first day) and the metadata describing files, folders, versions, and permissions (small, constantly read and written, needs strong consistency).

Treating both as one undifferentiated blob-and-record problem is the most common design mistake here — the right architecture pulls them apart early and lets each scale on its own terms.

Requirements

Functional

  • Upload and download files of varying size, from a few KB to several GB.
  • Organize files into a folder hierarchy, per user, with sharing (a folder shared with another user or a public link).
  • Sync changes across a user’s multiple devices — a file added on desktop shows up on mobile without a manual refresh.
  • Version history: recover a previous version of a file after it’s been overwritten.
  • Resume interrupted uploads without re-transferring already-sent bytes.

Non-functional

  • Durability first: a file the user believes is saved must not be lost, full stop — this is the property users trust the product for, above speed or features.
  • Deduplicate identical content across users to control storage cost, without leaking the existence of another user’s identical file.
  • Upload/download throughput should scale independently of metadata operation throughput — browsing a folder shouldn’t be slower because someone else is uploading a large video.
  • Sync latency (change on device A visible on device B) under a few seconds for small files/metadata changes; large file transfer time is bounded by bandwidth, not by the system’s own overhead.

Capacity Estimation

Assume 10M users, averaging 50 GB stored each (a realistic blend of light and heavy users), and 200,000 uploads/hour at peak.

  • Total storage: 10,000,000 × 50 GB ≈ 500 PB in principle, but real-world dedup and the fact that most users are well under the average brings actual stored bytes down substantially — we’ll use the stated 500 TB unique-content scale target as the practical planning number after estimating meaningful dedup savings from common files (OS files, popular documents, stock photos) offsetting outliers with genuinely large unique libraries.
  • Upload QPS at peak: 200,000/hour ÷ 3,600 ≈ 56/s for full-file upload initiations; each large file is chunked (see Deep Dive) into a handful of pieces, so actual chunk-upload QPS runs several times higher, in the few-hundred/s range.
  • Average file size: assume a blended average of 4 MB across all uploads (many small documents and photos, some large videos pulling the average up). At 56 uploads/s × 4 MB ≈ 224 MB/s ≈ 1.8 Gbps sustained ingest bandwidth at peak, which is a meaningful but entirely manageable number for a modern object storage backend and CDN-fronted architecture.
  • Metadata operations: every upload, rename, move, share, and folder listing hits the metadata store. Assume metadata ops run roughly 20x file operations (users browse and check status far more than they upload) — at 56/s uploads that’s ~1,100 metadata ops/s at peak, small enough that metadata service capacity is rarely the constraint, but latency-sensitive enough that it still needs its own dedicated scaling story (see Scaling).
  • Chunk size and count: with 4 MB average chunks, a 500 TB unique-content store holds roughly 500,000,000 MB ÷ 4 MB ≈ 125M chunks, each needing a metadata row (chunk hash, size, storage location) — a modest, well-indexable table by any modern database’s standards.
Back-of-the-envelope
Users / avg stored
10M / ~50GB
before dedup
Unique content stored
~500TB
after dedup savings
Upload QPS (peak)
~56/s files, ~200-300/s chunks
200K uploads/hour
Peak ingest bandwidth
~1.8 Gbps
56/s x 4MB avg file
Metadata ops (peak)
~1,100/s
~20x file op volume

API Design

POST /v1/uploads HTTP/1.1
Content-Type: application/json

{
  "fileName": "quarterly-report.pdf",
  "folderId": "fld_a91c2",
  "fileSize": 18874368,
  "contentHash": "sha256:9f3a...c012"
}

HTTP/1.1 200 OK
Content-Type: application/json

{
  "uploadId": "up_7c4e21",
  "chunkSize": 4194304,
  "existingChunks": ["chunk_a1", "chunk_b7"],
  "uploadUrls": [
    "https://storage.example.com/put/up_7c4e21/chunk/2",
    "https://storage.example.com/put/up_7c4e21/chunk/3"
  ]
}
PUT /put/up_7c4e21/chunk/2 HTTP/1.1
Content-Length: 4194304
Content-Type: application/octet-stream

<binary chunk data>

HTTP/1.1 200 OK
POST /v1/uploads/up_7c4e21/complete HTTP/1.1

HTTP/1.1 201 Created
Content-Type: application/json

{
  "fileId": "file_9e102c",
  "versionId": "ver_1",
  "status": "committed",
  "syncedAt": "2026-08-25T09:14:02Z"
}
GET /v1/sync/changes?since=cursor_88f2a1 HTTP/1.1

HTTP/1.1 200 OK
Content-Type: application/json

{
  "changes": [
    { "fileId": "file_9e102c", "type": "created", "versionId": "ver_1", "path": "/reports/quarterly-report.pdf" }
  ],
  "nextCursor": "cursor_88f2b4",
  "hasMore": false
}

The upload-initiation response returning existingChunks is the API surface for deduplication: the client hashes chunks locally, tells the server what it’s about to send, and the server says which chunks it already has (from any user) so the client skips re-uploading them entirely.

Data Model

CREATE TABLE files (
  id           BIGINT PRIMARY KEY,
  owner_id     BIGINT NOT NULL,
  folder_id    BIGINT NOT NULL,
  name         VARCHAR(255) NOT NULL,
  current_version BIGINT NOT NULL,
  created_at   TIMESTAMPTZ NOT NULL DEFAULT now(),
  deleted_at   TIMESTAMPTZ
);

CREATE TABLE file_versions (
  id          BIGINT PRIMARY KEY,
  file_id     BIGINT NOT NULL REFERENCES files(id),
  size_bytes  BIGINT NOT NULL,
  created_at  TIMESTAMPTZ NOT NULL DEFAULT now(),
  content_hash VARCHAR(64) NOT NULL
);

CREATE TABLE version_chunks (
  version_id  BIGINT NOT NULL REFERENCES file_versions(id),
  chunk_index INT NOT NULL,
  chunk_hash  VARCHAR(64) NOT NULL,
  PRIMARY KEY (version_id, chunk_index)
);

CREATE TABLE chunks (
  hash         VARCHAR(64) PRIMARY KEY,
  size_bytes   INT NOT NULL,
  storage_key  VARCHAR(255) NOT NULL,
  ref_count    INT NOT NULL DEFAULT 1
);

CREATE INDEX idx_files_folder ON files (folder_id, name);
CREATE INDEX idx_versions_file ON file_versions (file_id, created_at DESC);

The separation between file_versions (what a version is, logically) and version_chunks plus chunks (what bytes it’s made of, physically, shared across versions and users) is the load-bearing structure of the whole schema. A new version that only changes a few bytes in a large file reuses almost all of its version_chunks rows from the previous version — chunks.ref_count tracks how many versions across the whole system point at a given chunk, and a chunk is only deleted from physical storage once its ref count hits zero.

Table Access pattern
files / file_versions Point lookups and folder-scoped range scans, strongly consistent
chunks Point lookup by hash on every upload (dedup check), extremely read-heavy
version_chunks Read on download/reconstruction, write once per version created

High-Level Architecture

High-level architecture

The metadata service and blob storage gateway are entirely separate services with independent scaling paths, as the Problem section sets up. A client asks the metadata service what to do (create a file record, list a folder), and separately streams actual bytes to the blob gateway, which checks the chunk index for dedup before ever touching the underlying object store.

Deep Dive

Chunking

Files are split into fixed-size chunks (4 MB is a reasonable default — large enough to keep per-chunk metadata overhead low, small enough that a dropped connection loses at most one chunk’s worth of progress) before upload. Each chunk is hashed (SHA-256) client-side, and the chunk hash — not the file — is the unit of both deduplication and resumable upload.

Deduplication

Because chunks are content-addressed (keyed by hash, not by uploader or file), two different users uploading byte-identical files — extremely common for OS installers, popular PDFs, stock assets — share the same underlying chunk in chunks/storage_key, with ref_count tracking how many version_chunks rows point at it. The upload-initiation API tells the client which chunks the server already has, so identical content is never re-transferred over the network, not just never re-stored — this saves both storage and bandwidth.

Decision

Content-addressed chunk storage with server-side existence check

Client hashes locally, server confirms which hashes already exist before any bytes are sent

An alternative is server-side dedup after full upload (accept the bytes, hash them, discard if a duplicate exists) — simpler to implement, but wastes exactly the upload bandwidth we’re trying to save, which defeats much of the point for large files. Checking before transfer costs one extra round trip on the upload-initiation call but is trivial next to the bandwidth saved on a multi-gigabyte duplicate.

A privacy nuance: telling a client “we already have this chunk” reveals, in principle, that some other user has uploaded byte-identical content — this is a real but narrow leak (an attacker would need to already know or guess the exact bytes to exploit it, since the check is by hash, not by content guessing) and is a widely accepted trade-off in production systems, but it’s worth stating explicitly rather than discovering during a security review.

Metadata and sync

The metadata service is a fairly conventional strongly-consistent relational store — file/folder hierarchy, permissions, version pointers — because correctness here (not showing a user a stale folder listing, not losing track of a share permission) matters more than raw throughput, and the volume (per Capacity Estimation, roughly 1,100 ops/s at peak) doesn’t demand exotic scaling.

Sync works via a per-user change log with a cursor, structurally identical to the pattern in the chat system’s message catch-up: each device tracks the last cursor it’s seen, and on reconnect (or via a long-lived push channel) asks “what changed since this cursor,” getting a durable, ordered answer rather than needing to diff entire folder trees. This is what makes catching up after being offline for a week and getting a live update while connected the same code path, rather than two separate mechanisms.

Resumable uploads

Because chunking already breaks a large file into independently-uploadable pieces, resumability falls out naturally: if a connection drops mid-upload, the client re-queries /v1/uploads/{uploadId} for which chunks are confirmed-received and resumes from the first missing one, rather than needing a separate resumable-upload protocol layered on top.

Failure Handling

Failure Detection Response
Chunk upload fails mid-transfer Client-side timeout / connection error Client retries that one chunk only; no data loss beyond the single in-flight chunk
Object storage write succeeds but metadata commit fails Transaction/timeout mismatch between blob gateway and metadata service Orphaned chunk stays in storage (harmless, low cost) until a garbage-collection sweep finds it has zero references and reclaims it; the version is simply never marked committed, so the client retries the complete call
Metadata database primary down Health check / replication lag Promote replica; brief write unavailability for new uploads/renames, existing downloads (served from object storage directly, not metadata) are unaffected
Dedup chunk index unavailable Timeout on hash lookup Skip the dedup check and upload the chunk anyway (treated as new); costs some redundant storage temporarily, never blocks the upload — correctness of file storage takes priority over the storage-cost optimization
Sync notification service down Push channel disconnects Client falls back to polling /v1/sync/changes on an interval; sync latency degrades from seconds to that poll interval, nothing is lost since the change log itself is durable
Ref count reaches zero incorrectly (race between concurrent version deletes) Reconciliation job cross-checks ref_count against actual version_chunks references Garbage collection sweep double-checks live references before physically deleting a chunk, never trusting ref_count alone for an irreversible delete

Scaling

Blob storage scaling

This is the easy part precisely because we separated it from metadata: object storage (S3-compatible or equivalent) is designed to scale to exabytes with the storage provider handling replication, and chunk content-addressing means the blob gateway is stateless and horizontally scalable — any instance can serve any chunk, since chunks are looked up by hash, not routed by any notion of ownership.

Metadata database scaling

Sharded by owner_id (or, better, by a hash of folder_id for workspaces with shared team folders), since almost every metadata query is naturally scoped to one user’s or one team’s data — folder listings, version history, and permission checks rarely need to join across unrelated users’ data. This mirrors the same sharding logic used in the chat system’s conversation-scoped data.

Dedup index scaling

The chunks table, keyed by hash, shards trivially and evenly by hash prefix (hashes are already uniformly distributed by construction), so this never develops the kind of hot-key problem that owner-scoped sharding can when one user or team is unusually active — every chunk lookup is a clean, evenly distributed point read.

CDN for downloads

Frequently accessed files (a shared public link that goes viral, a commonly referenced team asset) benefit from a CDN layer in front of object storage for downloads, since download traffic for popular content is exactly the kind of read-heavy, cacheable pattern CDNs are built for — this offloads the bulk of download bandwidth away from origin storage entirely.

Trade-off

Fixed-size chunking

  • Simple, fast, cheap to compute client-side
  • Dedup breaks entirely for a file with a small edit in the middle — every chunk boundary after the edit shifts
  • Fine for append-only or rarely-edited large files (video, backups)

Content-defined (rolling-hash) chunking

  • Preserves dedup across small in-place edits
  • More CPU-intensive client-side to compute rolling hash boundaries
  • Better fit for frequently-edited documents

Recommendation — Start with fixed-size chunking — it's simpler and covers the common case (most files are uploaded once, rarely edited in place) — and add content-defined chunking later specifically for document-editing workloads if dedup rate data justifies the added client CPU cost.

Observability

  • Dedup hit rate (fraction of upload-initiation chunk checks that come back “already exists”) — the direct measurement of whether the storage-cost optimization this design leans on is actually paying off in production, and a sudden drop can indicate a client-side hashing bug rather than a genuine change in content uniqueness.
  • Upload completion rate (initiated versus completed) — a gap here indicates client-side upload failures or abandoned uploads, and is the leading indicator for user-facing “my file didn’t save” complaints before they turn into support tickets.
  • Sync cursor lag per user/device — how far behind a device’s last-seen cursor is from the current change-log head, which is the direct measurement of real-world sync latency, not just a synthetic test of it.
  • Orphaned chunk count from the garbage collection sweep — a rising trend suggests a systemic issue in the metadata-commit path (chunks being written but not consistently referenced), worth investigating before it becomes a meaningful storage-cost problem.

Every upload session carries a trace id from initiation through chunk uploads to completion, so a specific user’s “my upload got stuck” report can be traced through the full multi-request flow rather than reconstructed from separate chunk-level logs.

Trade-offs

The foundational trade-off in this design is separating metadata from blob storage entirely, rather than storing file bytes alongside their metadata in one system. This means every operation touches two systems instead of one, and keeping them consistent (a version record pointing at chunks that genuinely exist and are complete) requires the two-phase commit-style flow described in Failure Handling. The payoff is that each system scales on the axis it actually needs to — metadata for consistency and low-latency small reads/writes, blob storage for raw capacity and throughput — and neither has to compromise its design to accommodate the other’s very different access pattern.

The second trade-off is content-addressed deduplication itself, which trades a narrow, hash-gated information leak (see Deep Dive) and real implementation complexity (ref counting, careful garbage collection) for substantial storage cost savings at scale. For a consumer-facing product with a large user base and heavily overlapping common content, this trade is almost always worth it; for a system storing exclusively unique, sensitive per-user content with no realistic overlap, the complexity might not earn its keep.

Final Architecture

Final architecture with chunk dedup and sync change log

The finished system’s core decision — pulling metadata and blob storage apart into independently scaled services — pays off throughout: dedup and chunk indexing live entirely on the storage side and never touch the consistency-sensitive metadata path, sync is a durable, cursor-based change log rather than a diffing exercise, and garbage collection runs as a conservative, reference-verifying background process that never risks the one property (durability) this whole system exists to guarantee.