skip to content

Cluster Shape & Capacity

How many nodes, how many partitions or queues, what disks, which failure domains and how clients reach them — the shape fixed before traffic arrives. Asked because those choices are painful to undo.

part ofBroker & streaming operationsoverview, primer and where to startread it →
on this pageshow

questions

25

Why is a messaging client configured with one entry address instead of the address of every node in the cluster?

level: juniorimportance: must knowfreq 62%

answer

  1. a way in, not a directory
  2. the cluster describes itself
  3. the list serves the first hop only
  4. later connections use announced addresses

basics

~20 s

The entry address only has to reach any one member, because the cluster then describes itself and hands back the per-node addresses later connections use. A hand-written list of every node duplicates membership the cluster already knows, and goes stale.

solid answer

~50 s

A client needs one working way in, not a directory. The entry address — one address, or a short list of two or three for redundancy — is used to open a first connection to any reachable member and ask the cluster to describe itself. On clusters whose work is owned by particular members, the answer is the per-node addresses each member announces, plus which member currently serves which part of the work; every connection after that goes to those announced addresses. So enumerating all nodes in client configuration buys nothing durable: the extra entries only improve the odds that the *first* hop lands somewhere alive, they do not change where traffic goes afterwards, and the list drifts out of step the moment the cluster gains or loses a member. Some platforms answer every operation at a single address, and there the question does not arise.

go deeper

for a junior

Remember the shape: one address to get in, then the cluster tells the client where everything is. A short list exists only so the first connection has a spare.

for a middle

Explain that the entry address serves one connection and the announced per-node addresses serve all the rest, and that the two fail with different symptoms — cannot start, against starts then stalls.

for a senior

Show that you use the symptom to pick the suspect: a client that opens a connection and then times out is an announced-address problem, and editing the entry list is wasted effort.

for a principal

Frame it as a duplication question. Membership copied by hand into client configuration is a second source of truth that nothing maintains; the estate rule should keep that copy as small as redundancy allows.

## What the entry address is actually for A broker or streaming cluster is several cooperating machines, and a client has to start somewhere. The address a client is configured with — the **entry address** — has exactly one job: open a first connection to *any* single reachable member and ask the cluster to describe itself. It is a way in, not a directory of the cluster. The answer to that first question is the cluster's own view of itself. On platforms where a stream is split into partitions and each partition's work is owned by a particular member, that view is the **advertised address** every member announces for itself, together with which member currently serves which part of the work. From that point the entry address has done its job, and every subsequent connection the client opens targets an advertised address the cluster supplied. ## Why a longer entry list buys nothing durable - The cluster already holds its own membership and will hand it over on request; a list written into client configuration is a **second copy** of that membership, and two copies drift. - Only the first hop consults the list. Adding a tenth address does not change where the eleventh request goes — that is decided by what the cluster announced. - If what the cluster announces is wrong or unreachable from the client, no length of entry list rescues it. The client will connect and then fail on everything after. - Membership changes. A node replaced next quarter has an address nobody edits back into a hundred client configurations, so the long list rots into a list that is mostly right and quietly wrong. - The list is not useless, though: if the single member you named happens to be down, the client cannot even begin. Two or three entries, sitting in **different failure domains**, make the first hop survivable. That is redundancy for one connection, not a membership register. ## Two addresses that are easy to confuse | | the entry address | the advertised address | |---|---|---| | who sets it | the client's own configuration | each node, in the cluster's configuration | | how many | one, or a short list of two or three | one per node, per client-facing view | | what uses it | the first connection of the connect step | every connection the client makes afterwards | | symptom when wrong | the client cannot start at all | the client starts, then fails on all work | That bottom row is the whole practical value of the distinction. "Cannot connect at all" and "connected, then nothing works" are two different faults with two different owners, and reading the symptom tells you which address to look at. ## Where designs genuinely differ This is not one model dressed in neutral words; platforms in this class really do differ. - On **split-stream** platforms, work belongs to a specific member, so the client must learn per-node addresses and keep connections to several members at once. The two-step connect is unavoidable. - On **shared-queue** platforms, competing readers drain one work list and any member may be able to serve the client, so the second step can be thin or entirely absent. - Some hosted and single-address offerings answer **every** operation at one published address, and the client never sees per-node addresses at all. There, one entry address is not a convenience, it is the whole story. A candidate who says "you always get a list of nodes back" is describing one design and asserting it as the class. ## What to do in practice 1. Configure two or three entry addresses, chosen so they are not all in the same failure domain, and stop there. 2. Do not treat the entry list as documentation of the cluster's membership; nothing keeps it true. 3. When a client can open a connection but cannot do work, ignore the entry list and look at what the cluster is announcing about its members. 4. When a client cannot open anything at all, the entry addresses — or the path to them from that client — are the first suspect. The underlying idea generalises beyond messaging: a system that knows its own membership should be asked for it, rather than having it copied by hand into every caller. The copy is always the thing that goes stale first.

  • If one entry address is enough in principle, why do teams configure two or three?
    Redundancy for the first connection only. If the single member you named is unavailable, the client cannot even begin — it has no other way to ask the cluster anything. Two or three entries in different failure domains make that first hop survivable. Beyond that, extra entries add nothing: they do not influence any later connection.
  • Does the client keep the entry address after the first connection succeeds?
    Yes, it keeps it, but it carries no work traffic. The client re-enters through it when it has no usable view of the cluster — at startup, after restart, or when it has lost contact with the members it knew. Routine work goes to the addresses the cluster announced, and clients refresh that view periodically rather than re-entering each time.
  • Is there any case where the entry address really is the only address a client ever uses?
    Yes. Some platforms, particularly hosted ones, answer every operation at a single published address and never expose per-node addresses. There the connect step has no second act, and the client's whole reachability story is that one address. Assuming that model — or assuming its opposite — is the usual source of confusion when people move between platforms.

