【发布时间】:2020-03-02 17:27:58
【问题描述】:
我有一个大的 ConcurrentHashMap (cache.getCache()),其中保存了我的所有数据(大约 500+ MB 大小,但随着时间的推移会增长)。客户端可以通过使用普通 java HttpServer 实现的 API 访问它。
这是简化的代码:
JsonWriter jsonWriter = new JsonWriter(new OutputStreamWriter(new BufferedOutputStream(new GZIPOutputStream(exchange.getResponseBody())))));
new GsonBuilder().create().toJson(cache.getCache(), CacheContainer.class, jsonWriter);
还有一些客户端发送的过滤器,因此它们实际上并没有每次都获取所有数据,但是 HashMap 会不断更新,因此客户端必须经常刷新才能获得最新数据。这是低效的,所以我决定使用 WebSockets 将数据更新实时推送到客户端。
为此我选择了 Undertow,因为我可以简单地从 Maven 导入它,而且我不需要在服务器上进行额外的配置。
在 WS 连接上,我将通道添加到 HashSet 并发送整个数据集(客户端在获取初始数据之前发送带有一些过滤器的消息,但我从示例中删除了这部分):
public class MyConnectionCallback implements WebSocketConnectionCallback {
CacheContainer cache;
Set<WebSocketChannel> clients = new HashSet<>();
BlockingQueue<String> queue = new LinkedBlockingQueue<>();
public MyConnectionCallback(CacheContainer cache) {
this.cache = cache;
Thread pusherThread = new Thread(() -> {
while (true) {
push(queue.take());
}
});
pusherThread.start();
}
public void onConnect(WebSocketHttpExchange webSocketHttpExchange, WebSocketChannel webSocketChannel) {
webSocketChannel.getReceiveSetter().set(new AbstractReceiveListener() {
protected void onFullTextMessage(WebSocketChannel channel, BufferedTextMessage message) {
clients.add(webSocketChannel);
WebSockets.sendText(gson.toJson(cache.getCache()), webSocketChannel, null);
}
}
}
private void push(String message) {
Set<WebSocketChannel> closed = new HashSet<>();
clients.forEach((webSocketChannel) -> {
if (webSocketChannel.isOpen()) {
WebSockets.sendText(message, webSocketChannel, null);
} else {
closed.add(webSocketChannel);
}
}
closed.foreach(clients::remove);
}
public void putMessage(String message) {
queue.put(message);
}
}
每次更改缓存后,我都会获取新值并将其放入队列中(我不直接序列化 myUpdate 对象,因为在 updateCache 方法中还有其他逻辑)。只有一个线程负责更新缓存:
cache.updateCache(key, myUpdate);
Map<Key,Value> tempMap = new HashMap<>();
tempMap.put(key, cache.getValue(key));
webSocketServer.putMessage(gson.toJson(tempMap));
我用这种方法看到的问题:
- 在初始连接时,整个数据集被转换为字符串,我担心过多的请求会导致服务器 OOM。 WebSockets.sendText 只接受 String 和 ByteBuffer
- 如果我先将通道添加到客户端设置然后发送数据,则可能在发送初始数据之前推送到客户端,客户端将处于无效状态
- 如果我先发送初始数据,然后将通道添加到客户端集合中,那么在发送初始数据过程中来的推送消息将丢失,客户端将处于无效状态
我为问题 #2 和 #3 提出的解决方案是将消息放入队列中(我会将 Set<WebSocketChannel> 转换为 Map<WebSocketChannel,Queue<String>> 并仅在客户端收到初始消息后将消息发送到队列中数据集,但我欢迎在这里提出任何其他建议。
至于问题 #1,我的问题是通过 WebSocket 发送初始数据的最有效方式是什么?例如,使用 JsonWriter 直接写入 WebSocket。
我意识到客户端可以使用 API 进行初始调用并订阅 WebSocket 以进行更改,但是这种方法使客户端负责拥有正确的状态(他们需要订阅 WS、排队 WS 消息、获取初始数据使用 API,然后在获取初始数据后将排队的 WS 消息应用到他们的数据集),我不想把控制权交给他们,因为数据很敏感。
【问题讨论】: