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

如何向Flux推送待发布结果?Reactor事件处理类实现求助

嘿,我明白你遇到的问题了——你想实现一个能从外部异步推送处理结果的EventProcessor,用Flux让订阅者实时收到通知对吧?你之前用Flux.generate的思路踩了个小坑,我来给你捋清楚并给出正确的实现方式。

问题根源:Flux.generate不适合异步推送场景

你遇到的IllegalStateException是因为Flux.generate的设计意图是同步、按需生成元素:每次有订阅者订阅时,Reactor会调用你的generator回调函数,你必须在这个回调里立即调用sink的方法(next/complete)来生成元素,绝对不能把sink保存下来后续复用。它是给那种“订阅者要一个元素,我生成一个”的场景用的,完全不匹配你从外部事件触发异步推送的需求。

正确方案:用EmitterProcessor实现手动推送

Reactor专门提供了用于手动推送元素的处理器,最适合你这个场景的是EmitterProcessor(支持背压,还能配置缓存最近元素给新订阅者)。下面是修正后的完整实现:

import reactor.core.publisher.EmitterProcessor
import reactor.core.publisher.Flux
import reactor.core.publisher.FluxSink

class EventProcessor {
    // 创建EmitterProcessor,可根据需求调整初始缓冲区大小
    private val eventProcessor = EmitterProcessor.create<Result>()
    // 获取用于推送元素的Sink实例
    private val eventSink: FluxSink<Result> = eventProcessor.sink()

    // 对外暴露的Flux,供客户端订阅
    val flux: Flux<Result> = eventProcessor

    fun onUserEvent1(e: Event) {
        val result = process(e)
        // 推送处理后的结果给所有订阅者
        eventSink.next(result)
    }

    fun onUserEvent2(e: Event) {
        val result = process(e)
        eventSink.next(result)
    }

    private fun process(e: Event): Result {
        // 这里替换成你的实际事件处理逻辑
        return Result("Processed event type: ${e.type}")
    }
}

// 假设的Event和Result数据类(根据你的实际需求调整)
data class Event(val type: String)
data class Result(val content: String)

关键细节说明

  • EmitterProcessor本身就是一个Flux,同时提供了FluxSink用于外部手动推送元素,完美适配你的“事件触发→处理→推送结果”流程。
  • 如果你的场景完全不需要背压支持(比如订阅者处理速度肯定能跟上推送速度),也可以用更轻量的DirectProcessor,但要注意它没有缓冲区,推送速度过快可能会丢数据。
  • 记得在EventProcessor生命周期结束时(比如应用关闭),调用eventSink.complete()来结束Flux,或者用eventSink.error(Throwable)传递错误信息给订阅者。

客户端订阅示例

客户端只需要订阅EventProcessor.flux就能实时收到结果:

val processor = EventProcessor()
processor.flux.subscribe { result ->
    println("Received result: ${result.content}")
}

// 模拟触发事件
processor.onUserEvent1(Event("LOGIN"))
processor.onUserEvent2(Event("PURCHASE"))

这样修改后,你的需求就完全满足了——外部事件触发后,处理结果会立即推送给所有订阅者,完全符合响应式编程的设计思路。

内容的提问来源于stack exchange,提问作者Martin Drozdik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:23:34