时间:2025-05-18 19:53
人气:
作者:admin
我们都知道tcp,ip,http,https,websocket等等协议,今天了解一个新的协议SSE协议(Server-Sent Events)
SSE(Server-Sent Events) 是一种允许服务器主动向客户端推送数据的轻量级协议,基于 HTTP 长连接,实现 单向通信(服务器→客户端)。它是 W3C 标准,浏览器原生支持,无需额外插件(如 EventSource API)
核心特点与工作原理
GET 请求建立连接,服务器返回特殊格式的文本流(text/event-stream),连接保持打开状态,直到服务器主动关闭或超时。\n 分隔,支持事件类型、数据内容、重试时间等字段,例如:data: Hello, SSE! // 数据内容
event: customEvent // 自定义事件类型(可选)
id: 123 // 消息ID(可选)
retry: 5000 // 重连时间(毫秒,可选)
\n
适用于无需双向通信,仅需服务器单向推送数据。【比如现在的 gpt,豆包这个问答形式】
前端客户端可以使用原生的 EventSource API:
// 创建EventSource实例,连接服务器
const eventSource = new EventSource('/sse-endpoint');
// 监听默认事件("message")
eventSource.onmessage = (event) => {
console.log('Received:', event.data);
};
// 监听自定义事件(如"customEvent")
eventSource.addEventListener('customEvent', (event) => {
console.log('Custom Event:', event.data);
});
// 处理错误
eventSource.onerror = (error) => {
console.error('SSE Error:', error);
// 浏览器会自动重连,无需手动处理
};
服务端可用的就太多了。(本文以SpringBoot3.4.2为例子)
在知道这个协议之前,我们想要达到gpt这种问答形式,输出内容是一点一点拼接的,该怎么弄呢?是不是还可以用websocket。
| 特性 | SSE | WebSocket |
|---|---|---|
| 通信方向 | 单向(服务器→客户端) | 双向(全双工) |
| 协议 | 基于 HTTP(升级为长连接) | 独立协议(ws:// 或 wss://) |
| 二进制支持 | 仅文本(text/event-stream) |
支持文本和二进制 |
| 自动重连 | 浏览器内置 | 需手动实现 |
| 复杂度 | 简单(服务端实现轻量) | 较复杂(需处理握手、心跳) |
| 适用场景 | 服务器单向推送数据 | 双向交互(聊天、实时协作) |
下面结合Spring Boot 简单用一下SSE
// sse协议测试
@PostMapping(value = "/chat", produces = "text/event-stream;charset=UTF-8")
public SseEmitter streamSseMvc() {
// 感谢评论区:键盘三个键 指出该timeout问题。
// 有一点需要注意的是:这里的time_out参数,是SseEmitter(session会话)的存活时间. 这一点需要注意一下
SseEmitter emitter = new SseEmitter(30_000L);
// 模拟发送消息
System.out.println("SSE connection started");
ScheduledFuture<?> future = service.scheduleAtFixedRate(() -> {
try {
String message = "Message at " + System.currentTimeMillis();
emitter.send(SseEmitter.event().data(message));
} catch (IOException e) {
try {
emitter.send(SseEmitter.event().name("error").data(Map.of("error", e.getMessage())));
} catch (IOException ex) {
// ignore
}
emitter.completeWithError(e);
}
}, 0, 5, TimeUnit.SECONDS);
emitter.onCompletion(() -> {
System.out.println("SSE connection completed");
});
emitter.onTimeout(() -> {
System.out.println("SSE connection timed out");
emitter.complete();
});
emitter.onError((e) -> {
System.out.println("SSE connection error: " + e.getMessage());
emitter.completeWithError(e);
});
return emitter;
}
在SpringBoot中,用SseEmitter就可以达到这个效果了,它也和Websocket一样有onXXX这种类似的方法。上面是使用一个周期性的任务,来模拟AI生成对话的效果的。emitter.send(SseEmitter.event().data(message)); 这个就是服务端向客户端推送数据。
简单示例:就问一句话
申请deepseekKey这里就略过了,各位读者自行去申请。【因为deepseek官网示例是用的okhttp,所以我这里也用okhttp了】
我们先准备一个接口
@RestController
@RequestMapping("/deepseek")
public class DeepSeekController {
@Resource
private DeepSeekUtil deepSeekUtil;
/**
* 访问deepseek-chat
*/
@PostMapping(value = "/chat", produces = "text/event-stream;charset=UTF-8")
public SseEmitter chatSSE() throws IOException {
SseEmitter emitter = new SseEmitter(60000L);
deepSeekUtil.sendChatReqStream("123456", "你会MySQL数据库吗?", emitter);
return emitter; // 这里把该sse对象返回了
}
private boolean notModel(String model) {
return !"deepseek-chat".equals(model) && !"deepseek-reasoner".equals(model);
}
}
可以看到我们创建了一个SseEmitter对象,传给了我们自定义的工具
@Component
public class DeepSeekUtil {
public static final String DEEPSEEK_CHAT = "deepseek-chat";
public static final String DEEPSEEK_REASONER = "deepseek-reasoner";
public static final String url = "https://api.deepseek.com/chat/completions";
// 存储每个用户的消息列表
private static final ConcurrentHashMap<String, List<Message>> msgList = new ConcurrentHashMap<>();
// 1.调用api,然后以以 SSE(server-sent events)的形式以流式发送消息增量。消息流以 data: [DONE] 结尾。
public void sendChatReqStream(String uid, String message, SseEmitter sseEmitter) throws IOException {
// 1.构建一个普通的聊天body请求
AccessRequest tRequest = buildNormalChatRequest(uid, message);
OkHttpClient client = new OkHttpClient().newBuilder()
.build();
// 封装请求体参数
MediaType mediaType = MediaType.parse("application/json; charset=utf-8");
RequestBody body = RequestBody.create(JSON.toJSONString(tRequest), mediaType);
// 构建请求和请求头
Request request = new Request.Builder()
.url(url)
.method("POST", body)
.addHeader("Content-Type", "application/json")
.addHeader("Accept", "text/event-stream")
// 比如你的key是:s-123456
// .addHeader("Authorization", "Bearer s-123456")
.addHeader("Authorization", "Bearer 你的key")
.build();
// 创建一个监听器
SseChatListener listener = new SseChatListener(sseEmitter);
RealEventSource eventSource = new RealEventSource(request, listener);
eventSource.connect(client);
}
private AccessRequest buildNormalChatRequest(String uid, String message) {
// 这里,我们messages,添加了一条“你会MySQL数据库吗?",来达到一种对话具有上下文的效果
List<Message> messages = msgList.computeIfAbsent(uid, k -> new ArrayList<>());
messages.add(new Message("user", message));
/*
[
{"system", "你好, 我是DeepSeek-AI助手!"},
{"user", "你会MySQL数据库吗?"}
]
*/
AccessRequest request = new AccessRequest();
request.setMessages(messages);
request.setModel(DEEPSEEK_CHAT);
request.setResponse_format(Map.of("type", "text"));
request.setStream(true); // 设置为true
request.setTemperature(1.0);
request.setTop_p(1.0);
return request;
}
@PostConstruct
public void init() {
List<Message> m = new ArrayList<Message>();
m.add(new Message("system", "你好, 我是DeepSeek-AI助手!"));
// 初始化消息列表
msgList.put("123456", m);
}
}
// 请求体,参考deepseek官网
public class AccessRequest {
private List<Message> messages;
private String model; // 默认模型为deepseek-chat
private Double frequency_penalty = 0.0;
private Integer max_tokens;
private Double presence_penalty = 0.0;
//{
// "type": "text"
/
【从0到1构建一个ClaudeAgent】协作-Agent团队