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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 01:37:16