It is the address of a building's front desk. You need one that works to get through the door; the desk then tells you which floor the person you want is on, and you walk there yourself. Memorising every room number in advance does not save you the walk, and the list is wrong as soon as someone moves desks.

saying these in an interview costs you the question

  • Thinks the client must list every node or it cannot connect
  • Believes all client traffic keeps flowing through the entry address
  • Treats the entry list as an accurate record of cluster membership
  • Assumes a longer entry list survives node replacement
  • Says every platform hands back per-node addresses
open as a page

Why does a broker's append path ask more of a volume's sustained throughput than of its seek performance?

level: juniorimportance: must knowfreq 58%

basics

~20 s

Broker storage is dominated by ordered appends, and most reads arrive moments later asking for the same bytes in the same order, so the volume is asked for steady sequential bytes per second rather than fast random seeks.

open as a page

Three copies of one partition sit on three record-serving nodes in the same rack — what does that placement protect against?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Node-level failures only: a disk, a process, one machine. The rack is a single power and network domain, so one rack event takes all three copies at once — three copies against one failure class, one copy against another.

open as a page

In a broker cluster, what work does the metadata role do that a record-serving node does not?

level: juniorimportance: must knowfreq 68%

basics

~20 s

A cluster splits two jobs. Record-serving nodes accept writes and serve reads for the partitions or queues they own. The metadata role holds membership, stream and partition metadata and configuration, and admits every structural change.

open as a page

A stream is split into a fixed number of partitions and a shared queue is not split at all - what caps reader parallelism in each?

level: juniorimportance: must knowfreq 70%

basics

~20 s

On platforms that split a stream into partitions, the count is the parallelism ceiling: each partition is normally served to one reader at a time, so extra readers idle. A shared queue has no such count; competing readers simply take more work.

open as a page

Nodes are split evenly across exactly two failure domains — why is a third domain still worth adding?

level: middleimportance: must knowfreq 57%

basics

~20 s

An even split leaves each survivor holding exactly half, and half is never a majority. Losing either domain therefore leaves no majority of the copy set or of the membership holding cluster metadata. A third domain makes one survivor decisive, even if it only breaks ties.

open as a page

Why is the membership that holds a cluster's metadata role kept small and odd-sized rather than spanning every node?

level: middleimportance: must knowfreq 58%

basics

~20 s

Every change to cluster state must be agreed by a majority of that membership, so each extra member raises the number that must agree and the work per change without adding proportional failure tolerance. Odd sizes avoid paying for a member that tolerates nothing extra.

open as a page

Why is a stream's partition count treated as a one-way door, when raising it later is routine and lowering it is not?

level: middleimportance: must knowfreq 62%

basics

~20 s

The change is asymmetric. Raising the partition count on a stream is a supported, routine operation; lowering it generally is not offered at all, so coming down means creating a second stream at the smaller count and moving writers and readers across.

open as a page

A stream takes 20,000 records per second averaging 2 KB and four independent reader sets read every record - what bandwidth must the cluster carry?

level: middleimportance: must knowfreq 62%

basics

~20 s

Rate times average record size gives the write side: 20,000 x 2 KB is about 40 MB per second in. Each independent reader set that reads every record adds the same again, so four make roughly 160 MB per second out.

open as a page

If the membership behind a cluster's metadata role loses the majority it needs while every record-serving node stays healthy, what continues and what stops?

level: seniorimportance: must knowfreq 60%

basics

~20 s

Structural change stops: no new streams, no configuration change, no partition moves, no recorded change of ownership. On many platforms traffic to partitions whose ownership is already settled keeps flowing, though some designs fail closed instead. The cluster is frozen in its current shape.

open as a page

After a client reaches the entry address, what does the cluster hand back, and which addresses do later connections use?

level: middleimportance: should knowfreq 56%

basics

~20 s

The cluster answers the first connection with its own view of itself: the address each member announces, and which member currently serves which work. Every later connection targets those announced addresses, not the address the client originally dialled.

open as a page

