Table of Contents
Table of Contents
Diskless Kafka partition scaling is straightforward to describe but difficult to operate: add broker capacity without moving durable bytes, while keeping partition ownership, cache behavior, metadata, and client traffic within known limits. A bucket with free space does not prove that a partition move is safe. The decision depends on which bytes move, which bytes stay shared, and what the control plane must rebuild when a broker or storage path is impaired.
The practical starting point is to separate the workload into three things:
- Partition metadata: leaders, replicas or ownership records, offsets, and controller state.
- Compute: request handling, protocol work, scheduling, cache management, and network processing.
- Durable bytes: the retained stream data and the write-ahead path that makes an acknowledgement durable.
The central thesis is simple: partition scaling becomes predictable when these three dimensions are measured independently. A diskless architecture can change the data-movement part of the operation, but it does not remove the need to test metadata pressure, cache churn, object-storage requests, or recovery behavior.
1What partition scaling means in a diskless Kafka design
In a local-disk Kafka cluster, a partition assignment has a physical consequence. A broker is responsible for serving the partition and for retaining its local log segments. Adding a broker can therefore require moving partition data, replicating it, flushing it, and managing the extra network and disk work while the cluster is already busy.
A diskless design changes the physical consequence, not the logical contract. Producers still address topics and partitions. Consumers still track offsets and fetch records. Controllers still coordinate ownership and metadata. What changes is where durable stream data lives and how a broker obtains the bytes it serves.
That distinction gives operators a more useful scaling question:
When a partition changes ownership, are we moving durable bytes, changing metadata, warming a cache, or doing all three?
Answer the question for each implementation and workload. Treat a partition move as a measured workflow with an input, an observable transition, and a stop condition.
| Scaling dimension | What changes | Evidence to collect |
|---|---|---|
| Metadata | Ownership, leader or assignment records, offsets, and controller work | Metadata update latency, controller queue, quorum health, and recovery time |
| Compute | Requests, active partitions, protocol operations, and scheduling | CPU, memory, request latency, network throughput, and throttling |
| Durable bytes | Local segments, shared objects, write-ahead log (WAL) records, or a combination | Bytes written, uploaded, fetched, replayed, and retained |
| Cache state | Hot data, historical ranges, and prefetch state | Hit ratio, evictions, warm-up time, and cold-read latency |
| Network path | Client traffic, broker-to-storage traffic, and control traffic | Bytes and latency by zone, endpoint, region, and failure domain |
This table is also a migration boundary. A capacity plan that reports only broker count leaves the most important questions unanswered: whether another broker must copy data, whether a cold consumer will evict hot data, and whether object-storage request limits become the actual bottleneck.
For a neutral baseline, document the current cluster first. Record partition distribution, producer and consumer rates, peak shape, retention policy, replay behavior, compression, client acknowledgement settings, and the recovery objective. Keep averages and peaks separate. A scale-out triggered by a short burst is governed by the burst path, while retained bytes follow the longer storage window.
2The mechanism: brokers, cache, metadata, and object storage
A diskless Kafka data path has at least two lanes. The first carries Kafka requests and control signals. The second carries durable stream data through a write buffer or WAL and object storage. Operators need visibility into both lanes because a healthy request path can hide an upload backlog, and a healthy bucket can hide a controller or cache bottleneck.
The write path can be reasoned about as a sequence:
- A producer sends records to the partition leader.
- The broker appends the records to the configured durable write path.
- The broker acknowledges according to its durability contract.
- The system uploads, compacts, or indexes data in shared object storage.
- A replacement broker reads the required ranges and rebuilds serving state.
Each step has a different signal. Acknowledgement latency reflects the write path and durability choice. Upload lag reflects the path to object storage. Recovery time reflects how quickly metadata, cache, and durable ranges become available to the replacement broker. Combining these into one “broker latency” number hides the decision you need to make.
Reads have the same separation. Tailing consumers usually ask for the newest data, while catch-up consumers request older ranges. A replay may drive object-store reads and cache fills without increasing producer traffic. Measure these classes separately and decide whether they share cache or bandwidth budgets.
Metadata deserves its own capacity budget. Partition count affects controller records, ownership changes, offset tracking, ACL evaluation, and the amount of state a broker must load before it can serve traffic. More partitions do not automatically mean more durable bytes, and more retained bytes do not automatically mean more metadata. The two dimensions can scale at different rates, which is why a single “partitions per broker” limit is an incomplete readiness test.
A useful production dashboard groups signals by mechanism:
- Ownership and control: assignment changes, controller latency, quorum health, metadata load time, and failed transitions.
- Request handling: produce and fetch latency, active connections, CPU, memory, and network throughput.
- Durable path: WAL or buffer queue, flush latency, upload lag, object-store request latency, retries, and errors.
- Read path: tailing latency, catch-up latency, cache hit ratio, evictions, and consumer lag.
- Recovery: bytes to replay, metadata rebuild duration, time to resume normal traffic, and client error rate.
This is also where the term “Kafka on S3” needs precision. Object storage can be the durable layer, but it does not make request paths identical. Object size, request rate, cache policy, endpoint placement, and the chosen WAL or write buffer shape the cost and latency envelope. A design review should ask which data is served from memory, which data is fetched from object storage, and which path carries an acknowledgement.
Kafka Tiered Storage is a useful comparison point, but it answers a different storage question. A tiered design commonly keeps an active local tier and moves older log data to a remote tier. A shared-storage design changes the ownership boundary for durable data, so a broker replacement does not automatically imply copying the full retained log to a new local disk. KIP-1150 describes a diskless-topic direction in the Kafka community; verify proposal status and supported behavior against the Kafka release under evaluation.
3Failure, cost, and compatibility checks
Scaling is safe only when failure behavior is part of the plan. Test the path that matters to the application, not only the control-plane API that starts the operation. A scale-out can succeed from the controller’s perspective while consumers experience cache misses, object-store retries, or a temporary loss of locality.
Use a worksheet with one row per failure or transition:
| Scenario | Measurements | Gate before production |
|---|---|---|
| Broker replacement | Metadata rebuild, bytes replayed, cache warm-up, client errors | Service resumes within the stated recovery objective and the scope is observable |
| Scale-out during a write burst | Produce latency, WAL or buffer queue, upload lag, and controller work | The burst remains within the application service-level objective (SLO) and no queue grows without a bound |
| Replay surge | Fetch rate, cache occupancy, object-store bandwidth, and consumer lag | Historical reads do not starve tailing traffic or exceed the agreed cost envelope |
| Object-storage degradation | Request latency, throttling, retries, and acknowledged-but-not-uploaded data | Backpressure, alerting, and recovery behavior are explicit |
| Zone or endpoint failure | Traffic by boundary, access errors, and failover path | The storage and client paths have a tested fallback |
| Rollback | Ownership state, offsets, dual-write or replay behavior, and operator steps | A rollback point is defined before traffic is moved |
Cost follows the same paths. Keep compute, object storage, requests, transfer, observability, and migration operations as separate lines. Cross-zone, cross-region, private-endpoint, and egress charges use provider-specific rules; use the current pricing page for the regions and paths in scope. Do not turn a remembered rate into a universal claim.
Compatibility is another scaling input. Inventory client versions, producer acknowledgement settings, idempotence, transactions, compaction, consumer groups, Kafka Streams, Connect, admin tooling, authentication, quotas, and monitoring integrations. A platform that removes data movement but changes a client contract has not delivered a low-risk scaling path.
Plan the migration itself as a workload. During a cutover, the source and target may both generate reads, writes, network traffic, and operational events. Measure the extra path, define the offset or replay strategy, and reserve capacity for rollback. “The object store is ready” is a storage check, not a migration gate.
If the team needs more detail on retention and remote-storage boundaries, the long-retention FinOps framework and the remote-log storage checklist provide adjacent planning questions. They should complement the partition-scaling worksheet, not replace it.
4How a shared-storage Kafka option changes the operating model
Once the neutral framework is explicit, a Kafka-compatible shared-storage platform can be evaluated against the same measurements. AutoMQ keeps Kafka protocol and semantics at the client boundary while using a Shared Storage architecture for durable stream data. Its S3Stream storage layer, WAL options, cache behavior, and object-storage path give operators concrete signals to place on the worksheet.
The operating change is about ownership. Brokers continue to handle protocol requests, partition leadership, scheduling, quotas, and cache management. Durable stream data is held in shared storage rather than being tied to one broker’s local disk. A scale-out operation can therefore focus on ownership and traffic placement, while the test still verifies cache warm-up, object-store access, controller work, and recovery.
AutoMQ’s architecture overview describes this separation. The WAL storage documentation is relevant when the workload requires a specific write-latency or failure-domain choice. Record the selected backend, its placement, its recovery path, and the measurements that justify cache and bandwidth settings.
The architecture does not remove network or object-storage planning. Measure producer and consumer locality, broker-to-storage paths, endpoint policies, request rates, retries, and egress rules. Validate identity and access management (IAM) permissions and failure behavior in the cloud or S3-compatible service you intend to use. The platform team still owns dashboards, thresholds, drills, and rollback gates.
This is the useful distinction between a product statement and an operating model. A shared-storage platform can change which bytes must move during a partition reassignment. The team still has to prove that metadata transitions, cache behavior, storage requests, and client compatibility remain inside the production envelope.
5Decision checklist and FAQ
Bring these questions to a design review:
- What is changing? Is the operation adding compute, moving ownership, changing retention, or changing the durable storage path?
- Where are the bytes? Can the team identify the WAL or buffer, cache, object-store objects, and ranges needed for recovery?
- Which signals decide? Are metadata latency, request latency, upload lag, cache behavior, object-store requests, and network boundaries visible?
- What fails first? Has the team tested broker replacement, replay surge, object-storage degradation, zone or endpoint failure, and rollback?
- What is compatible? Are client behavior, transactions, compaction, Connect, Streams, authentication, quotas, and monitoring covered?
- What stops the rollout? Is there a stop condition, rollback point, dashboard owner, and review date?
5.1Does diskless Kafka eliminate partition reassignment?
No. Topics and partitions remain logical Kafka objects, and ownership still has to be coordinated. The potential change is whether reassignment requires copying retained bytes between broker-local disks. Verify the implementation’s data path and measure the metadata, cache, and recovery work.
5.2Is partition count the same as storage capacity?
No. Partition count affects metadata and scheduling. Storage capacity follows retained bytes, write rate, read rate, object layout, and retention. They should be planned as separate dimensions and then tested together under peak and recovery conditions.
5.3Does object storage make scaling free?
No. Object-storage requests, network paths, cache misses, compute, observability, and migration operations all have costs. A useful forecast keeps these lines separate and uses provider-specific inputs for the selected regions and endpoints.
5.4What should a pilot prove?
A pilot should prove the workload contract, the data path, and the failure gates: normal produce and fetch behavior, scale-out and scale-in, replay, broker replacement, object-store degradation, client compatibility, rollback, and the measurements that tell an operator to stop.
Partition scaling becomes manageable when “add a broker” stops being a proxy for every other operation. Separate metadata, compute, durable bytes, cache, and network paths; then attach a measurement and a failure gate to each one. If the current Kafka plan still assumes that every scaling event must move a broker’s retained log, start an AutoMQ evaluation with the workload measurements and rollback criteria in hand.
