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

如何在Spring Reactive中通过RSocket/WebSocket/Http永久保持响应式流连接?

可行实现方案推荐

结合你的需求(异步非阻塞、数据库变更实时推送给客户端、支持数据编辑),以下是几个经过验证的方案,适配你已尝试的技术栈:

1. MongoDB Change Streams + RSocket/WebFlux SSE

这是替代MongoDB capped集合的最优方案,支持普通集合的所有CRUD操作监听,不需要固定容量,完全契合你的需求。

实现思路

MongoDB的Change Streams可以监听指定集合的所有变更事件(插入、更新、删除等),Spring Data MongoDB Reactive提供了原生的Reactive API来消费这些事件,再通过RSocket或WebFlux SSE推送给客户端。

核心代码示例

服务层监听变更

@Service
public class EntityChangeService {
    private final ReactiveMongoTemplate mongoTemplate;

    public EntityChangeService(ReactiveMongoTemplate mongoTemplate) {
        this.mongoTemplate = mongoTemplate;
    }

    // 过滤出插入和更新事件,返回实体的Flux流
    public Flux<Entity> getEntityUpdates() {
        return mongoTemplate.changeStream("entities", Entity.class)
                .filter(event -> 
                        ChangeStreamOperationType.INSERT.equals(event.getOperationType()) ||
                        ChangeStreamOperationType.UPDATE.equals(event.getOperationType()))
                .map(ChangeStreamEvent::getBody);
    }
}

RSocket控制器(双向低延迟场景)

@Controller
public class EntityRSocketController {
    private final EntityChangeService changeService;

    public EntityRSocketController(EntityChangeService changeService) {
        this.changeService = changeService;
    }

    @MessageMapping("entity.real-time.updates")
    public Flux<Entity> streamEntityUpdates() {
        return changeService.getEntityUpdates();
    }
}

WebFlux SSE控制器(浏览器客户端场景)

@RestController
@RequestMapping("/api/entities")
public class EntitySseController {
    private final EntityChangeService changeService;

    public EntitySseController(EntityChangeService changeService) {
        this.changeService = changeService;
    }

    @GetMapping(value = "/updates", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<Entity> streamRealTimeUpdates() {
        return changeService.getEntityUpdates();
    }
}

2. PostgreSQL LISTEN/NOTIFY + WebFlux/R2DBC

PostgreSQL原生支持LISTEN/NOTIFY机制,可通过触发器在数据变更时发送通知,结合R2DBC的Reactive驱动实现非阻塞监听,再推送给客户端。

实现思路

  1. 创建PostgreSQL触发器,在实体表的插入/更新操作后发送包含变更数据的通知
  2. 服务器端用R2DBC监听通知频道,解析通知后获取最新实体数据,生成Flux流推送给客户端

核心代码示例

数据库触发器(SQL)

-- 创建通知函数
CREATE OR REPLACE FUNCTION notify_entity_change()
RETURNS TRIGGER AS $$
DECLARE
    payload JSON;
BEGIN
    IF TG_OP = 'INSERT' OR TG_OP = 'UPDATE' THEN
        payload = row_to_json(NEW);
        PERFORM pg_notify('entity_change_channel', payload::TEXT);
    END IF;
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

-- 绑定触发器到实体表
CREATE TRIGGER entity_change_trigger
AFTER INSERT OR UPDATE ON entities
FOR EACH ROW EXECUTE FUNCTION notify_entity_change();

服务层监听通知

@Service
public class EntityChangeService {
    private final ConnectionFactory connectionFactory;
    private final DatabaseClient databaseClient;

    public EntityChangeService(ConnectionFactory connectionFactory, DatabaseClient databaseClient) {
        this.connectionFactory = connectionFactory;
        this.databaseClient = databaseClient;
    }

    public Flux<Entity> getEntityUpdates() {
        return Flux.usingWhen(
                connectionFactory.create(),
                connection -> {
                    // 监听通知频道
                    connection.createStatement("LISTEN entity_change_channel").execute();
                    // 消费通知,查询最新实体数据
                    return Flux.from(connection.getNotifications())
                            .map(notification -> notification.getParameter())
                            .flatMap(payload -> {
                                Long entityId = JSON.parseObject(payload).getLong("id");
                                return databaseClient.select()
                                        .from(Entity.class)
                                        .matching(Criteria.where("id").is(entityId))
                                        .one();
                            });
                },
                Connection::close
        );
    }
}

3. Redis Stream + 数据库变更触发

如果需要更高的可靠性(避免服务器下线丢失消息),可以用Redis Stream作为中间件,数据库变更时将事件写入Stream,服务器端Reactive消费Stream并推送给客户端。

实现思路

  1. 用数据库触发器或Change Streams(MongoDB)将变更事件写入Redis Stream
  2. 服务器端用Reactive Redis API消费Stream,生成实体Flux流推送给客户端

核心代码示例

写入Redis Stream(以MongoDB Change Streams为例)

@Service
public class MongoToRedisStreamService {
    private final ReactiveMongoTemplate mongoTemplate;
    private final ReactiveRedisTemplate<String, String> redisTemplate;
    private final ObjectMapper objectMapper;

    public MongoToRedisStreamService(ReactiveMongoTemplate mongoTemplate, 
                                     ReactiveRedisTemplate<String, String> redisTemplate,
                                     ObjectMapper objectMapper) {
        this.mongoTemplate = mongoTemplate;
        this.redisTemplate = redisTemplate;
        this.objectMapper = objectMapper;
    }

    @PostConstruct
    public void startPublishingChanges() {
        mongoTemplate.changeStream("entities", Entity.class)
                .filter(event -> INSERT.equals(event.getOperationType()) || UPDATE.equals(event.getOperationType()))
                .flatMap(event -> {
                    try {
                        String payload = objectMapper.writeValueAsString(event.getBody());
                        return redisTemplate.opsForStream().add("entity:changes", Map.of("data", payload));
                    } catch (JsonProcessingException e) {
                        return Mono.error(e);
                    }
                })
                .subscribe();
    }
}

消费Redis Stream并推送

@Service
public class EntityChangeService {
    private final ReactiveRedisTemplate<String, String> redisTemplate;
    private final ObjectMapper objectMapper;

    public EntityChangeService(ReactiveRedisTemplate<String, String> redisTemplate, ObjectMapper objectMapper) {
        this.redisTemplate = redisTemplate;
        this.objectMapper = objectMapper;
    }

    public Flux<Entity> getEntityUpdates() {
        StreamOffset<String> offset = StreamOffset.fromStart("entity:changes");
        return redisTemplate.opsForStream()
                .read(StreamReadOptions.empty().block(Duration.ZERO), offset)
                .map(StreamMessage::getValue)
                .map(map -> map.get("data"))
                .flatMap(payload -> {
                    try {
                        return Mono.just(objectMapper.readValue(payload, Entity.class));
                    } catch (JsonProcessingException e) {
                        return Mono.error(e);
                    }
                });
    }
}

方案对比

方案优势劣势适用场景
MongoDB Change Streams集成简单、无额外组件、支持所有集合操作仅限MongoDB单MongoDB场景、快速落地
PostgreSQL LISTEN/NOTIFY轻量无中间件、原生支持需编写SQL触发器、通知数据有限单PostgreSQL场景、轻量级需求
Redis Stream高可靠、持久化消息、支持分布式消费需Redis中间件高可靠性要求、分布式系统

内容的提问来源于stack exchange,提问作者Paulo Rodrigues

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 10:27:28