【问题标题】:How to call SseEmitter just once from spring endpoint如何从spring端点只调用一次SseEmitter
【发布时间】:2020-10-16 21:32:14
【问题描述】:

我正在使用 JavaScript 收听我的 /score 端点,如下所示:

var sse = new EventSource('/score');
sse.onmessage = function (evt) {
    var el = document.getElementById('scores');
    el.appendChild(document.createTextNode(evt.data));
    el.appendChild(document.createElement('br'));
};

但由于某种原因,它就像每秒调用的端点。

EventSource.onmessage 文档说:

是一个EventHandler在收到消息事件时调用,即消息来自源时

这是我的/score 端点:

private ExecutorService executorService = Executors.newCachedThreadPool();

@GetMapping("/score")
public SseEmitter getScore() {
 final SseEmitter sseEmitter = new SseEmitter();

 executorService.submit(() -> {
  try {

   //System.out.println(text);
   //sseEmitter.send(text);
   sseEmitter.send("ok");
   sseEmitter.complete();


  } catch (Exception e) {
   e.printStackTrace();
  }
 });

 return sseEmitter;
}

只有在我手动发送请求时才能触发它?

【问题讨论】:

    标签: javascript java spring-boot spring-mvc


    【解决方案1】:

    经过大量阅读后,我设法在每个请求上触发 /score 端点,但我不得不大量更改服务器端。

     @Autowired
        private MessageProcessor processor;
    
     @GetMapping(path = "/score", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
        public Flux<MyObject> receive() {
            return Flux.create(sink -> {
                processor.register(sink::next);
            });
        }
    

    现在我返回一个 Flux 对象而不是 SseEmitter ,因为使用 Emitter 我需要不断地向客户端发送响应。

    我还创建了另一个端点 /send ,在那里我使用 POST 发送我的对象

     @PostMapping("/send")
        public String send(@RequestBody MyObject event) {
            LOGGER.info("Received '{}'", event);
            processor.process(event);
            return "Done";
        }
    

    Client 端没有任何更改,/receiveEventSource 之间的管道不应从客户端终止。我只添加了 JSON 解析,因为现在我有了一个自定义对象。

    eventSource = new EventSource("/score");
    
    eventSource.onmessage = function (evt) {
      var obj =JSON.parse (evt.data);
            
      var el = document.getElementById('scores');
       el.appendChild(document.createTextNode(obj.id+' '+obj.name+' '+obj.desc));
       el.appendChild(document.createElement('br'));
    };
    

    关键部分是使用 MessageProcessor

    @Service
    public class MessageProcessor {
    
     private static final Logger LOGGER = LoggerFactory.getLogger(MessageProcessor.class);
    
     private List < Consumer < MyObject >> listeners = new CopyOnWriteArrayList < > ();
    
     public void register(Consumer < MyObject > listener) {
      listeners.add(listener);
      LOGGER.info("Added a listener, for a total of {} listener{}", listeners.size(), listeners.size() > 1 ? "s" : "");
     }
    
     // TODO FBE implement unregister
    
     public void process(MyObject event) {
      System.out.println("Processing: " + event);
      listeners.forEach(c -> c.accept(event));
     }
    }
    

    Webflux的另一个例子可以在here找到

    输出

    • 客户

    • 服务器(我发布的 4 个对象)

    【讨论】:

      【解决方案2】:

      那是因为您从未调用 EventSource.close() 来关闭连接。由于您的 eventSource 没有终止条件,您的浏览器将重复调用 /score

      【讨论】:

      • 我正在寻找服务器端解决方案,如果我关闭它然后触发event 怎么办? ,我想不断地收听端点
      • 您发布的代码是客户端问题。您的服务器代码按预期运行。它发送一个“ok”,然后 sse 连接从服务器端终止。您不断看到重复“ok”的原因是从客户端 EventSource 重复重试的结果。
      • 我无法调用 close(),因为我需要继续收听该端点,如果我不能只调用一次,可能无法处理来自服务器端的消息
      • 所以您正在寻找服务器端解决方案来解决什么问题?你不清楚解释你想做什么。
      • 同上,这是客户端问题。您的服务器代码仅发送一次 ok 然后终止。但是您的客户不断尝试通过调用/score 重新连接。您的服务器无法控制您的客户端发出多少请求。如果您只想获得一次ok,则事件流不适合并坚持使用 Http 1.1 的常规 req/resp
      猜你喜欢
      • 1970-01-01
      • 2019-04-01
      • 1970-01-01
      • 1970-01-01
      • 2012-06-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-06-28
      相关资源
      最近更新 更多