它是绝对正确的设计,它值得为所有可能的操作节省资源,并且只为每个客户机使用一个连接。
但是,不要实现一个轮子并使用提供所有这些通信的协议。
-
RSOCK有一个
请求-响应
这个模型允许您今天进行最常见的客户机-服务器交互。
-
RSOCK有一个
请求流
通信模型,这样您就可以满足所有的需要,并异步地返回一个事件流,重用同一个连接。rsocket将所有逻辑流映射到网络连接和网络连接,因此您自己不会感到这样做的痛苦。
-
rsocket有更多的交互模型,比如
火忘
和
溪流
如果是这样的话
以两种方式发送数据流。
如何在弹簧中使用RSocket
这样做的一个选择是使用RSotoJava实现RSoT协议。RSOCK Java是建立在项目反应堆之上的,因此它自然适合于春季WebF通量生态系统。
不幸的是,没有与春季生态系统的特色集成。幸运的是,我花了几个小时来提供
RSocket Spring Boot Starter
它将SpringWebFlux与RSocket集成,并将WebSocket RSocket服务器与WebFlux HTTP服务器一起公开。
为什么RSocket是更好的方法?
基本上,rsocket隐藏了自己实现相同方法的复杂性。对于RSocket,我们不必关心交互模型定义作为自定义协议和在爪哇中的实现。rsocket为我们提供数据到特定的逻辑通道。它提供了一个将消息发送到同一个WS-Connection的内置客户端,因此我们不必为此发明自定义实现。
因为rsocket只是一个协议,它不提供任何消息格式,所以这个挑战是针对业务逻辑的。但是,有一个rsocket-rpc项目,它提供了一个协议缓冲区作为消息格式,并重用了与grpc相同的代码生成技术。因此,使用rsocket-rpc,我们可以很容易地为客户机和服务器构建一个API,并且完全不考虑传输和协议抽象。
同一个RSocket弹簧引导集成提供了
example
以及rsocket-rpc的用法。
好吧,它还没说服我,我还想有一个自定义的WebSocket服务器
因此,为了达到这个目的,你必须自己去实现这个地狱。我以前已经做过一次,但是我不能指出那个项目,因为它是一个企业项目。
不过,我可以共享一些代码示例,这些示例可以帮助您构建适当的客户机和服务器。
服务器端
处理程序和打开逻辑订阅服务器映射
必须考虑的第一点是,一个物理连接中的所有逻辑流都应存储在某个位置:
class MyWebSocketRouter implements WebSocketHandler {
final Map<String, EnumMap<ActionMessage.Type, ChannelHandler>> channelsMapping;
@Override
public Mono<Void> handle(WebSocketSession session) {
final Map<String, Disposable> channelsIdsToDisposableMap = new HashMap<>();
...
}
}
上面的示例中有两个地图。第一个是您的路由映射,它允许您根据传入的消息参数来标识路由。第二个是为请求流用例创建的(在我的例子中,它是活动订阅的映射),因此您可以发送一个消息框架来创建订阅,或者订阅特定的操作并保留该订阅,这样一旦执行了取消订阅操作,如果存在订阅,您将被取消订阅。
使用处理器进行消息多路复用
为了从所有逻辑流发回消息,必须将消息多路复用到一个流。例如,使用reactor,可以使用
UnicastProcessor
:
@Override
public Mono<Void> handle(WebSocketSession session) {
final UnicastProcessor<ResponseMessage<?>> funIn = UnicastProcessor.create(Queues.<ResponseMessage<?>>unboundedMultiproducer().get());
...
return Mono
.subscriberContext()
.flatMap(context -> Flux.merge(
session
.receive()
...
.cast(ActionMessage.class)
.publishOn(Schedulers.parallel())
.doOnNext(am -> {
switch (am.type) {
case CREATE:
case UPDATE:
case CANCEL: {
...
}
case SUBSCRIBE: {
Flux<ResponseMessage<?>> flux = Flux
.from(
channelsMapping.get(am.getChannelId())
.get(ActionMessage.Type.SUBSCRIBE)
.handle(am) // returns Publisher<>
);
if (flux != null) {
channelsIdsToDisposableMap.compute(
am.getChannelId() + am.getSymbol(), // you can generate a uniq uuid on the client side if needed
(cid, disposable) -> {
...
return flux
.subscriberContext(context)
.subscribe(
funIn::onNext, // send message to a Processor manually
e -> {
funIn.onNext(
new ResponseMessage<>( // send errors as a messages to Processor here
0,
e.getMessage(),
...
ResponseMessage.Type.ERROR
)
);
}
);
}
);
}
return;
}
case UNSABSCRIBE: {
Disposable disposable = channelsIdsToDisposableMap.get(am.getChannelId() + am.getSymbol());
if (disposable != null) {
disposable.dispose();
}
}
}
})
.then(Mono.empty()),
funIn
...
.map(p -> new WebSocketMessage(WebSocketMessage.Type.TEXT, p))
.as(session::send)
).then()
);
}
正如我们从上面的示例中看到的,这里有很多东西:
-
信息应包括路线信息
-
消息应包含与之相关的唯一流ID。
-
用于消息复用的独立处理器,其中错误也应为消息
-
每个通道都应该存储在某个地方,在这种情况下,我们都有一个简单的用例,其中每个消息都可以提供
Flux
或者只是一个
Mono
(在Mono的情况下,它可以在服务器端实现得更简单,因此您不必保留唯一的流ID)。
-
此示例不包括消息编码解码,因此此挑战留给您。
客户端
客户也不是那么简单:
句柄会话
为了处理连接,我们必须分配两个处理器,以便进一步使用它们来多路复用和解复用消息:
UnicastProcessor<> outgoing = ...
UnicastPorcessor<> incoming = ...
(session) -> {
return Flux.merge(
session.receive()
.subscribeWith(incoming)
.then(Mono.empty()),
session.send(outgoing)
).then();
}
将所有逻辑流保留在某个位置
所有创建的流
单声道
或
通量
应该存储在某个地方,以便我们能够区分流消息与哪个相关:
Map<String, MonoSink> monoSinksMap = ...;
Map<String, FluxSink> fluxSinksMap = ...;
由于Monosink,我们必须保留两个映射,而FluxSink没有相同的父接口。
消息路由
在上面的示例中,我们只是考虑了客户端的初始部分。现在我们必须构建一个消息路由机制:
...
.subscribeWith(incoming)
.doOnNext(message -> {
if (monoSinkMap.containsKey(message.getStreamId())) {
MonoSink sink = monoSinkMap.get(message.getStreamId());
monoSinkMap.remove(message.getStreamId());
if (message.getType() == SUCCESS) {
sink.success(message.getData());
}
else {
sink.error(message.getCause());
}
} else if (fluxSinkMap.containsKey(message.getStreamId())) {
FluxSink sink = fluxSinkMap.get(message.getStreamId());
if (message.getType() == NEXT) {
sink.next(message.getData());
}
else if (message.getType() == COMPLETE) {
fluxSinkMap.remove(message.getStreamId());
sink.next(message.getData());
sink.complete();
}
else {
fluxSinkMap.remove(message.getStreamId());
sink.error(message.getCause());
}
}
})
上面的代码示例显示了如何路由传入消息。
多路传输请求
最后一部分是消息复用。为此,我们将讨论可能的发送方类impl:
class Sender {
UnicastProcessor<> outgoing = ...
UnicastPorcessor<> incoming = ...
Map<String, MonoSink> monoSinksMap = ...;
Map<String, FluxSink> fluxSinksMap = ...;
public Sender () {
//在此处创建WebSocket连接并放置前面提到的代码
}
Mono<R> sendForMono(T data) {
//generate message with unique
return Mono.<R>create(sink -> {
monoSinksMap.put(streamId, sink);
outgoing.onNext(message); // send message to server only when subscribed to Mono
});
}
Flux<R> sendForFlux(T data) {
return Flux.<R>create(sink -> {
fluxSinksMap.put(streamId, sink);
outgoing.onNext(message); // send message to server only when subscribed to Flux
});
}
}
自定义实现的总结
-
中坚分子
-
没有实施反压力支持,因此这可能是另一个挑战。
-
很容易射到自己的脚
外卖
-
请使用rsocket,不要自己发明协议,这很难!!!!
-
从关键人物那里了解更多关于rsocket的信息-
https://www.youtube.com/watch?v=WVnAbv65uCU
-
从我的一次谈话中了解更多关于rsocket的信息-
https://www.youtube.com/watch?v=XKMyj6arY2A
-
在rsocket之上构建了一个名为proteus的特色框架-您可能对此感兴趣-
https://www.netifi.com/
-
从rsocket协议的核心开发人员了解更多关于proteus的信息-
https://www.google.com/url?sa=t&source=web&rct=j&url=https://m.youtube.com/watch%3Fv%3D_rqQtkIeNIQ&ved=2ahUKEwjpyLTpsLzfAhXDDiwKHUUUA8gQt9IBMAR6BAgNEB8&usg=AOvVaw0B_VdOj42gjr0YrzLLUX1E