RxGo设置StopOnError后Observable流未停止仍输出后续值问题
核心原因
rxgo.WithErrorStrategy(rxgo.StopOnError) 配置的生效边界是下游操作符的处理链路,不会干预你自定义Producer函数内部的执行逻辑。
你写的Producer是独立运行的goroutine,代码里顺序写死了4次发送动作:发1、发2、发错误、发3,这部分逻辑完全由你自己的代码控制,rxgo框架不会在你发送完错误事件后强行中断Producer函数的执行,所以值3会被正常写入next通道,最终被Observe()遍历到。StopOnError的实际作用是:当流中出现错误事件后,后续所有链式操作符(比如Map、Filter、FlatMap等)不会再处理错误事件之后收到的元素,而不是从Producer侧阻断数据发送。
效果验证
你可以给Observable加一个Map操作符测试,就能看到StopOnError实际是生效的:
observable.Map(func(ctx context.Context, i interface{}) (interface{}, error) { fmt.Printf("处理元素: %v\n", i) return i, nil })
运行后你会发现,值3不会进入Map的处理逻辑,这就是StopOnError在起作用——它截断了错误事件之后的元素向下游传递的路径,只是没有阻止Producer侧主动写入数据。
正确实现错误后停止发送的方式
- 方式一:在自定义Producer逻辑中主动控制流程,发送完错误事件后直接return,不要执行后续发送代码。
- 方式二:监听Producer入参的ctx,流触发错误后ctx会被主动取消,你可以在发送逻辑里判断ctx状态,收到取消信号就终止发送,参考代码:
func(ctx context.Context, next chan<- rxgo.Item) { // 封装安全发送方法,ctx取消就不再发送 safeSend := func(item rxgo.Item) bool { select { case <-ctx.Done(): return false case next <- item: return true } } safeSend(rxgo.Of(1)) safeSend(rxgo.Of(2)) safeSend(rxgo.Error(errors.New("unknown"))) safeSend(rxgo.Of(3)) // 此时ctx已取消,不会真正发送3 }
内容的提问来源于stack exchange,提问作者ohlab
相关产品推荐
相关产品推荐

