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包会比较大。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