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

Akka框架下如何从服务端主动关闭WebSocket连接?

Akka服务端主动关闭WebSocket连接实现问题及解决方案

问题现象

实现Akka框架服务端主动关闭WebSocket连接功能时,使用PoisonPill向Source.actorRef发送消息触发关闭,运行时报错如下:

PoisonPill message sent to StageActor(akka://system/system/Materializers/StreamSupervisor-1/$$c-actorRefSource) will be ignored, since it is not a real Actor.Use a custom message type to communicate with it instead.

错误原因

Akka Stream的Source.actorRef返回的ActorRef不是原生的Akka Actor,是Stream阶段的虚拟代理Actor,不支持PoisonPill、Kill这类原生Actor的系统控制消息,发送这类消息会直接被忽略,对应上述报错提示。

原错误实现代码

fun main() {

    val system = ActorSystem.create("system")
    val materializer = Materializer.createMaterializer(system)

    val http = Http.get(system)

    val routeFlow = path("ws") {
        get {
            parameter("as") { sys ->
                parameter("id") { id ->
                    parameter("v") { v ->
                        handleWebSocketMessages(mainFlow(sys, id, v))
                    }
                }
            }
        }
    }.flow(system, materializer)

    http.newServerAt("0.0.0.0", 8080).withMaterializer(materializer).bindFlow(routeFlow).thenRun {
        println("Web socket server is running at localhost:8080")
    }
}

fun mainFlow(sys: String, id: String, v: String): Flow<Message, Message, NotUsed> {
    val source = Source.actorRef<Message>(
        { Optional.empty() },
        { Optional.empty() },
        100,
        OverflowStrategy.fail())
        .mapMaterializedValue {
            it.tell(PoisonPill.getInstance(), ActorRef.noSender())
        }

    val sink = Flow.create<Message>()
        .map {
            
        }
        .to(Sink.onComplete {
            
        })

    return Flow.fromSinkAndSource(sink, source)
}

错误的关闭触发代码段:

.mapMaterializedValue {
    it.tell(PoisonPill.getInstance(), ActorRef.noSender())
}

正确实现方案

Source.actorRef的第一个入参是完成匹配函数,需要自定义关闭信号,当收到匹配的消息时返回CompletionStrategy触发Source正常完成,Source终止后WebSocket连接会自动关闭。

fun mainFlow(connectionsController: ActorRef, sys: String, id: String, v: String, system: ActorSystem): Flow<Message, Message, NotUsed> {
    val source = Source.actorRef<Message>(
        {
            if (it is Done) {
                Optional.of(CompletionStrategy.immediately())
            } else {
                Optional.empty()
            }
        },
        { Optional.empty() },
        16,
        OverflowStrategy.fail())
        .mapMaterializedValue {
            it.tell(Done.done(), ActorRef.noSender())
        }

    val sink = Flow.create<Message>()
        .map {
            
        }
        .to(Sink.onComplete {
            
        })

    return Flow.fromSinkAndSource(sink, source)
}

实现说明

  • 使用Akka自带的Done消息作为关闭信号,也可自定义其他消息类型,只要在完成匹配函数中添加对应判断逻辑即可
  • 当Source.actorRef收到Done消息时,会触发立即完成策略,终止Source流
  • Akka HTTP的WebSocket处理逻辑会在流终止时自动关闭对应的WebSocket连接,满足服务端主动关闭的需求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 12:15:03