You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot SSE动态响应实现:定时空响应+实时新用户推送

How to Combine Periodic Heartbeats and Real-Time Database Updates for SSE in Spring Boot

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 empty UserDto is 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 07:34:02