如何向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
相关产品推荐
相关产品推荐

