A reassignment plan moving a partition (or queue) to another cluster node is accepted in a second — why is the move not finished?
answer
- a declaration, not a completed move
- bytes still have to travel
- retained history, not just new records
- retained bytes over allowed bandwidth
- ownership settles after catch-up
basics
~20 sA reassignment plan only declares which node should hold which partition (or queue). Accepting it starts ordinary copy traffic that streams the unit's retained records to the new node, and ownership settles only once that copy has caught up.
solid answer
~40 sA reassignment plan is a declared mapping from each unit of ownership — a partition, or on queue-shaped brokers a queue — to the set of nodes that should hold it. Submitting it is a small metadata write, which is why it is accepted instantly. What it triggers is a background transfer: the destination node creates an empty copy and fetches the unit's retained records from a node that already holds them, while producers keep writing, so the target keeps moving. Only when the new copy's catch-up distance reaches roughly zero does it join the caught-up set and the old copy become droppable. The duration is retained bytes divided by the bandwidth the move is actually allowed, not the time the plan took to accept.
go deeper
Remember the one-line version: the plan is a declaration and the data still has to be copied. Being able to say that ownership changes only after the new copy catches up is most of the mark at this level.
Explain the mechanics: an empty copy is created, retained records are fetched while writes continue, catch-up distance closes, and only then does the old copy become droppable. Give the duration as retained bytes over allowed bandwidth.
Show that you plan around it. Say how you would size the transfer before starting, how you would read progress as remaining catch-up rather than elapsed time, and when you would judge a move unable to converge.
Frame it as a recurring estate cost. Every hardware change, growth step or rebalance of capacity buys another multi-hour copy operation, and that number is what makes a storage design where moves cost no bytes worth pricing.
## A plan is a declaration, not a transfer A **reassignment plan** is a declared mapping: for each *unit of ownership* — a partition, or on queue-shaped brokers a queue — it names the set of cluster nodes that should hold a copy of that unit. Submitting the plan writes that intent into the cluster's own metadata. It is a tiny write, it is accepted about as fast as any other metadata change, and it says nothing at all about where the records currently sit. The records are the point. A unit that has been taking traffic for a week is not a name; it is however many gigabytes of retained records are sitting on the disks of the nodes that hold it today. A plan naming a node that has never held that unit is a plan to put those bytes on that node's disk, and nothing but ordinary copy traffic can put them there. ## What the cluster does after it accepts the plan 1. The destination node creates an empty copy of the unit. 2. It begins fetching records from a node that already holds one, usually starting at the oldest record still inside the retention window and working forward. 3. Producers keep writing to the unit throughout, so the end the destination is chasing keeps moving. 4. When the destination's **catch-up distance** — how far behind the current leader its copy still is, in records, bytes or seconds — reaches roughly zero, that copy joins the **caught-up set**. 5. Only then does the plan's final state become real: a copy on a node the plan drops can be removed, and if the plan also changes which node serves the unit, that handover can happen. Steps 2 to 4 are where all the time goes, and they are exactly the part that is invisible if you only watch whether the plan was accepted. ## What actually sets the duration | Input | Why it matters | |---|---| | Retained bytes on the unit | The dominant term: this is what has to travel. | | Bandwidth the move is allowed | The **move throttle** — the bandwidth ceiling on the migration — divides into those bytes. | | Incoming write rate | The new copy has to out-run it, or catch-up distance never closes. | | Units moving concurrently | Parallel moves share the same source disks and the same links. | Two consequences surprise people. First, a busy unit can fail to converge at all: if the allowed copy rate sits below the rate at which new records arrive, catch-up distance grows rather than shrinks and the move runs forever. Second, adding nodes does not make one unit's move faster — more nodes let more units move in parallel, but a single unit's transfer is bounded by the two nodes at either end and whatever bandwidth is permitted between them. ## The unit stays available the whole time During the copy, the nodes holding the unit today keep serving it. Writers keep writing, readers keep reading, and the migration is background traffic layered on top. That is the good news, and it is also the entire problem: the copy shares the same disks, the same host cache and the same network links as the live traffic, which is why a bandwidth ceiling on the move exists at all. ## Reading progress honestly - Track **remaining catch-up** per unit and its trend, not elapsed time. A flat trend means stuck, not slow. - Compare bytes already transferred against bytes still to go; elapsed minutes tell you nothing about a move whose rate is being capped. - Expect the tail of a move to be the quick part: once the historical records are across, the copy only has to keep pace with live writes. - A move that has been "nearly done" for hours is usually a converge problem, not a slow disk. ## Where designs genuinely differ Not every platform pays this cost. On designs with **detached storage** — where a unit's records live on shared or remote storage rather than a node's local disk — changing which node owns a unit is a metadata change with no bytes in flight, and a move that takes hours elsewhere takes seconds. Some platforms keep a locally durable recent tail and offload only older records, so a move carries the tail and nothing more. The shape of the copy differs too: where one design has a follower stream records from a leader, another writes each record to a quorum of nodes at ingest and has no single follower stream to resume. The safely general statement, and the one to give in an interview, is that accepting the plan is instant, moving the data is not, and how much data has to move depends entirely on where the platform keeps it.
- How would you estimate, before starting, how long one unit's move will take?Take the retained bytes currently held for that unit and divide by the bandwidth the move is actually permitted after the ceiling and after competing traffic — not the link's rated speed. Then sanity-check convergence: if the incoming write rate is close to that allowed copy rate, the new copy never closes its catch-up distance and the estimate is meaningless until the ceiling is raised or the load drops.
- Does the unit go unavailable while its records are being copied?No. The nodes that hold it today keep serving reads and writes for the whole transfer, which is precisely why the copy is background traffic and why it competes with production for disks and links. The only moment anything changes hands is at the end, when the new copy has caught up and ownership is handed over — and that step is short.
- Why does a move often slow down rather than fail outright when capacity is tight?Because the transfer is rate-limited rather than transactional: it fetches what it can, whenever it can. Tight capacity or a low bandwidth ceiling simply shrinks the rate, so remaining catch-up falls more slowly, or stops falling once new writes arrive as fast as the copy proceeds. Nothing reports an error; the plan just sits in its intermediate state.
Filing a change-of-address form takes a minute and is accepted immediately; the van full of furniture still has to drive across the country, on the same roads everyone else is using. The paperwork is the reassignment plan; the van is the copy traffic.
saying these in an interview costs you the question
- Says the move is done because the plan was accepted.
- Assumes only newly written records travel, not the retained history.
- Thinks the new node can serve the unit the moment it is listed in the plan.
- Estimates the move from record rate instead of retained bytes.
- Believes adding cluster nodes speeds up a single unit's transfer.
- Assumes the unit is offline while its records are being copied.