如何在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驱动实现非阻塞监听,再推送给客户端。
实现思路
- 创建PostgreSQL触发器,在实体表的插入/更新操作后发送包含变更数据的通知
- 服务器端用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并推送给客户端。
实现思路
- 用数据库触发器或Change Streams(MongoDB)将变更事件写入Redis Stream
- 服务器端用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
相关产品推荐
相关产品推荐

