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.

Advertisement

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.

A Lagom 1.6 service: descriptor, sharded entity, journal, read side and topicService descriptorpathCall / withTopicsServiceCallCluster ShardingentityRefFor(typeKey, id)askEventSourcedBehaviorcommand -> events -> statepersistJournalevents tagged by shardeventStream(tag, offset)ReadSideProcessorbuilds query tablesTopicProducerpublishes to KafkaOffset storeper tag, per processorOther servicessubscribe.atLeastOnceEach tag shard is processed by one worker at a time across the cluster; the offset store makes restarts resume, not repeat from zero.
One Lagom service. Commands reach a sharded entity, events land in the journal with a shard tag, and two independent stream consumers, a read side and a topic producer, follow the tags with their own offsets.

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.

Advertisement

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

SymptomCauseFix
Ask timeouts after a deployEntities replaying long event histories on first touchSnapshots every N events; warm-up; longer timeout for first command
Two nodes both think they own shardsNetwork partition without a downing strategyConfigure the Akka 2.6 split brain resolver; never use auto-down
Read side stops advancing for one tagHandler throws on one event foreverMake handlers total; log and park poison events; alert on per-tag lag
Events vanish from a projectionnumShards changed, so old and new events use different tagsTreat the shard count as fixed, or migrate offsets deliberately
Duplicate trial in billingAt-least-once topic consumer without deduplicationIdempotent handler keyed by event or entity version
Unrecoverable entityEvent class changed incompatiblySerializer 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 pieceReplacementNotes
Service descriptor and clientAkka HTTP or Pekko HTTP routes, or Play controllers, plus an OpenAPI contractCallers move to plain HTTP clients
Persistent entity (typed)The same EventSourcedBehavior on Cluster ShardingDrop the Lagom tagger adapter, keep identical tag names
Read-side processorAkka Projections or Pekko ProjectionsOffsets live in a different store; migrate or rebuild
TopicProducerA projection whose handler produces to KafkaKeep the same topic, key and message format
Service locatorKubernetes DNS or Akka DiscoveryUsually already in place in production
Dev mode runAllDocker Compose or TestcontainersThe 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

  1. Inventory every Lagom service: entities, tags and shard counts, read sides, topics, and which services consume them.
  2. Add dependency scanning on the assembled classpath and put the Lagom services behind a gateway with documented HTTP contracts.
  3. Add per-tag lag, ask latency and shard rebalance metrics, and confirm a split brain resolver is configured.
  4. Choose Akka or Apache Pekko after reviewing licence terms, and spike one entity with its tags unchanged.
  5. Migrate one service end to end, rebuilding one projection into a new table and comparing it with the old one before switching.
  6. Set a retirement date for the last Lagom dependency and track it like any other security deadline.
Key takeaway: Lagom packaged Akka 2.6 sharding, persistence and streams behind service descriptors, entities, read sides and topics, and it reached end of life on 1 July 2024. Understand the tagged event streams that hold it together, because they decide ordering, lag and replay. Operate it defensively with classpath scanning, a gateway and per-tag metrics, then migrate service by service to plain Akka or Apache Pekko, keeping tag names and rebuilding projections deliberately.