Scaling Millions of IoT Devices with Akka Cluster Sharding and Entity Passivation
Learn how Akka Cluster Sharding with entity passivation lets you address millions of IoT devices as actors without keeping them all in memory, and how to tune the idle timeout for your workload.
08 Jul 2025, 06:44 UTC

The problem: two million devices, one heap
Your fleet has roughly two million devices. Each sends a temperature reading every few minutes to once an hour, and each reading must be validated against that device's recent history. The natural model is one actor per device: it owns the device's state, handles its messages one at a time, and persists changes. The problem is arithmetic — two million resident actors, each with state and mailbox overhead, will not fit on any heap you can realistically buy.
The useful takeaway: Akka Cluster Sharding plus entity passivation keeps the one-actor-per-device model while you only pay for devices that are actually talking. Idle actors are stopped, and their state is reloaded from a journal on the next message. The decision you are signing up for is where to set the idle timeout, and it is worth making with data.
How sharding routes by entity ID
Cluster Sharding gives you a ShardRegion per entity type — a local entry point that knows nothing about specific devices. You supply a message extractor that derives the entity ID (here, deviceId) from each command and maps it to a shard with a stable hash, typically (math.abs(deviceId.hashCode) % numberOfShards).toString. Every node computes the same mapping, so the cluster agrees on placement without a lookup service.
When the first message for a device arrives, the shard spawns that device's actor on whichever node currently owns its shard; later messages are forwarded to the same actor. As nodes join or leave, whole shards move and entities restart on the new owner. You never hard-code which node owns a device.
Two sizing notes from long-standing documentation guidance: set number-of-shards well above your maximum expected node count (roughly 10x is the common rule) so rebalancing stays even, and pick the value before launch — the mapping is hash-based, so changing it later reassigns entities to different shards.
Passivation: the memory-versus-churn lever
Passivation means the shard stops an entity actor after a configured idle timeout, reclaiming its memory. In typed sharding this is driven by a settings key; in classic sharding the shard sends a Passivate message the entity must answer by stopping. Key names and defaults differ across Akka versions, so check reference.conf for the version you run before tuning anything.
Passivation pairs naturally with Akka Persistence (event sourcing): the entity appends each reading to a journal plugin such as JDBC or Cassandra and replays its events to rebuild state on reactivation. Without persistence, passivation discards state — fine for session-like entities where the next message carries everything needed.
You can estimate residency with simple arithmetic. If devices report every 10 minutes on average and the idle timeout is 5 minutes, most entities stay resident between readings. If devices report hourly, a 5-minute timeout keeps roughly one twelfth of the fleet in memory at steady state. The timeout is the knob; the reporting distribution is the input you need to measure.
Worked example: a device telemetry service
Each device is an event-sourced entity that persists ReadingRecorded events and answers RecordReading commands:
object Device {
sealed trait Command
final case class RecordReading(deviceId: String, celsius: Double,
replyTo: ActorRef[Ack]) extends Command
final case class Ack(deviceId: String)
final case class ReadingRecorded(celsius: Double)
final case class State(lastReadingCelsius: Option[Double] = None)
def apply(entityId: String): Behavior[Command] =
EventSourcedBehavior[Command, ReadingRecorded, State](
persistenceId = PersistenceId("Device", entityId),
emptyState = State(),
commandHandler = (_, cmd) => cmd match {
case RecordReading(_, celsius, replyTo) =>
Effect.persist(ReadingRecorded(celsius))
.thenRun(_ => replyTo ! Ack(entityId))
},
eventHandler = (state, evt) =>
state.copy(lastReadingCelsius = Some(evt.celsius))
)
}
Initialization runs once per node; the extractor maps each command's deviceId into one of 2,000 shards:
val TypeKey = EntityTypeKey[Device.Command]("device")
val region: ActorRef[ShardingEnvelope[Device.Command]] =
ClusterSharding(system).init(
Entity(TypeKey)(entityContext => Device(entityContext.entityId))
.withMessageExtractor(messageExtractor) // deviceId -> shard via hash-mod
)
Configuration (Akka 2.6/2.7-era typed sharding keys — verify against your version's reference.conf):
akka.cluster.sharding {
number-of-shards = 2000 # ~10x the largest cluster you expect
passivate-idle-entity-after = 5 m # stop entities idle longer than this
}
A dashboard either asks the region for a device's latest reading or, better for read-heavy access, uses a projection that consumes the journal into a queryable store so dashboards never touch the sharded actors. Treat the code as a map of the moving parts rather than copy-paste material: extractor and passivation APIs shifted between Akka versions.
Trade-offs and limits
- Churn on the other side of the timeout. Every reactivation pays actor spawn plus journal replay. If a device's event history grows long, replay dominates — keep events per entity bounded or add snapshots so reactivation stays cheap.
- At-most-once delivery. Ordinary sharded messages can be lost across rebalances or node failures. Make
RecordReadingidempotent (for example, ignore a reading whose timestamp was already recorded) or add at-least-once delivery on the sending side. - Sharding presumes a healthy cluster. Seed-node formation, failure detection, and a split-brain resolver are prerequisites, not optional extras.
- Licensing. Since late 2022 Akka is under a Business Source License with a free-use grant for smaller organizations and conversion to Apache 2.0 per release. Confirm current terms for your organization, or evaluate Apache Pekko, the Apache-2.0 fork with equivalent sharding APIs.
How to check it on your setup
- Build a two- or three-node local cluster (distinct ports on one machine suffice), send one message, and confirm from logs or sharding stats that the entity spawned on demand.
- Wait past the idle timeout, confirm the entity stops, send again, and confirm it restarts and replays its state.
- Find the passivation key and its default in your version's
reference.confbefore overriding it. - Load-test with your real reporting distribution: sweep the timeout (for example 1, 5, and 15 minutes) and record heap footprint and reactivation latency.
Model entities as transient, persisted resources rather than permanent residents, tune the idle timeout against measured traffic, and one-actor-per-device scales to millions without a heroic heap.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.