Distributed Systems: Sharding, Kafka, Retry Storms & Resilience
Scaling a system is only half the problem.
The harder problem is what happens when the system comes under pressure.
A database slows down.
Kafka consumers fall behind.
A dependency starts timing out.
Retries increase.
Shared resources become exhausted.
One small problem can turn into a much larger outage.
This is why modern distributed systems need both scalability and resilience.
Scalability helps a system handle more workload.
Resilience determines how well it behaves when something goes wrong.
Database Sharding: Scaling Beyond One Database
A single database is often the simplest place to start.
As traffic and data grow, however, the database can become a bottleneck.
Common warning signs include:
- Slow queries
- High CPU usage
- Large indexes
- Write contention
- Storage pressure
- Replication lag
Before introducing sharding, teams typically have simpler options to consider:
- Query optimization
- Index tuning
- Caching
- Read replicas
- Partitioning
- Vertical scaling
When these approaches are no longer enough, sharding can provide another scaling path.
What Is Database Sharding?
Sharding distributes data across multiple database nodes.
Instead of one database handling the entire workload, each shard manages a portion of the data.
This allows the system to scale horizontally.
For example:
Users → Shard A | Shard B | Shard C
Each shard handles part of the overall storage and workload.
The Trade-Off
Sharding improves scalability, but it also introduces complexity.
A production system may need to handle:
- Shard-key selection
- Request routing
- Hot shards
- Cross-shard queries
- Data rebalancing
- Operational overhead
A poor shard key can create an uneven workload and turn one shard into the new bottleneck.
Key takeaway: Sharding is a scaling strategy, not a free performance upgrade.
Kafka Lag: The Bottleneck May Be Somewhere Else
Kafka is designed for high-throughput event processing.
But increasing Kafka lag does not automatically mean Kafka is the problem.
Consider this flow:
Kafka → Consumer → Database
Now imagine the database becomes slow.
Consumers take longer to process events.
Kafka continues receiving events.
The backlog grows.
Consumer lag increases.
The bottleneck may actually be downstream.
Common Causes of Kafka Lag
When investigating lag, look beyond the Kafka brokers.
Possible causes include:
- Slow consumers
- Slow databases
- Poor partition keys
- Hot partitions
- Large messages
- Consumer rebalancing
- Retry amplification
- Downstream backpressure
Hot Partitions
Kafka parallelism depends heavily on partitioning.
If a partition key distributes traffic poorly, some partitions can become much busier than others.
This creates an uneven workload.
Some consumers may be overloaded while others have relatively little work.
Key takeaway: Kafka lag should be investigated across the entire processing pipeline.
Kafka Offsets: When Should Consumers Commit?
Another important part of Kafka consumer design is deciding when to commit an offset.
If an offset is committed before downstream processing finishes, Kafka may consider the message processed even though the database or API operation has not completed successfully.
That can create a data-loss risk.
But committing later introduces another possibility.
The database operation may succeed, while the consumer crashes before committing the offset.
The same event can then be processed again.
This creates the possibility of duplicate processing.
Why Idempotency Matters
Idempotency makes repeated execution safer.
This is particularly important for operations such as:
- Payments
- Orders
- Inventory updates
- Account changes
- Other business operations with side effects
The important principle is:
If duplicate processing is possible, the downstream operation should be designed to handle it safely where practical.
Retry Storms: When Recovery Creates More Load
Retries are useful when failures are temporary.
But retries can become dangerous when the downstream system is already overloaded.
Imagine a database starts timing out.
The application retries.
The database now receives additional requests.
Load increases.
Latency increases.
More requests timeout.
More retries begin.
The result is a feedback loop:
Failure → Retry → More Load → More Failure
This is commonly described as a retry storm.
The retry mechanism that was supposed to improve reliability can actually make the incident worse.
How to Control Retries
Retries should be treated as a controlled resource rather than an unlimited recovery mechanism.
Retry Budgets
Limit how many retries the system can generate.
Exponential Backoff
Increase the delay between attempts to reduce immediate pressure.
Jitter
Add randomness to retry timing so many clients do not retry simultaneously.
Circuit Breakers
Temporarily stop requests to an unhealthy dependency and fail fast.
Idempotency
Reduce the risk of harmful duplicate side effects when operations are repeated.
Load Shedding
During severe overload, reject or defer work instead of allowing pressure to spread deeper into the system.
The goal isn't to eliminate retries.
The goal is to prevent retries from becoming another source of failure.
Designing for Failure
Distributed systems depend on multiple services, networks and infrastructure components.
Failures are therefore part of normal operation.
Examples include:
- Network timeouts
- Database slowdowns
- Dependency failures
- Traffic spikes
- Resource exhaustion
- Consumer failures
The important question is not:
"Can this component fail?"
It probably can.
The better question is:
"What happens when it fails?"
Cascading Failures
A slow dependency can cause requests to remain active for longer.
Those requests consume resources.
More traffic arrives.
Timeouts increase.
Retries begin.
More requests reach the already unhealthy dependency.
Eventually, unrelated services can also experience resource pressure.
A local problem becomes a system-wide problem.
This is failure propagation.
The objective of resilience engineering is to limit that propagation.
Bulkheads: Controlling the Blast Radius
The bulkhead pattern focuses on resource isolation.
Imagine several workloads sharing the same resources.
If one workload becomes unhealthy, it may consume a disproportionate amount of:
- Threads
- Database connections
- Worker capacity
- CPU
- Memory
Other workloads then suffer even though they are healthy.
Bulkheads create separation between workloads.
For example:
Payments → Payment Resource Pool
Notifications → Notification Resource Pool
User APIs → User API Resource Pool
If notifications become unhealthy, they should not be able to consume every resource required by payments or user-facing APIs.
Why Bulkheads Matter
Bulkheads help achieve controlled degradation.
Instead of:
One failure → Everything fails
the goal becomes:
One failure → One area degrades → Other areas remain available
The Trade-Off
Isolation requires resource planning.
Too little isolation increases the blast radius.
Too much isolation can reduce overall resource efficiency.
The right balance depends on the workload and architecture.
How These Patterns Connect
These concepts are not isolated techniques.
They interact during real incidents.
Consider a simplified system:
Client → API → Kafka → Consumer → Database Shards
Now introduce a database slowdown.
The consumer takes longer to process messages.
Kafka lag increases.
The application may start retrying failed operations.
Those retries generate additional database traffic.
The database experiences even more pressure.
If the shard distribution is poor, one shard may become a hotspot.
If processing is not idempotent, duplicate operations may create additional problems.
If resources are shared across unrelated workloads, the failure can spread further.
This is why resilience must be considered as a system-level problem, not just a collection of individual patterns.
A Practical Resilience Checklist
When reviewing a distributed architecture, ask:
1. Where is the bottleneck?
Which component is most likely to limit throughput?
2. What happens when it slows down?
Does the system queue work, apply backpressure, shed load or continue sending requests?
3. Can retries increase the problem?
Look for unlimited or duplicated retry mechanisms across multiple layers.
4. Can an operation execute twice?
If yes, consider whether idempotency or deduplication is required.
5. Can one workload consume shared resources?
Check thread pools, connection pools, worker capacity and other shared resources.
6. What is the blast radius?
If one component fails, determine how many other components can be affected.
These questions can reveal resilience problems before they become production incidents.
Scalability vs. Resilience
Scalability and resilience solve different problems.
Scalability asks:
Can the system handle more workload?
Resilience asks:
Can the system continue operating when something goes wrong?
Sharding helps distribute database workload.
Kafka partitions help distribute event processing.
Retry budgets and backoff help control failure-generated traffic.
Idempotency helps make repeated processing safer.
Circuit breakers help prevent repeated calls to unhealthy dependencies.
Bulkheads help isolate resources.
Together, these patterns help create systems that can degrade in a controlled way instead of failing everywhere at once.
Key Takeaways
- Sharding enables horizontal database scaling but introduces operational and query complexity.
- Kafka lag does not always indicate a Kafka problem; consumers and downstream systems may be the bottleneck.
- Offset timing creates a trade-off between potential data loss and duplicate processing.
- Idempotency is important when operations can be executed more than once.
- Retries can amplify overload when they are uncontrolled.
- Backoff and jitter help reduce synchronized retry traffic.
- Circuit breakers can stop repeated calls to unhealthy dependencies.
- Bulkheads isolate resources and reduce failure propagation.
- Resilience is about controlling the impact of failure, not eliminating failure completely.
Final Thought
A distributed system doesn't need to be perfect.
It needs to fail predictably and safely.
One slowdown shouldn't become an outage.
One retry shouldn't become a storm.
One unhealthy workload shouldn't consume everything.
Scale gives a system capacity.
Resilience controls the blast radius when that capacity comes under pressure.
What do you focus on first when reviewing a production architecture: scalability, retries, messaging, or failure isolation?
