如何在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
相关产品推荐
相关产品推荐

