
golang-samber-ro
热门使用 samber/ro 在 Golang 中实现响应式流和事件驱动编程——ReactiveX 实现,包含 150+ 类型安全操作符、冷/热可观察对象、5 种主题类型(Publish、Behavior、Replay、Async、Unicast)、通过 Pipe 声明式管道、40+ 插件(HTTP、cron、fsnotify、JSON、日志)、自动背压、错误传播和 Go 上下文集成。适用于使用或采用 samber/ro 时、代码库导入 github.com/samber/ro 时,或在 Go 中构建异步事件驱动管道、实时数据处理、流或响应式架构时。不适用于有限切片转换(→ 参见 `samber/cc-skills-golang@golang-samber-lo` 技能)。
使用 samber/ro 在 Golang 中实现响应式流和事件驱动编程——ReactiveX 实现,包含 150+ 类型安全操作符、冷/热可观察对象、5 种主题类型(Publish、Behavior、Replay、Async、Unicast)、通过 Pipe 声明式管道、40+ 插件(HTTP、cron、fsnotify、JSON、日志)、自动背压、错误传播和 Go 上下文集成。适用于使用或采用 samber/ro 时、代码库导入 github.com/samber/ro 时,或在 Go 中构建异步事件驱动管道、实时数据处理、流或响应式架构时。不适用于有限切片转换(→ 参见 `samber/cc-skills-golang@golang-samber-lo` 技能)。
角色: 你是一名 Go 工程师,当数据流异步或无限时,你会选择响应式流。你使用 samber/ro 构建声明式管道,而不是手动编写 goroutine/channel 代码,但你知道何时简单的切片加 samber/lo 就足够了。
思考模式: 在设计高级响应式管道或选择冷/热可观察对象、主题和组合操作符时,使用 ultrathink。错误的架构会导致资源泄漏或事件丢失。
samber/ro — Go 的响应式流
Go 实现的 ReactiveX。基于泛型、类型安全、可组合的管道,用于异步数据流,具有自动背压、错误传播、上下文集成和资源清理。150+ 操作符,5 种主题类型,40+ 插件。
官方资源:
本技能并非详尽无遗。请参考库文档和代码示例以获取更多信息。Context7 可作为发现平台提供帮助。对于 Go 包文档、版本、符号和已知漏洞,→ 参见 samber/cc-skills-golang@golang-pkg-go-dev 技能。
为什么选择 samber/ro(流 vs 切片)
Go 的 channel + goroutine 在处理复杂异步管道时变得笨重:手动关闭 channel、冗长的 goroutine 生命周期、跨嵌套 select 的错误传播,以及缺乏可组合的操作符。samber/ro 通过声明式、可链式调用的流操作符解决了这些问题。
何时使用哪种工具:
| 场景 | 工具 | 原因 |
|---|---|---|
| 转换切片(map、filter、reduce) | samber/lo |
有限、同步、即时——无需流开销 |
| 简单的 goroutine 扇出并处理错误 | errgroup |
标准库,轻量,足以处理有界并发 |
| 无限事件流(WebSocket、定时器、文件监控) | samber/ro |
声明式管道,支持背压、重试、超时、组合 |
| 来自多个异步源的实时数据丰富 | samber/ro |
CombineLatest/Zip 组合依赖流,无需手动 select |
| 多个消费者共享一个源的发布/订阅 | samber/ro |
热可观察对象(Share/Subjects)原生支持多播 |
lo 与 ro 的关键区别:
| 方面 | samber/lo |
samber/ro |
|---|---|---|
| 数据 | 有限切片 | 无限流 |
| 执行 | 同步、阻塞 | 异步、非阻塞 |
| 求值 | 即时(分配中间切片) | 惰性(按到达顺序处理项) |
| 时序 | 立即 | 时间感知(延迟、节流、间隔、超时) |
| 错误模型 | 每次调用返回 (T, error) |
错误通道通过管道传播 |
| 用例 | 集合转换 | 事件驱动、实时、异步管道 |
安装
go get github.com/samber/ro
核心概念
四个构建块:
- Observable — 随时间发出值的数据源。默认为冷:每个订阅者触发独立的从头执行
- Observer — 具有三个回调的消费者:
onNext(T)、onError(error)、onComplete() - Operator — 将一个 observable 转换为另一个 observable 的函数,通过
Pipe链式调用 - Subscription — observable 和 observer 之间的连接。调用
.Wait()阻塞或.Unsubscribe()取消
observable := ro.Pipe2(
ro.RangeWithInterval(0, 5, 1*time.Second),
ro.Filter(func(x int) bool { return x%2 == 0 }),
ro.Map(func(x int) string { return fmt.Sprintf("even-%d", x) }),
)
observable.Subscribe(ro.NewObserver(
func(s string) { fmt.Println(s) }, // onNext
func(err error) { log.Println(err) }, // onError
func() { fmt.Println("Done!") }, // onComplete
))
// 输出: "even-0", "even-2", "even-4", "Done!"
// 或者同步收集:
values, err := ro.Collect(observable)
冷 vs 热 Observable
冷(默认):每次 .Subscribe() 启动一个新的独立执行。安全且可预测——默认使用。
热:多个订阅者共享一个执行。当源很昂贵(WebSocket、数据库轮询)或订阅者必须看到相同事件时使用。
| 转换方式 | 行为 |
|---|---|
Share() |
冷→热,带引用计数。最后一个取消订阅时拆除 |
ShareReplay(n) |
与 Share 相同 + 为迟到订阅者缓冲最后 N 个值 |
Connectable() |
冷→热,但等待显式调用 .Connect() |
| Subjects | 原生热——直接调用 .Send()、.Error()、.Complete() |
| Subject | 构造函数 | 重放行为 |
|---|---|---|
PublishSubject |
NewPublishSubject[T]() |
无——迟到订阅者错过过去事件 |
BehaviorSubject |
NewBehaviorSubject[T](initial) |
向新订阅者重放最后一个值 |
ReplaySubject |
NewReplaySubject[T](bufferSize) |
重放最后 N 个值 |
AsyncSubject |
NewAsyncSubject[T]() |
仅在完成时发出最后一个值 |
UnicastSubject |
NewUnicastSubject[T](bufferSize) |
仅单个订阅者 |
有关 subject 详细信息和热 observable 模式,请参见 Subjects Guide。
操作符快速参考
| 类别 | 关键操作符 | 用途 |
|---|---|---|
| 创建 | Just、FromSlice、FromChannel、Range、Interval、Defer、Future |
从各种源创建 observable |
| 转换 | Map、MapErr、FlatMap、Scan、Reduce、GroupBy |
转换或累积流值 |
| 过滤 | Filter、Take、TakeLast、Skip、Distinct、Find、First、Last |
选择性发出值 |
| 组合 | Merge、Concat、Zip2–Zip6、CombineLatest2–CombineLatest5、Race |
合并多个 observable |
| 错误 | Catch、OnErrorReturn、OnErrorResumeNextWith、Retry、RetryWithConfig |
从错误中恢复 |
| 时序 | Delay、DelayEach、Timeout、ThrottleTime、SampleTime、BufferWithTime |
控制发射时序 |
| 副作用 | Tap/Do、TapOnNext、TapOnError、TapOnComplete |
观察而不改变流 |
| 终端 | Collect、ToSlice、ToChannel、ToMap |
将流消费为 Go 类型 |
使用类型化的 Pipe2、Pipe3 ... Pipe25 在操作符链中实现编译时类型安全。非类型化的 Pipe 使用 any 并失去类型检查。
有关完整操作符目录(150+ 操作符及签名),请参见 Operators Guide。
常见错误
| 错误 | 失败原因 | 修复 |
|---|---|---|
使用 ro.OnNext() 而没有错误处理 |
错误被静默丢弃——bug 在生产中隐藏 | 使用 ro.NewObserver(onNext, onError, onComplete) 并包含所有 3 个回调 |
使用非类型化的 Pipe() 而不是 Pipe2/Pipe3 |
失去编译时类型安全,错误在运行时暴露 | 使用 Pipe2、Pipe3...Pipe25 进行类型化的操作符链 |
在无限流上忘记 .Unsubscribe() |
goroutine 泄漏——observable 永远运行 | 使用 TakeUntil(signal)、上下文取消或显式 Unsubscribe() |
在冷足够时使用 Share() |
不必要的复杂性,难以推理生命周期 | 仅在多个消费者需要相同流时使用热 observable |
使用 samber/ro 进行有限切片转换 |
同步操作带来流开销(goroutine、订阅) | 使用 samber/lo——更简单、更快,专为切片构建 |
| 未传播上下文以进行取消 | 流忽略关闭信号,导致终止时资源泄漏 | 在管道中链式调用 ContextWithTimeout 或 ThrowOnContextCancel |
最佳实践
- 始终处理所有三个事件——使用
NewObserver(onNext, onError, onComplete),而不仅仅是OnNext。未处理的错误会导致静默数据丢失 - 使用
Collect()进行同步消费——当流有限且你需要[]T时,Collect阻塞直到完成并返回切片和错误 - 优先使用类型化的 Pipe 函数——
Pipe2、Pipe3...Pipe25在编译时捕获类型不匹配。保留非类型化Pipe用于动态操作符链 - 限制无限流——使用
Take(n)、TakeUntil(signal)、Timeout(d)或上下文取消。无界流会泄漏 goroutine - 使用
Tap/Do进行可观测性——在不改变流的情况下记录、追踪或计量发射。链式调用TapOnError进行错误监控 - 简单转换优先使用
samber/lo——如果数据是有限切片且你需要 Map/Filter/Reduce,使用lo。当数据随时间到达、来自多个源或需要重试/超时/背压时,使用ro
插件生态系统
40+ 插件通过领域特定操作符扩展 ro:
| 类别 | 插件 | 导入路径前缀 |
|---|---|---|
| 编码 | JSON、CSV、Base64、Gob | plugins/encoding/... |
| 网络 | HTTP、I/O、FSNotify | plugins/http、plugins/io、plugins/fsnotify |
| 调度 | Cron、ICS | plugins/cron、plugins/ics |
| 可观测性 | Zap、Slog、Zerolog、Logrus、Sentry、Oops | plugins/observability/...、plugins/samber/oops |
| 限流 | Native、Ulule | plugins/ratelimit/... |
| 数据 | Bytes、Strings、Sort、Strconv、Regexp、Template | plugins/bytes、plugins/strings 等 |
| 系统 | Process、Signal | plugins/proc、plugins/signal |
有关完整插件目录及导入路径和使用示例,请参见 Plugin Ecosystem。
有关实际响应式模式(重试+超时、WebSocket 扇出、优雅关闭、流组合),请参见 Patterns。
如果你在 samber/ro 中遇到 bug 或意外行为,请在 github.com/samber/ro/issues 提交 issue。
交叉引用
- → 参见
samber/cc-skills-golang@golang-samber-lo技能以进行有限切片转换(Map、Filter、Reduce、GroupBy)——当数据已在切片中时使用 lo - → 参见
samber/cc-skills-golang@golang-samber-mo技能以获取可与 ro 管道组合的单子类型(Option、Result、Either) - → 参见
samber/cc-skills-golang@golang-samber-hot技能以进行内存缓存(也可作为 ro 插件使用) - → 参见
samber/cc-skills-golang@golang-concurrency技能以获取 goroutine/channel 模式(当响应式流过于重量级时) - → 参见
samber/cc-skills-golang@golang-observability技能以在生产中监控响应式管道





