Java环境下Akka WebSocket认证成功后无法接收数据流问题
Akka WebSocket无法接收认证后金融数据流问题
我借助公开金融数据流学习Akka WebSocket,WebSocket已完成认证,但无法接收后续数据流。参考Akka官方文档编写两段Java代码,执行后仅输出认证成功消息Got message: {"status_code":200,"message":"Authorized"},无后续数据。已确认WebSocket服务器正常(可接收消息),修改代码尝试保持Source开启后结果仍一致。
原始代码
import akka.Done; import akka.NotUsed; import akka.actor.ActorSystem; import akka.http.javadsl.Http; import akka.http.javadsl.model.ws.Message; import akka.http.javadsl.model.ws.TextMessage; import akka.http.javadsl.model.ws.WebSocketRequest; import akka.stream.Materializer; import akka.stream.SystemMaterializer; import akka.stream.javadsl.Flow; import akka.stream.javadsl.Sink; import akka.stream.javadsl.Source; import java.util.concurrent.CompletionStage; public class EodHistoricalDataTest { public static void main(String[] args) { final ActorSystem system = ActorSystem.create("eodhistoricaldata-example"); final Materializer materializer = SystemMaterializer.get(system).materializer(); Http http = Http.get(system); Sink<Message, CompletionStage<Done>> printSink = Sink.foreach((message) -> System.out.println("Got message: " + message.asTextMessage().getStrictText()) ); final String subscribeMessage = "{\"action\":\"subscribe\",\"symbols\":\"AMZN\" }"; final Source<Message, NotUsed> initialSource = Source.single(TextMessage.create(subscribeMessage)); final Flow<Message, Message, NotUsed> flow = Flow.fromSinkAndSource( printSink, initialSource); http.singleWebSocketRequest( WebSocketRequest.create("wss://ws.eodhistoricaldata.com/ws/us?api_token=demo"), flow, materializer); } }
修改后的尝试代码
public class EodHistoricalDataTest { public static void main(String[] args) { final ActorSystem system = ActorSystem.create("eodhistoricaldata-example"); final Materializer materializer = SystemMaterializer.get(system).materializer(); Http http = Http.get(system); final Message subscribeMessage = TextMessage.create("{\"action\":\"subscribe\",\"symbols\":\"AMZN\" }"); Sink<Message, CompletionStage<Done>> printSink = Sink.foreach((message) -> System.out.println("Got message: " + message.asTextMessage().getStrictText()) ); final Flow<Message, Message, NotUsed> flow = Flow.fromSinkAndSource( printSink, Source.single(subscribeMessage)); http.singleWebSocketRequest( WebSocketRequest.create("wss://ws.eodhistoricaldata.com/ws/us?api_token=demo"), flow, materializer); } }
问题根源
你用Source.single()发送订阅消息,发送完成后Source会立即关闭。Akka WebSocket的Flow.fromSinkAndSource中,发送端Source关闭会触发整个WebSocket连接的终止信号,服务器后续的数据流就无法再推送给客户端了。服务器认证成功后需要客户端保持连接活跃,才能持续推送数据。
解决方案
让发送端Source保持活跃状态,不要在发送订阅消息后就关闭。可以用Source.concat把订阅消息和一个永不结束的Source拼接,比如Source.maybe(),这样发送完订阅后Source会一直等待,不会触发连接关闭。
修复后的代码
import akka.Done; import akka.NotUsed; import akka.actor.ActorSystem; import akka.http.javadsl.Http; import akka.http.javadsl.model.ws.Message; import akka.http.javadsl.model.ws.TextMessage; import akka.http.javadsl.model.ws.WebSocketRequest; import akka.stream.Materializer; import akka.stream.SystemMaterializer; import akka.stream.javadsl.Flow; import akka.stream.javadsl.Sink; import akka.stream.javadsl.Source; import java.util.concurrent.CompletionStage; public class EodHistoricalDataTest { public static void main(String[] args) { final ActorSystem system = ActorSystem.create("eodhistoricaldata-example"); final Materializer materializer = SystemMaterializer.get(system).materializer(); Http http = Http.get(system); Sink<Message, CompletionStage<Done>> printSink = Sink.foreach((message) -> System.out.println("Got message: " + message.asTextMessage().getStrictText()) ); final Message subscribeMessage = TextMessage.create("{\"action\":\"subscribe\",\"symbols\":\"AMZN\" }"); // 发送订阅消息后保持Source活跃,不关闭连接 final Source<Message, NotUsed> keepAliveSource = Source.single(subscribeMessage) .concat(Source.maybe()); // Maybe Source会一直等待,不会主动终止 final Flow<Message, Message, NotUsed> flow = Flow.fromSinkAndSource( printSink, keepAliveSource); http.singleWebSocketRequest( WebSocketRequest.create("wss://ws.eodhistoricaldata.com/ws/us?api_token=demo"), flow, materializer); } }
补充说明
Source.maybe()创建的Source永远不会主动终止,这样发送端Flow不会发出终止信号,WebSocket连接会持续打开,服务器就能推送后续数据流。- 若后续需要动态发送其他消息,可以把
Source.maybe()换成Source.queue,实现消息的动态发送。 - 确认订阅消息格式符合服务器要求(比如符号、action字段是否正确),避免因格式问题导致服务器不推送数据。
内容的提问来源于stack exchange,提问作者blue-sky
相关产品推荐
相关产品推荐

