advanced 3 min answer

Your analytics platform ingests change data from an operational estate that was split into 480 logical shards in 2021 and then spread across a larger number of physical database hosts in 2023, the arrangement Notion has described publicly. Your connectors, offsets and target tables are all configured per host. What does that physical move do to the downstream pipeline, and what design turns the next one into a configuration change rather than a re-ingest?

notionshardingcdcingestion contractbackfill
Show the full answer Hide the answer

The situation they were in

Notion's published account has two moves in it. In 2021 the product's Postgres was split into 480 logical shards, a number chosen to be far larger than the host count. In 2023 those logical shards were redistributed across more physical hosts. The point of the design is that growth is a move of shards between machines rather than a re-hash of keys, so the application's routing never changes. Notion later published, in 2024, a data lake built on change data capture into Hudi tables on object storage, driven by a change stream dominated by updates rather than inserts.

The question is what the 2023 move does to a downstream pipeline that was built against the wrong abstraction.

What a physical resharding does downstream

If the ingestion is keyed to hosts, doubling the host count means:

  • New replication slots with no history. Each new host has its own write-ahead log starting at the moment it was created. A connector pointed at it either re-snapshots the shards it now carries, which is a full re-ingest of that data, or starts from "now" and loses every change in the gap, silently.
  • Checkpoints that do not transfer. Offsets and watermarks recorded per host are meaningless on a different host, so there is no safe resume point for the shards that moved.
  • A partition column whose domain changed. If the target table is partitioned or tagged by source host, history is now split across two naming schemes and any query grouping by it is wrong for the transition period.
  • Duplicates or losses at the boundary, depending on whether the row keys are unique across the whole estate or only within a host.

What makes the next one a configuration change

Write every downstream contract against the logical shard and never against the host.

  1. Stream name, checkpoint key, and any partition or tag column derive from the logical shard id. A host is then a routing detail the pipeline never records.
  2. Primary keys are unique across the whole estate, so a row that moves hosts is an upsert on the same key rather than a new row.
  3. Snapshot and stream converge per logical shard, with the stream started before the snapshot, so a relocated shard resumes from its own recorded position.
  4. The host mapping lives in one configuration object that the connector reads, and changing it is a restart.

What it costs

Hundreds of logical shards means hundreds of small streams: per-stream connector overhead, hundreds of monitors, and a target table receiving many small files that needs real compaction capacity. You pay operational surface to buy a contract that does not move. That is a good trade at 480 shards and a poor one at 4.

When not to copy this, and where it would be a mistake

A company with one 200 GB Postgres should not pre-shard into 480 of anything; the abstraction costs more than the problem. The transferable rule is not the number, it is the question: which abstraction is my downstream contract written against, and how hard is that one to change? Choose the one you can over-provision cheaply. The second half of Notion's story transfers even less: the lake was built because their change stream is update-heavy, and an append-mostly workload needs none of the merge machinery that justified it.