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

如何用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());
}

问题原因分析

  1. 异步订阅脱离响应链:generateFirstNotification里用了subscribe(),会让历史消息生产逻辑脱离WebFlux响应式上下文,变成独立异步线程执行,无法和Kafka消费流的启动顺序同步,可能出现消费流已启动但历史消息还未发送到Kafka的情况,导致历史消息被遗漏。
  2. 冷流启动时机问题:kafkaReceiver.receive()是冷流,只有WebFlux订阅该Flux时才会启动监听。你先调用生产方法再返回消费Flux,可能历史消息发送时消费流还未启动,导致消息无法被捕获。
  3. 异常处理失效: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 04:24:54