Spring Boot SSE动态响应实现:定时空响应+实时新用户推送
Hey there! Let's tackle this problem step by step. First off, your concept isn't wrong—you absolutely can combine periodic heartbeat messages (empty UserDto every 20 seconds) with real-time updates when new users are added to the database. Your current implementation uses a 20-second poll, which is inefficient and not truly real-time. Here's how to fix it:
1. Shift from Polling to Real-Time Database Change Listening
Instead of checking the database every 20 seconds, we'll set up a way to listen for actual user creation events and push them immediately. For Spring Boot, the cleanest way to do this is using Spring's event system combined with JPA entity listeners (or reactive database change streams if you're using R2DBC).
Step 1: Create a User Creation Event
First, define an event to carry the new user data:
public class UserCreatedEvent { private final UserDto userDto; public UserCreatedEvent(UserDto userDto) { this.userDto = userDto; } public UserDto getUserDto() { return userDto; } }
Step 2: Add an Entity Listener to Trigger Events
Attach a listener to your User entity to fire an event whenever a new user is persisted:
@Entity @EntityListeners(UserEntityListener.class) public class User { // Your existing entity fields and methods } @Component public class UserEntityListener { @Autowired private ApplicationEventPublisher eventPublisher; @PostPersist public void onUserCreated(User user) { // Convert your User entity to UserDto (add your conversion logic here) UserDto userDto = mapUserToDto(user); eventPublisher.publishEvent(new UserCreatedEvent(userDto)); } private UserDto mapUserToDto(User user) { UserDto dto = new UserDto(); // Populate dto fields from user entity return dto; } }
2. Create the Two Flux Streams
Now we'll build two separate streams: one for real-time user updates, and one for the periodic heartbeat.
Real-Time Update Stream
This stream listens for UserCreatedEvent instances and emits the corresponding UserDto:
private Flux<UserDto> createRealTimeUserStream(ApplicationEventPublisher eventPublisher) { return Flux.create(sink -> { // Define a listener that pushes new users to the stream ApplicationListener<UserCreatedEvent> eventListener = event -> { sink.next(event.getUserDto()); }; // Register the listener eventPublisher.addApplicationListener(eventListener); // Clean up the listener when the stream is closed (prevents memory leaks) sink.onDispose(() -> eventPublisher.removeApplicationListener(eventListener)); }); }
Periodic Heartbeat Stream
This stream emits an empty UserDto every 20 seconds:
private Flux<UserDto> createHeartbeatStream() { return Flux.interval(Duration.ofSeconds(20)) .map(tick -> new UserDto()); // Empty UserDto for heartbeat }
3. Merge the Streams
Combine the two streams so both real-time updates and heartbeats are sent to the client:
@Autowired private ApplicationEventPublisher eventPublisher; public Flux<UserDto> fetchUserEvent() { Flux<UserDto> realTimeStream = createRealTimeUserStream(eventPublisher); Flux<UserDto> heartbeatStream = createHeartbeatStream(); // Merge the streams: real-time updates are pushed immediately, heartbeats run on schedule return Flux.merge(realTimeStream, heartbeatStream); }
Key Notes
- Why this works:
Flux.merge()combines multiple streams into one, emitting items from any of the source streams as they become available. So when a new user is created, it's pushed right away; every 20 seconds, an emptyUserDtois sent regardless of activity. - Reactive Databases: If you're using a reactive database like PostgreSQL with R2DBC, you can use database-native change streams (e.g., PostgreSQL's
LISTEN/NOTIFY) instead of JPA entity listeners for even more efficient real-time updates. - Avoiding Memory Leaks: The
sink.onDispose()callback ensures we remove the event listener when the client disconnects, preventing unused listeners from cluttering memory.
内容的提问来源于stack exchange,提问作者Madhu Sudhan Reddy

