|
1
2
3
4
|
package com.huaheng.pc.stompwebsocket.service;
import com.alibaba.fastjson.JSON;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
|
5
|
import com.huaheng.common.utils.Wrappers;
|
|
6
7
|
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.huaheng.pc.stompwebsocket.WebsocketConstants;
|
|
8
|
import com.huaheng.pc.stompwebsocket.domain.Message;
|
|
9
10
11
12
13
|
import com.huaheng.pc.stompwebsocket.domain.StompMessage;
import com.huaheng.pc.stompwebsocket.domain.StompPayload;
import com.huaheng.pc.stompwebsocket.mapper.StompMessageMapper;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.simp.SimpMessageSendingOperations;
|
|
14
|
import org.springframework.scheduling.annotation.Async;
|
|
15
16
17
18
19
20
|
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
|
|
21
|
@Service
|
|
22
|
public class StompWebSocketServiceImpl extends ServiceImpl<StompMessageMapper, StompMessage> implements StompWebSocketService {
|
|
23
24
25
26
27
|
@Autowired
private SimpMessageSendingOperations messaging;
/**
|
|
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
|
* 快速推荐消息到用户
* @param user
* @param title
* @param content
*/
public void quickSendToUser(String user, String title, String content){
StompPayload<Message> p = new StompPayload<>();
p.addUser(user);
Message msg = new Message();
msg.setTitle(title);
msg.setContent(content);
p.setData(msg);
p.setQos(1);
send(WebsocketConstants.TOPIC_MESSAGE, p);
}
public void quickBroadcast(String title, String content){
StompPayload<Message> p = new StompPayload<>();
Message msg = new Message();
msg.setTitle(title);
msg.setContent(content);
p.setData(msg);
p.setQos(0);
send(WebsocketConstants.TOPIC_MESSAGE, p);
}
/**
|
|
55
|
* 发送消息
|
|
56
|
*
|
|
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
|
* @param topic
* @param payload
* @param <T>
*/
public <T> void send(String topic, StompPayload<T> payload) {
List<String> sendTo = payload.getSendTo();
if (sendTo == null) {
//没有指定sendTo则广播发送
sendMessageBroadcast(topic, payload.getData());
} else {
//遍历发送消息给指定用户
for (String loginName : payload.getSendTo()) {
sendMessageToUser(topic, payload, loginName);
}
}
}
/**
* 广播发送消息
*
* @param topic
* @param data
*/
public <T> void sendMessageBroadcast(String topic, T data) {
messaging.convertAndSend(topic, data);
}
/**
* 给指定用户发消息
*
* @param topic
* @param payload
|
|
91
|
* @param user loginName
|
|
92
93
94
|
*/
public <T> void sendMessageToUser(String topic, StompPayload<T> payload, String user) {
HashMap<String, Object> headers = new HashMap<>();
|
|
95
|
if (payload.isQos1()) {
|
|
96
97
98
99
100
101
102
103
104
105
106
107
|
StompMessage msg = new StompMessage();
msg.setUser(user);
msg.setTopic(topic);
msg.setMsg(JSON.toJSONString(payload.getData()));
msg.setCreated(new Date());
msg.setLastUpdated(new Date());
save(msg);
headers.put(WebsocketConstants.HEADER_MEG_ID, msg.getId().toString());
}
messaging.convertAndSendToUser(user, topic, payload.getData(), headers);
}
|
|
108
|
public void confirmMsg(String msg_id) {
|
|
109
110
111
|
removeById(msg_id);
}
|
|
112
113
|
/**
* 将离线后的消息发送给用户
|
|
114
|
*
|
|
115
116
|
* @param user
*/
|
|
117
118
119
120
121
122
123
124
125
126
127
128
129
130
|
@Async
public void sendOfflineMessageToUser(String user) {
try {
Thread.sleep(2000);
LambdaQueryWrapper<StompMessage> queryWrapper = Wrappers.lambdaQuery();
queryWrapper.eq(StompMessage::getUser, user);
List<StompMessage> msgList = list(queryWrapper);
for (StompMessage msg : msgList) {
messaging.convertAndSendToUser(user, msg.getTopic(), msg.getMsg());
removeById(msg.getId());
}
} catch (Exception e) {
System.out.println("===========================================================【end】===\n\n\n");
System.out.println(" ======> e.getMessage() = " + e.getMessage());
|
|
131
132
133
|
}
}
}
|