Spring4项目H2变更后通过SSE实时推送数据到Angular的实现咨询
实现方案
基于Spring 4原生的事件驱动机制即可实现同应用内跨层通知,无需引入额外中间件或数据库触发器,具体步骤如下:
1. 自定义数据变更事件
首先创建自定义事件类,继承Spring的ApplicationEvent,用于携带新增/更新的Message实体:
public class MessageChangeEvent extends ApplicationEvent { private final Message message; public MessageChangeEvent(Object source, Message message) { super(source); this.message = message; } public Message getMessage() { return message; } }
2. 写库层注入事件发布器
在执行H2插入、更新操作的Service/DAO层中,注入ApplicationEventPublisher,完成数据库操作后发布事件:
@Service public class MessageService { @Autowired private MessageRepository messageRepository; @Autowired private ApplicationEventPublisher eventPublisher; // 你原本的新增/更新消息方法 public void saveMessage(Message message) { // 执行H2写入操作 messageRepository.save(message); // 写入完成后发布变更事件 eventPublisher.publishEvent(new MessageChangeEvent(this, message)); } }
你的定时线程调用该saveMessage方法时,会自动发布事件。
3. 管理活跃的SseEmitter实例
创建一个单例的Emitter管理组件,用于存储所有当前活跃的SSE连接,注意用线程安全容器避免并发问题:
@Component public class SseEmitterManager { // 用线程安全的List存储所有活跃的Emitter private final CopyOnWriteArrayList<SseEmitter> emitters = new CopyOnWriteArrayList<>(); // 添加新的Emitter public void addEmitter(SseEmitter emitter) { emitters.add(emitter); // Emitter完成/超时/出错时自动移除 emitter.onCompletion(() -> emitters.remove(emitter)); emitter.onTimeout(() -> emitters.remove(emitter)); emitter.onError((e) -> emitters.remove(emitter)); } // 向所有活跃Emitter推送消息 public void broadcast(Message message) { emitters.forEach(emitter -> { try { emitter.send(SseEmitter.event() .data(message) .id(String.valueOf(message.getId()))); } catch (Exception e) { // 发送失败直接移除失效的Emitter emitters.remove(emitter); } }); } }
4. 编写事件监听器
创建事件监听器,接收Message变更事件后调用Emitter管理器广播消息:
@Component public class MessageChangeListener { @Autowired private SseEmitterManager emitterManager; // Spring 4.2及以上支持@EventListener注解,更早版本可实现ApplicationListener<MessageChangeEvent>接口 @EventListener public void handleMessageChange(MessageChangeEvent event) { emitterManager.broadcast(event.getMessage()); } }
5. 改造原有SSE接口
修改你的控制器代码,发送完历史消息后将Emitter加入管理器,保持长连接等待后续推送:
@RestController public class SseController { @Autowired private MessageRepository messageRepository; @Autowired private SseEmitterManager emitterManager; // 全局复用线程池,不要每次请求都创建新的 private static final ExecutorService SSE_EXECUTOR = Executors.newCachedThreadPool(); @GetMapping("/sse") public SseEmitter streamSseMvc() { // 设置超时时间,比如1小时,也可以设为-1永不超时,根据业务调整 SseEmitter emitter = new SseEmitter(3600 * 1000L); emitterManager.addEmitter(emitter); SSE_EXECUTOR.execute(() -> { try { List<Message> messages = messageRepository.findAll(); for (int i = 0; i < messages.size(); i++) { SseEmitter.SseEventBuilder event = SseEmitter.event() .data(messages.get(i)) .id(String.valueOf(i)); emitter.send(event); Thread.sleep(1000); } } catch (Exception ex) { emitter.completeWithError(ex); } }); return emitter; } }
内容的提问来源于stack exchange,提问作者Master-Antonio
相关产品推荐
相关产品推荐

