能否在Lagom中直接发布实体更新流为WebSocket服务的Source?如何实现?
我想请教一下,能不能不需要额外的Topic,直接把实体更新流发布成如下签名的服务返回的Source?
public ServiceCall<NotUsed, Source<EntityPublicEvent, ?>> entityUpdates()
我的目标是在Angular 5里通过WebSocket消费这些事件。
我现在的实现代码是这样的:
@Override public ServiceCall<String, Source<EntityPublicEvent, ?>> getUpdateNotifications() { return request -> CompletableFuture.completedFuture(updatesTopic.subscriber().map(this::convertEvent)); }
目前Angular已经成功连接,服务也被调用并返回了Source,但Angular完全收不到任何更新消息。日志里出现了这段内容,我猜这可能是收不到消息的原因:
[info] akka.actor.RepointableActorRef [sourceThread=currency-impl-application-akka.actor.default-dispatcher-3, akkaTimestamp=06:59:55.778UTC, akkaSource=akka://currency-impl-application/system/distributed PubSubMediator, sourceActorSystem=currency-impl-application] - Message [tcs.currency.impl.CurrencyEvent$Created] without sender to Actor[akka://currency-impl-application/system/distributedPubSubMediator#128487187] was not delivered. [1] dead letters encountered. This logging can be turned off or adjusted with configuration settings 'akka.log-dead-letters' and 'akka.log-dead-letters-during-shutdown'.
首先明确告诉你:完全可以不用额外Topic直接发布实体更新流,不过得先搞定当前的死信问题——这就是Angular收不到消息的核心原因。
先搞懂为什么会出现死信?
从日志内容来看,你的CurrencyEvent$Created消息发给了分布式PubSubMediator,但没被成功投递,大概率是这两个原因:
- 你的
updatesTopic订阅者没有和返回给客户端的Source正确绑定,导致消息发出来后没人接收,最终变成死信。 - 当前实现用了
String类型的请求参数,但你目标是NotUsed,这种参数不匹配可能导致订阅逻辑和Source的生命周期脱节,流根本没真正激活。
怎么实现无额外Topic的实体更新流?
你可以直接从实体的原生事件流(比如Akka Persistence的事件日志,或者内部的事件发布器)创建Source,跳过中间的PubSub Topic,具体步骤如下:
把服务签名调整到你想要的格式
先把服务改成目标签名,不需要客户端传任何参数,直接返回事件流:@Override public ServiceCall<NotUsed, Source<EntityPublicEvent, ?>> entityUpdates() { return request -> { // 直接从实体事件源获取流,不用绕Topic Source<EntityPublicEvent, ?> eventSource = getEntityEventStream() .map(this::convertEvent); return CompletableFuture.completedFuture(eventSource); }; }直接绑定实体的事件流
根据你用的事件存储方式选对应的实现:- 如果是Akka Persistence实体,直接订阅它的事件日志:
private Source<CurrencyEvent, ?> getEntityEventStream() { // 替换成你的实体ID和实际Journal配置 return PersistenceQuery.get(system) .journalFor("akka.persistence.journal.inmem") .eventsByPersistenceId("currency-1", 0, Long.MAX_VALUE) .map(EventEnvelope::event) .cast(CurrencyEvent.class); } - 如果用的是Akka的全局EventStream,直接订阅事件类型:
private Source<CurrencyEvent, ?> getEntityEventStream() { return Source.actorRef(100, OverflowStrategy.dropHead()) .mapMaterializedValue(actorRef -> { system.eventStream().subscribe(actorRef, CurrencyEvent.class); return actorRef; }); }
- 如果是Akka Persistence实体,直接订阅它的事件日志:
如果还要保留Topic,怎么修复当前问题?
要是你不想放弃Topic,那得确保订阅和发布的逻辑都正确:- 检查
updatesTopic的创建方式,必须用DistributedPubSub的topic方法:Topic updatesTopic = DistributedPubSub.get(system).topic("entity-updates-topic"); - 发布事件的时候,要发给Topic的Publisher,而不是直接给Mediator:
// 正确的事件发布方式 DistributedPubSub.get(system).publisherFor("entity-updates-topic").tell(event, ActorRef.noSender()); - 确保返回的Source是活跃的,不要在转换过程中意外终止流。
- 检查
最后验证Angular的WebSocket连接
Angular端要确保用正确的WebSocket订阅逻辑,比如用RxJS的webSocket:
import { webSocket } from 'rxjs/webSocket'; // 替换成你的服务WebSocket端点 const eventSocket$ = webSocket('ws://your-server-domain/entity-updates'); eventSocket$.subscribe( event => console.log('收到实体更新:', event), error => console.error('WebSocket错误:', error), () => console.log('WebSocket连接关闭') );
内容的提问来源于stack exchange,提问作者monad

