advanced 3 min answer

Vitess - the MySQL sharding layer built at YouTube and later a CNCF project - reshards a live keyspace by having each target shard build its tables with a copy phase and then follow the source's binlog from the position where the copy ended. The same stream is exposed to applications as VStream. What does that design buy and where would copying it be a mistake?

youtubevitessvreplicationreshardingbinlog
Show the full answer Hide the answer

The situation they were in

Sharding behind a proxy means the shard map must be allowed to change. "Stop writes, dump, load, point the app at the new shards" is available to a small system and not to a large one, and the whole value of putting a routing layer in front of MySQL evaporates if resharding needs a maintenance window.

What they chose

Table, then stream, then table. Vitess's documentation on the life of a VReplication stream describes the shape: each target primary runs one stream per source shard whose key range overlaps its own. A stream given a GTID starts straight in binlog streaming; a request for a whole table starts in copy mode, because the older binlogs that would be needed to build the table from the log alone may no longer exist. The copy proceeds per table in cycles of copy, catch-up and fast-forward, and once replication lag is small the stream stays in binlog streaming. The same machinery feeds either another tablet during resharding or a gRPC stream for VStream clients, and if the target shard fails over the new primary resumes from where the previous one stopped.

The interesting part is the handoff, not the copy. The copy is a snapshot, the binlog is the changelog, and the join between them is a recorded position - which is why the copy is interleaved with catch-up rather than being one long dump. That is stream-table duality used as an operational mechanism rather than as a mental model.

Why it fit their constraints

Row-based binlogs already existed for replication, so the change stream cost nothing new. Because a stream resumes from a stored position, a target failover does not restart hours of copying. And because the same stream has an external API, the change-capture capability that was built for migrations becomes available to applications without a second system.

What it cost them

Consumers of that stream are coupled to the physical table shape: it is row-level change capture, not a domain event. The source must retain binlogs for longer than the slowest copy. And the operational surface is a workflow state machine with per-stream lag and per-table copy progress to watch, which is a thing an on-call engineer has to learn.

A planning number worth carrying: a 2 TB table copied at 50 MB/s is about 11 hours, so 24-hour binlog retention is already too tight if the copy stalls once.

When not to copy it

  • Not for moving one table once. A dump and a maintenance window cost hours. This costs a workflow engine, and building that to avoid a Sunday-morning window is the wrong trade.
  • Not as a team-to-team integration contract. It is the right mechanism for a data-movement tool and the wrong contract for a service boundary; publishing it turns a physical schema into a public API that nobody can refactor.
  • Not without measuring apply throughput against peak write rate first. If the source writes faster than the target applies, the catch-up phase never converges and the migration cannot finish. That single comparison decides whether the plan is feasible at all.
  • Not without a cutover plan. Any handoff of this shape needs a moment when the target is provably caught up and no new writes are landing on the source, however brief. If the system cannot tolerate that moment, the migration needs a different design, not a faster copy.