【问题标题】:Spring Webflux send event when any new dataSpring Webflux 在任何新数据时发送事件
【发布时间】:2021-02-09 10:48:40
【问题描述】:

我正在尝试学习 Spring webflux 和 R2DBC。我尝试的一个是简单的用例:

  1. 有一个book
  2. 创建一个提供文本流并返回Flux<Book>的API (/books)
  3. 我希望当我点击一次/books 时,保持浏览器打开,并且将任何新数据插入到book 表中,它会将新数据发送到浏览器。

场景 2,仍然来自 book 表:

  1. 有一个book
  2. 创建一个 API (/books/count),将 book 中的数据计数返回为 Mono<Long>
  3. 我希望当我点击一次/books/count 时,保持浏览器打开,并将任何新数据插入/删除到book 表中,它会将新计数发送到浏览器。

但它不起作用。在我插入新数据后,没有数据发送到我的任何端点。
我需要点击/books/books/count 来获取更新的数据。
我认为要做到这一点,我需要使用服务器发送事件吗?但是如何做到这一点并查询数据呢?我得到的大多数示例都是简单的 SSE,它每隔一定的时间间隔发送一次字符串。

有任何样本可以做到这一点吗?

这是我的 BookApi.java

@RestController
@RequestMapping(value = "/books")
public class BookApi {

    private final BookRepository bookRepository;

    public BookApi(BookRepository bookRepository) {
        this.bookRepository = bookRepository;
    }

    @GetMapping(produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<Book> getAllBooks() {
        return bookRepository.findAll();
    }

    @GetMapping(value = "/count", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Mono<Long> count() {
        return bookRepository.count();
    }
}

BookRepository.java (R2DBC)

import org.springframework.data.r2dbc.repository.R2dbcRepository;

public interface BookRepository extends R2dbcRepository<Book, Long> {
}

Book.java

@Table("book")
@Data
@AllArgsConstructor
@NoArgsConstructor
public class Book {

    @Id
    private Long id;

    @Column(value = "name")
    private String name;

    @Column(value = "author")
    private String author;

}

【问题讨论】:

  • Mono 是一项延迟项目。当延迟的项目被解决并发送回客户端时,订阅结束并且连接关闭。如果您需要发送多件物品,则需要退回助焊剂。我没有任何示例,但我可以告诉您,您需要订阅一个返回 Flux 的端点。
  • 这个Flux可以例如使用处理器/接收器projectreactor.io/docs/core/release/reference/#processors这样当有人插入一个项目时,在他们返回这个项目之前,他们会触发一个函数,该函数将从db获取最新计数,并将其放入 sink 中,然后 sink 会将新值推送给所有订阅的人。
  • 这是更高级的事情,在开始使用处理器之前,您需要了解接收器的概念。我建议在projectreactor.io/docs/core/release/reference/#producing 处运行反应器示例,了解如何以编程方式创建可以通过 Flux 生成项目的序列。

标签: spring-webflux spring-data-r2dbc


【解决方案1】:

使用处理器或接收器来处理 Book created 事件。

检查my example using reactor Sinks,并阅读this article for the details

或者使用可拖尾的 Mongo 文档。

一个tailable MongoDB文档可以自动完成这项工作,检查同一个repos的主分支。

我上面的例子使用了WebSocket协议,很容易切换到SSERSocket

【讨论】:

    猜你喜欢
    • 2019-10-19
    • 2021-04-22
    • 1970-01-01
    • 1970-01-01
    • 2020-11-22
    • 2019-01-26
    • 1970-01-01
    • 2019-02-09
    • 2020-07-25
    相关资源
    最近更新 更多