You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否在Lagom中直接发布实体更新流为WebSocket服务的Source?如何实现?

问题:无需额外Topic直接发布实体更新流为Akka HTTP 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,具体步骤如下:

  1. 把服务签名调整到你想要的格式
    先把服务改成目标签名,不需要客户端传任何参数,直接返回事件流:

    @Override 
    public ServiceCall<NotUsed, Source<EntityPublicEvent, ?>> entityUpdates() { 
        return request -> {
            // 直接从实体事件源获取流,不用绕Topic
            Source<EntityPublicEvent, ?> eventSource = getEntityEventStream()
                .map(this::convertEvent);
            return CompletableFuture.completedFuture(eventSource);
        };
    }
    
  2. 直接绑定实体的事件流
    根据你用的事件存储方式选对应的实现:

    • 如果是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;
              });
      }
      
  3. 如果还要保留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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 09:50:43