如何用Reactor Kafka在同一WebFlux端点实现消息生产与消费?
WebFlux + Reactor Kafka 单端点实现历史通知+实时监听问题
我需要用WebFlux和Reactor Kafka实现一个响应式端点,核心需求是:
- 先返回数据库中存储的用户历史通知
- 保持端点连接,持续监听并推送新产生的通知
目前我拆分出「生产历史消息」和「获取消息」两个端点时功能正常,但合并成单个端点后无法正常工作,求解决方案。
单端点实现代码
Controller
@GetMapping("/notifications/users/{userId}") Flux<ServerSentEvent> getEventsFlux(@PathVariable String userId) { Flux<ReceiverRecord<String, String>> kafkaFlux = kafkaReceiver.receive(); kafkaNotificationService.generateFirstNotification(userId); return kafkaFlux.filter(sN -> sN.value().toLowerCase().contains(userId)) .map(r -> ServerSentEvent.builder().event(r.value()).build()); }
Service
public void generateFirstNotification(String userId) { Flux<Notification> notifications = notificationRepository.findAllByUserId(UUID.fromString(userId)); sender.<Notification>send(notifications .map(i -> SenderRecord.create(new ProducerRecord<>(TOPIC, "Key", objectMapper.convertValue(i, JSONObject.class).toJSONString()), i))) .doOnError(e -> log.error("Send failed", e)) .subscribe(); }
正常工作的双端点代码
@GetMapping("/notifications/produce-messages/{userId}") void getEventsFlux(@PathVariable String userId) { kafkaService.generateFirstNotification(userId); } @GetMapping("/notifications/get-messages/{userId}") Flux<ServerSentEvent> getMessages(@PathVariable String userId) { Flux<ReceiverRecord<String, String>> kafkaFlux = kafkaReceiver.receive(); return kafkaFlux.filter(sN ->sN.value().toLowerCase().contains(userId)) .map(r -> ServerSentEvent.builder().event(r.value()).build()); }
问题原因分析
- 异步订阅脱离响应链:
generateFirstNotification里用了subscribe(),会让历史消息生产逻辑脱离WebFlux响应式上下文,变成独立异步线程执行,无法和Kafka消费流的启动顺序同步,可能出现消费流已启动但历史消息还未发送到Kafka的情况,导致历史消息被遗漏。 - 冷流启动时机问题:
kafkaReceiver.receive()是冷流,只有WebFlux订阅该Flux时才会启动监听。你先调用生产方法再返回消费Flux,可能历史消息发送时消费流还未启动,导致消息无法被捕获。 - 异常处理失效:
subscribe()的异常处理doOnError不在WebFlux响应链中,生产历史消息的异常无法传递到端点,也无法被正确处理。
解决方案
方案一:直接返回历史消息+拼接Kafka实时流(推荐,减少不必要的Kafka中转)
直接从数据库查询历史消息转为ServerSentEvent,再和Kafka实时消息流拼接,先返回历史消息,再持续推送新消息:
@GetMapping("/notifications/users/{userId}") Flux<ServerSentEvent<String>> getEventsFlux(@PathVariable String userId) { // 1. 从数据库获取历史通知,转为ServerSentEvent Flux<ServerSentEvent<String>> historyEvents = notificationRepository.findAllByUserId(UUID.fromString(userId)) .map(notification -> { String json = objectMapper.convertValue(notification, JSONObject.class).toJSONString(); return ServerSentEvent.<String>builder() .event("history_notification") .data(json) .build(); }); // 2. 监听Kafka的新通知,转为ServerSentEvent Flux<ServerSentEvent<String>> realtimeEvents = kafkaReceiver.receive() .filter(record -> record.value().toLowerCase().contains(userId)) .map(record -> { record.receiveAcknowledgement().acknowledge(); // 手动确认消费,避免重复推送 return ServerSentEvent.<String>builder() .event("new_notification") .data(record.value()) .build(); }); // 3. 先返回历史消息,再持续推送新消息 return Flux.concat(historyEvents, realtimeEvents); }
此方案可直接移除原generateFirstNotification方法,无需将历史消息发送到Kafka,减少中间环节。
方案二:保留Kafka中转历史消息的逻辑(适配原有业务)
如果必须将历史消息发送到Kafka再消费,需将生产逻辑和消费逻辑整合到同一响应链中,避免使用subscribe():
修改Service方法
public Flux<Void> generateFirstNotification(String userId) { Flux<Notification> notifications = notificationRepository.findAllByUserId(UUID.fromString(userId)); return sender.<Notification>send(notifications .map(i -> SenderRecord.create(new ProducerRecord<>(TOPIC, "Key", objectMapper.convertValue(i, JSONObject.class).toJSONString()), i))) .doOnError(e -> log.error("Send failed", e)) .then(); // 转为Flux<Void>,表示生产完成 }
修改Controller
@GetMapping("/notifications/users/{userId}") Flux<ServerSentEvent<String>> getEventsFlux(@PathVariable String userId) { // 先完成历史消息生产,再启动Kafka消费流 return kafkaNotificationService.generateFirstNotification(userId) .thenMany(kafkaReceiver.receive()) .filter(record -> record.value().toLowerCase().contains(userId)) .map(record -> { record.receiveAcknowledgement().acknowledge(); return ServerSentEvent.<String>builder() .event(record.value()) .build(); }); }
此方案保证历史消息生产完成后再启动消费流,确保消息能被正常捕获。
内容的提问来源于stack exchange,提问作者JLopez
相关产品推荐
相关产品推荐

