Java:如何实现SSE(Server-Sent Events)功能
·
SSE(Server-Sent Events)现在使用的范围非常广,现在提供3种实现SSE功能的方法,供以后参考。
1、使用Servlet 3.0
import java.io.PrintWriter;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import javax.servlet.http.HttpServlet;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
public class SseServlet extends HttpServlet {
private static final long serialVersionUID = -2006037095834838462L;
private final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("YYYY-MM-dd HH:mm:dd");
protected void doGet(HttpServletRequest request, HttpServletResponse response) {
String countStr = (String) request.getParameter("count");
String delayStr = (String) request.getParameter("delay");
String initDelayStr = (String) request.getParameter("initDelay");
int count = countStr == null ? 60 : Integer.parseInt(countStr);
long delay = delayStr == null ? 10000 : Long.parseLong(delayStr);
long initDelay = initDelayStr == null ? 5000 : Long.parseLong(initDelayStr);
response.setContentType("text/event-stream");
response.setCharacterEncoding("UTF-8");
response.setHeader("Cache-Control", "no-cache");
response.setHeader("Connection", "keep-alive");
response.setHeader("Access-Control-Allow-Orign", "*");
PrintWriter writer = null;
try {
writer = response.getWriter();
Thread.sleep(initDelay);
String now = formatter.format(LocalDateTime.now());
writer.write("id:0\n");
writer.write("event:test\n");
writer.write("data:welcome-" + now + "\n\n");
writer.flush();
for (int i = 1; i <= count; i++) {
if (writer.checkError()) {
System.out.println("write error");
return;
}
Thread.sleep(delay);
now = formatter.format(LocalDateTime.now());
writer.write("id:" + i + "\n");
writer.write("event:test\n");
writer.write("data:msg-" + now + "\n\n");
writer.flush();
}
} catch (Exception e) {
if (writer != null) {
writer.flush();
}
} finally {
if (writer != null) {
writer.close();
}
}
}
}
以上是使用servlet的同步API实现的SSE功能,应该也可以使用servlet的异步API实现这个功能。这种方式使用web容器的原生API实现SSE功能,不需要用到其它开源组件。
2、使用spring-webmvc
pom.xml
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@RestController
@RequestMapping(path = "/sse")
public class SseController {
private final ExecutorService es = Executors.newFixedThreadPool(10);
private final DateTimeFormatter f = DateTimeFormatter.ofPattern("YYY-MM-dd HH:mm:ss");
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter stream() {
SseEmitter sseEmitter = new SseEmitter(Long.MAX_VALUE);
es.submit(() -> {
try {
for (int i = 0; i < 60; i++) {
Thread.sleep(5000);
String now = f.format(LocalDateTime.now());
sseEmitter.send(SseEmitter.event().id(String.valueOf(i)).name("sse-event").data("==" + now));
}
sseEmitter.complete();
} catch (Exception e) {
e.printStackTrace();
}
});
return sseEmitter;
}
}
这种实现方式应该是最简单的方式,需要写的代码最少,但需要依赖spring-boot和spring-mvc。
3、使用spring-boot + JAX-RS 2.1
pom.xml
<dependencies>
<dependency>
<groupId>org.glassfish.jersey.media</groupId>
<artifactId>jersey-media-sse</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jersey</artifactId>
</dependency>
</dependencies>
import java.time.LocalTime;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import javax.ws.rs.GET;
import javax.ws.rs.Path;
import javax.ws.rs.Produces;
import javax.ws.rs.core.Context;
import javax.ws.rs.core.MediaType;
import javax.ws.rs.sse.OutboundSseEvent;
import javax.ws.rs.sse.Sse;
import javax.ws.rs.sse.SseEventSink;
@Path("/sse")
public class SseResource {
private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
@GET
@Path("/stream")
@Produces(MediaType.SERVER_SENT_EVENTS)
public void stream(@Context Sse sse, @Context SseEventSink sink) {
OutboundSseEvent.Builder builder = sse.newEventBuilder();
try {
// 发送欢迎消息
sink.send(builder.name("init").data("Connected!").build());
} catch (Exception e) {
return;
}
// 定时推送
scheduler.scheduleAtFixedRate(() -> {
if (sink.isClosed())
return; // 连接断开则停止
try {
String msg = "Time: " + LocalTime.now();
sink.send(builder.id(System.currentTimeMillis() + "").name("tick").data(msg).build());
} catch (Exception e) {
sink.close();
}
}, 1, 2, TimeUnit.SECONDS);
}
}
import org.glassfish.jersey.server.ResourceConfig;
import org.springframework.stereotype.Component;
import javax.ws.rs.ApplicationPath;
@Component
@ApplicationPath("/api")
public class JerseyConfig extends ResourceConfig {
public JerseyConfig() {
register(SseResource.class);
}
}
这种实现方式也是非常简单的方式,需要写的代码也较少,但需要引入JAX-RS 2.1相关的开源组件,最终打出的jar包会比较大。
更多推荐



所有评论(0)