Lagom was Lightbend's opinionated microservices framework for Scala and Java: service APIs described as Scala traits, state kept as event-sourced entities sharded across a cluster, query tables built from the event stream, and Kafka topics published from the same stream, all wired together and runnable with one command in development. It made a demanding architecture approachable.
It is also finished. Lagom reached end of life on 1 July 2024 and no longer receives security or functional patches; 1.6.7 is the last stable release. So this article is not a pitch. It is for engineers who own a Lagom system: how a Lagom 1.6 service actually works, where it fails, how to run it safely while it lives, and how to move it, piece by piece, onto the Akka 2.6 primitives it was built on or their open-source successor, Apache Pekko.
What Lagom gave you
A Lagom service has four parts. The service descriptor is a Scala trait that declares the HTTP calls and the Kafka topics the service exposes; Lagom generates the server routing and a typed client from it, so other services call it like a local method. Persistent entities hold write-side state as event-sourced actors, sharded so that each entity id lives on exactly one node. Read-side processors consume the tagged event stream and build query tables. Topic producers publish selected events to Kafka for other services.
Underneath, Lagom 1.6 is Akka 2.6: Cluster Sharding, Akka Persistence and Akka Streams, plus Play for HTTP. That matters for migration, because most of what you wrote is already Akka code with Lagom glue around it.
The descriptor: an API you can call like a method
Here is the descriptor of an accounts service. pathCall binds a URL pattern to a method, and withTopics declares a Kafka topic with a partition key, so all messages for one account land on one partition in order.
trait AccountService extends Service {
def openAccount(id: String): ServiceCall[OpenAccount, Done]
def getAccount(id: String): ServiceCall[NotUsed, AccountView]
def accountEvents: Topic[AccountMessage]
override final def descriptor: Descriptor = {
import Service._
named("account")
.withCalls(
pathCall("/api/accounts/:id/open", openAccount _),
pathCall("/api/accounts/:id", getAccount _)
)
.withTopics(
topic("account-events", accountEvents)
.addProperty(
KafkaProperties.partitionKeyStrategy,
PartitionKeyStrategy[AccountMessage](_.accountId)
)
)
.withAutoAcl(true)
}
}A consuming service binds a client to this trait and calls accountService.getAccount("acct_4821").invoke(). The Lagom service locator resolves the name account to an address. That convenience is also a coupling: both sides must share the API module, so changing a request type is a coordinated release.
The entity: event sourcing on Cluster Sharding
Since 1.6, Lagom recommends writing entities directly as Akka Persistence Typed behaviours, which is good news for migration. The behaviour receives a command, validates it against current state, persists zero or more events and replies; on recovery, it replays its events to rebuild state.
object AccountBehavior {
def create(entityContext: EntityContext[AccountCommand]): Behavior[AccountCommand] =
EventSourcedBehavior
.withEnforcedReplies[AccountCommand, AccountEvent, AccountState](
persistenceId = PersistenceId(entityContext.entityTypeKey.name, entityContext.entityId),
emptyState = AccountState.empty,
commandHandler = (state, cmd) => state.applyCommand(cmd),
eventHandler = (state, evt) => state.applyEvent(evt)
)
.withTagger(AkkaTaggerAdapter.fromLagom(entityContext, AccountEvent.Tag))
}
object AccountEvent {
// The shard count is part of your stored data: changing it later re-tags new events only.
val Tag: AggregateEventShards[AccountEvent] = AggregateEventTag.sharded[AccountEvent](numShards = 10)
}
// In the application loader:
clusterSharding.init(Entity(AccountState.typeKey)(ctx => AccountBehavior.create(ctx)))
// In the service implementation:
implicit val timeout: Timeout = Timeout(5.seconds)
override def openAccount(id: String) = ServiceCall { req =>
clusterSharding.entityRefFor(AccountState.typeKey, id)
.ask[Confirmation](reply => Open(req.owner, req.plan, reply))
.map {
case Accepted => Done
case Rejected(reason) => throw BadRequest(reason)
}
}withEnforcedReplies makes the compiler insist every command gets a reply, which removes a class of hung-ask bugs. AkkaTaggerAdapter.fromLagom applies Lagom's sharded tag, so each event is written with one of ten tags, chosen by hashing the entity id. Tags are how the read side and the topic producer find events later: each tag is one ordered stream, processed by one worker at a time across the cluster. The deeper mechanics are in the Akka Persistence article and Akka Cluster article.
Read sides and topics: two consumers of the same tags
A ReadSideProcessor declares which tags it consumes and a handler per event type, typically writing to Cassandra or a relational database through Lagom's Slick or JDBC read-side support, and Lagom stores its offset per tag so a restart resumes where it stopped. A topic producer does the same thing with a Kafka sink instead of a table.
// Producer side: publish from the tagged event stream, offsets tracked by Lagom.
override def accountEvents: Topic[AccountMessage] =
TopicProducer.taggedStreamWithOffset(AccountEvent.Tag) { (tag, fromOffset) =>
persistentEntityRegistry
.eventStream(tag, fromOffset)
.collect { case ev if isPublic(ev.event) => (toMessage(ev), ev.offset) }
}
// Consumer side, in another service: at-least-once, so the handler must be idempotent.
accountService.accountEvents.subscribe.atLeastOnce(
Flow[AccountMessage].mapAsync(parallelism = 1)(msg => billing.applyIdempotently(msg).map(_ => Done))
)Both are at-least-once. The offset is saved after the effect, so a crash between them replays a few events: the read-side handler must be an upsert and the Kafka consumer must deduplicate. Filtering to public events in the producer is deliberate: internal events are your persistence format, and publishing them directly turns every refactor into a breaking change for other teams.
Worked example: one request, end to end
Follow POST /api/accounts/acct_4821/open. Play routes it to openAccount, which asks the entity ref for acct_4821. Sharding hashes the id to a shard, finds the node hosting that shard, and delivers the Open command; if the entity is not in memory, it is started and its events are replayed first. The command handler checks the account is not already open, persists AccountOpened with tag AccountEvent3, say, and replies Accepted. Only after the journal write succeeds does the HTTP call return 200.
Some time later, typically well under a second, the read-side worker for tag 3 polls the journal from its stored offset, sees the event, upserts a row in the accounts view table and saves the new offset. Independently, the topic producer for tag 3 publishes an AccountMessage keyed by acct_4821 to Kafka, and the billing service's subscriber starts a trial. A GET issued immediately after the 200 may not see the account yet: that read-your-writes gap is the eventual consistency you accepted, and the UI must handle it.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Ask timeouts after a deploy | Entities replaying long event histories on first touch | Snapshots every N events; warm-up; longer timeout for first command |
| Two nodes both think they own shards | Network partition without a downing strategy | Configure the Akka 2.6 split brain resolver; never use auto-down |
| Read side stops advancing for one tag | Handler throws on one event forever | Make handlers total; log and park poison events; alert on per-tag lag |
| Events vanish from a projection | numShards changed, so old and new events use different tags | Treat the shard count as fixed, or migrate offsets deliberately |
| Duplicate trial in billing | At-least-once topic consumer without deduplication | Idempotent handler keyed by event or entity version |
| Unrecoverable entity | Event class changed incompatibly | Serializer migrations; never change a persisted event's meaning |
Serialization: the part that outlives the framework
Lagom 1.6 Scala services usually serialize events and commands with Play JSON formats registered in a JsonSerializerRegistry, and evolve them with JSON migrations that rewrite old payloads as they are read. The journal stores each event's bytes together with a serializer id and a manifest naming its type. Those rows are the most valuable thing the service owns, and they will be read by whatever replaces Lagom.
So treat serialization as a contract before you migrate anything. Write a test that reads a sample of real production events, including the oldest ones, through the serializer and asserts they decode. Keep that test running through the migration, against the new code. If the new system uses a different serializer, register a compatible one under the same id and manifest, or convert the journal in a controlled job. Never rename an event class without a manifest mapping, because the old rows still name the old class.
Running Lagom safely while it lives
An end-of-life framework is a dependency risk that grows every month. Lagom pins versions of Akka 2.6, Play and many transitive libraries, and fixes to them no longer flow through Lagom releases. Run dependency scanning on the final assembled classpath, not just your declared dependencies. Where a CVE lands in a transitive library, a dependency override can sometimes help, but test it thoroughly, because Lagom was never tested against it.
Freeze the architecture: no new Lagom services, and no new features that deepen the dependency, such as more Lagom-specific read-side processors. Put the Lagom services behind an API gateway so callers depend on HTTP contracts, not on the Lagom client. Measure per-tag read-side lag, ask latency, shard rebalance counts and journal write latency, so that you know the system's health before you start changing it.
Migrating: piece by piece, not big bang
Because Lagom 1.6 sits on Akka 2.6, the migration is mostly removing glue. Map each Lagom piece to the primitive underneath it.
| Lagom piece | Replacement | Notes |
|---|---|---|
| Service descriptor and client | Akka HTTP or Pekko HTTP routes, or Play controllers, plus an OpenAPI contract | Callers move to plain HTTP clients |
| Persistent entity (typed) | The same EventSourcedBehavior on Cluster Sharding | Drop the Lagom tagger adapter, keep identical tag names |
| Read-side processor | Akka Projections or Pekko Projections | Offsets live in a different store; migrate or rebuild |
| TopicProducer | A projection whose handler produces to Kafka | Keep the same topic, key and message format |
| Service locator | Kubernetes DNS or Akka Discovery | Usually already in place in production |
| Dev mode runAll | Docker Compose or Testcontainers | The real loss in developer experience |
Choose the destination first. Akka from 2.7 onwards is licensed under the Business Source License, so check its terms for your use. Apache Pekko is an Apache-licensed fork of Akka 2.6 with renamed packages, which makes it the closer step from Lagom's own Akka version. Either way the API concepts are the same.
Then go service by service with a strangler approach: stand up the new service sharing the same journal tables and serializers, move the HTTP routes behind the gateway, then move projections one at a time. Projections are the delicate step. A new projection does not know the Lagom offset store; either copy the offsets per tag into the new store, or rebuild the view table from the start of the journal into a new table and switch reads over. Rebuilding is slower but verifiable; prefer it where the journal is small enough.
Trade-offs that remain
Leaving Lagom does not remove event sourcing's costs: event schema versioning, replay time and eventual consistency stay with you. If a service never needed them, for example a simple CRUD service that adopted Lagom because the platform mandated it, the best migration may be to a plain HTTP service with a relational table, retiring its journal. Check the Akka Typed article for the behaviour model you keep, and the Play Framework article if Play becomes your HTTP layer.
What to do next
- Inventory every Lagom service: entities, tags and shard counts, read sides, topics, and which services consume them.
- Add dependency scanning on the assembled classpath and put the Lagom services behind a gateway with documented HTTP contracts.
- Add per-tag lag, ask latency and shard rebalance metrics, and confirm a split brain resolver is configured.
- Choose Akka or Apache Pekko after reviewing licence terms, and spike one entity with its tags unchanged.
- Migrate one service end to end, rebuilding one projection into a new table and comparing it with the old one before switching.
- Set a retirement date for the last Lagom dependency and track it like any other security deadline.