【问题标题】:spring webflux: purely functional way to attach websocket adapter to reactor-netty serverspring webflux:将 websocket 适配器附加到 reactor-netty 服务器的纯功能方式
【发布时间】:2018-06-09 13:10:12
【问题描述】:

我无法找到将 WebSocketHandlerAdapter 附加到 reactor netty 服务器的方法。

要求: 我想启动一个 reactor netty 服务器并将 http (REST) 端点和 websocket 端点附加到同一台服务器。我已经阅读了文档和文档中提到的一些示例演示应用程序。他们展示了如何使用 newHandler() 函数将 HttpHandlerAdapter 附加到 HttpServer。但是当涉及到 websockets 时,他们会切换回使用 spring boot 和注释示例。我无法找到如何使用功能端点附加 websocket。

请为我指明如何实现这一点的正确方向。 1.如何将websocket适配器连接到netty服务器? 2.我应该使用HttpServer还是TcpServer?

注意: 1.我没有使用弹簧靴。 2. 我没有使用注释。 3. 尝试仅使用功能性 webflux 端点来实现这一点。

示例代码:

public HandlerMapping webSocketMapping() 
{
  Map<String, WebSocketHandler> map = new HashMap<>();
  map.put("/echo", new EchoTestingWebSocketHandler());
  SimpleUrlHandlerMapping mapping = new SimpleUrlHandlerMapping();
  mapping.setUrlMap(map);
  mapping.setOrder(-1);
  return mapping;
}
public WebSocketHandlerAdapter wsAdapter() 
{
  HandshakeWebSocketService wsService = new HandshakeWebSocketService(new ReactorNettyRequestUpgradeStrategy());
  return new WebSocketHandlerAdapter(wsService);
}

  protected void startServer(String host, int port) 
  {
    HttpServer server = HttpServer.create(host, port);
    server.newHandler(wsAdapter()).block();    //how do I attach the websocket adapter to the netty server
  }

【问题讨论】:

    标签: spring websocket reactor-netty


    【解决方案1】:

    不幸的是,如果不运行整个 SpringBootApplication,没有简单的方法可以做到这一点。否则,您将需要自己编写整个 Spring WebFlux 处理程序层次结构。考虑使用 SpringBootApplication 组合您的功能路由:

        @SpringBootApplication
        public class WebSocketApplication {
    
            public static void main(String[] args) {
                SpringApplication.run(WebSocketApplication.class, args);
            }
    
    
            @Bean
            public RouterFunction<ServerResponse> routing() {
                return route(
                        POST("/api/orders"),
                        r -> ok().build()
                );
            }
    
            @Bean
            public HandlerMapping wsHandlerMapping() {
                HashMap<String, WebSocketHandler> map = new HashMap<>();
    
                map.put("/ws", new WebSocketHandler() {
                    @Override
                    public Mono<Void> handle(WebSocketSession session) {
                        return session.send(
                                session.receive()
                                      .map(WebSocketMessage::getPayloadAsText)
                                      .map(tMessage -> "Response From Server: " + tMessage)
                                      .map(session::textMessage)
                        );
                    }
                });
    
                SimpleUrlHandlerMapping mapping = new SimpleUrlHandlerMapping();
                mapping.setUrlMap(map);
                mapping.setOrder(-1);
                return mapping;
            }
    
            @Bean
            HandlerAdapter wsHandlerAdapter() {
                return new WebSocketHandlerAdapter();
            }
        }
    

    如果SpringBoot infra不是这样的话

    尝试考虑与 ReactorNetty 直接交互。 Reactor Netty 围绕原生 Netty 提供了非常好的抽象,您可以以相同的功能方式与它进行交互:

    ReactorHttpHandlerAdapter handler =
                        new ReactorHttpHandlerAdapter(yourHttpHandlers);
    
                HttpServer.create()
                          .startRouterAndAwait(routes -> {
                                      routes.ws("/pathToWs", (in, out) -> out.send(in.receive()))
                                            .file("/static/**", ...)
                                            .get("**", handler)
                                            .post("**", handler)
                                            .put("**", handler)
                                            .delete("**", handler);
                                  }
                          );
    

    【讨论】:

    • 谢谢奥莱。我正在尝试将 webflux 集成到现有的 Spring 应用程序中,因此不能选择使用 SpringBoot 或注释。我必须要么使用 xml 配置应用程序上下文,要么使用功能路由。我会试试你的第二个建议。您能否详细说明这一行: (in, out) -> out.send(in.receive())) 。我假设您指的是反应堆 websocketinbound 和 webcoketoutboud。我应该在我当前的 websocket 处理程序中实现任何反应器接口吗?
    【解决方案2】:

    我是这样处理的。并使用原生 reactor-netty

    routes.get(rootPath, (req, resp)->{
            // doFilter check the error
            return this.doFilter(request, response, new RequestAttribute())
                    .flatMap(requestAttribute -> {
                        WebSocketServerHandle handleObject = injector.getInstance(GameWsHandle.class);
                        return response
                            .header("content-type", "text/plain")
                            .sendWebsocket((in, out) ->
                                this.websocketPublisher3(in, out, handleObject, requestAttribute)
                            );
                    });
        })
    
    private Publisher<Void> websocketPublisher3(WebsocketInbound in, WebsocketOutbound out, WebSocketServerHandle handleObject, RequestAttribute requestAttribute) {
            return out
                    .withConnection(conn -> {
                        // on connect
                        handleObject.onConnect(conn.channel());
                        conn.channel().attr(AttributeKey.valueOf("request-attribute")).set(requestAttribute);
                        conn.onDispose().subscribe(null, null, () -> {
                                conn.channel().close();
                                handleObject.disconnect(conn.channel());
                                // System.out.println("context.onClose() completed");
                            }
                        );
                        // get message
                        in.aggregateFrames()
                                .receiveFrames()
                                .map(frame -> {
                                    if (frame instanceof TextWebSocketFrame) {
                                        handleObject.onTextMessage((TextWebSocketFrame) frame, conn.channel());
                                    } else if (frame instanceof BinaryWebSocketFrame) {
                                        handleObject.onBinaryMessage((BinaryWebSocketFrame) frame, conn.channel());
                                    } else if (frame instanceof PingWebSocketFrame) {
                                        handleObject.onPingMessage((PingWebSocketFrame) frame, conn.channel());
                                    } else if (frame instanceof PongWebSocketFrame) {
                                        handleObject.onPongMessage((PongWebSocketFrame) frame, conn.channel());
                                    } else if (frame instanceof CloseWebSocketFrame) {
                                        conn.channel().close();
                                        handleObject.disconnect(conn.channel());
                                    }
                                    return "";
                                })
                                .blockLast();
                    });
        }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-11-20
      • 1970-01-01
      • 2019-01-03
      • 2022-06-05
      • 2018-10-26
      • 1970-01-01
      • 2019-04-24
      • 1970-01-01
      相关资源
      最近更新 更多