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

如何在Akka HTTP SSE中刷新Source以获取最新Payload?

问题分析

你当前的代码逻辑是仅在SSE连接建立时调用一次some_def(key),之后就一直重复发送第一次获取到的T实例,自然不会响应后续some_def返回的新数据。核心问题是Source.fromFuture只会执行一次异步调用,后续不会重新触发。

解决方案

要实现定期刷新数据,需要让some_def(key)被周期性调用,每次获取最新的结果再发送。可以用Source.tick来触发定时任务,每次tick时调用some_def并处理结果:

(path("events") & get) {
  complete {
    // 0秒延迟立即执行第一次调用,之后每2秒触发一次,获取最新数据
    Source.tick(0.seconds, 2.seconds, ())
      // 每次触发时重新调用some_def,将异步结果转为Source
      .flatMapConcat(_ => Source.fromFuture(some_def(key)))
      // 过滤掉错误结果,仅保留有效数据
      .collect { case Right(s: T) => s }
      // 转换为ServerSentEvent格式
      .map(s => ServerSentEvent(s.message))
      // 发送心跳包保持连接
      .keepAlive(1.second, () => ServerSentEvent.heartbeat)
  }
}
关键改动说明
  • 用Source.tick替代单次Future触发:确保每隔指定时间就重新获取最新数据
  • 每次tick时重新调用some_def(key):不再复用初始的单次结果
  • 用collect过滤异常结果:避免错误数据流入SSE流

内容的提问来源于stack exchange,提问作者Juhi Talukdar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 11:45:35