When a reader asks for a record written seconds ago, what decides whether the volume is touched at all?

level: middleimportance: should knowfreq 46%

basics

~20 s

Whether that record is still resident in memory somewhere: the operating system's file cache, the broker process's own memory, or neither. On designs that cache little, and for any reader that has fallen far behind, every such read goes to the volume.

open as a page

If raising a stream's partition count is routine, why is picking a very large number up front still a mistake?

level: middleimportance: should knowfreq 58%

basics

~20 s

Each partition carries fixed overhead on every node holding a copy of it, multiplied by the copy count and by every stream in the estate. A very large total also lengthens recovery and enlarges the metadata the coordination membership must agree on.

open as a page

A cluster sized from a daily average of 8,000 records per second sees 40,000 in its busiest ten minutes - which number should have sized it?

level: middleimportance: should knowfreq 55%

basics

~20 s

The busiest interval, not the daily mean. A peak-to-average ratio of five means the machines must carry 40,000 records per second, while the volume size still comes from the sustained rate accumulated over the retention span.

open as a page

A client connects to the entry address successfully, then every read and write times out — what explains that, and what do you change?

level: seniorimportance: should knowfreq 54%

basics

~20 s

The entry address is reachable from the client but the per-node addresses the cluster handed back are not, so every connection after the first goes nowhere. The fix is a second announced view whose addresses are reachable from where that client sits.

open as a page

Producers see write latency climb while the volume shows spare space and no errors — which storage decision should you suspect?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Suspect how the volume's speed was provisioned: a throughput allowance bought separately from size, or a volume shared with another workload. The append path is being made to wait, and waiting is not an error, a fault or a space alarm.

open as a page

Every partition showed a full copy set, yet one rack outage took partitions offline — how did copies end up sharing a domain?

level: seniorimportance: should knowfreq 46%

basics

~20 s

The layout was correct when built and drifted afterwards. Nodes replaced one at a time rejoined with a missing or wrong domain label, so the cluster placed copies as though every node were independent. Counts stayed full; the spread behind them did not.

open as a page

A single shared queue has no partition count to pick, so what decides how many competing readers it can usefully be planned for?

level: seniorimportance: should knowfreq 48%

basics

~20 s

The slowest resource every reader shares while handling a record - a database, a pool, an external service's allowed call rate. The queue imposes no count, so the useful number of competing readers is the one that saturates that resource and no more.

open as a page

Every record on a stream is kept on three record-serving nodes - why does the network between them carry more than the writers send?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Because a write is not finished when it is accepted. Where one node takes the write and sends it to the other two, internal traffic is about twice client ingress, and every stored copy also writes the full record bytes to its own volume.

open as a page

Would you standardise your broker estate on volumes local to each record-serving node or on network-attached volumes?

level: principalimportance: should knowfreq 42%

basics

~20 s

Neither wins outright. Local volumes give the highest and least variable throughput but die with the node; network-attached volumes survive the node and can be reattached, paying a network hop in the write path and a throughput allowance you must buy.

open as a page

What should an organisation's standing rule be for how many failure domains a cluster spans, and what does that rule foreclose?

level: principalimportance: should knowfreq 34%

basics

~20 s

Usually: every cluster spans three independent domains, with neither the copies of a partition nor the membership holding cluster metadata leaving a majority in any one. It forecloses single-rack clusters, two-domain locations without a tie-breaker, and any domain far enough away to sit in the write path.

open as a page

What does an operator gain or give up when the metadata role runs inside the cluster instead of in a separate coordination service or hidden by a provider?

level: principalimportance: should knowfreq 42%

basics

~20 s

An internal membership means one system to size, patch and reason about, but coordination competes with record traffic unless separated. A separate coordination service isolates cluster state at the cost of a second system with its own majority, upgrades and on-call. A hidden one removes both jobs and all visibility.

open as a page

As a platform lead, how would you set one default partition count for every new stream across a large estate?

level: principalimportance: should knowfreq 42%

basics

~20 s

Set a small default for the long tail of low-rate streams and publish one or two higher tiers a team opts into with a stated reader-parallelism need. A generous single default is multiplied by every stream anyone ever creates.

open as a page

A broker cluster is sized exactly to its measured peak rate - what routinely arrives that the machines have no room for, and what should the design target have been?

level: principalimportance: should knowfreq 44%

basics

~20 s

Three things reliably arrive on top of measured traffic: a reader replaying from the start, a writer retry surge re-sending records already accepted, and a dead node's share landing on its surviving peers. The design target is a fraction of what the machines can do, not all of it.

open as a page

Clients reach your clusters from inside the network, from a partner network and from laptops — what rule do you set for the addresses each cluster advertises, and what does each extra view cost?

level: principalimportance: nice to knowfreq 31%

basics

~20 s

The rule is that every client population gets a view whose addresses are usable from where it sits, and that the number of populations stays deliberately small. Each extra view is one more address per member to keep correct for the cluster's life.

open as a page