golang-samber-ro

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` 技能)。

2261Star
150Fork
更新于 2026/6/6
SKILL.md
只读
名称
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` 技能)。

角色: 你是一名 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

核心概念

四个构建块:

  1. Observable — 随时间发出值的数据源。默认为冷:每个订阅者触发独立的从头执行
  2. Observer — 具有三个回调的消费者:onNext(T)onError(error)onComplete()
  3. Operator — 将一个 observable 转换为另一个 observable 的函数,通过 Pipe 链式调用
  4. 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

操作符快速参考

类别 关键操作符 用途
创建 JustFromSliceFromChannelRangeIntervalDeferFuture 从各种源创建 observable
转换 MapMapErrFlatMapScanReduceGroupBy 转换或累积流值
过滤 FilterTakeTakeLastSkipDistinctFindFirstLast 选择性发出值
组合 MergeConcatZip2Zip6CombineLatest2CombineLatest5Race 合并多个 observable
错误 CatchOnErrorReturnOnErrorResumeNextWithRetryRetryWithConfig 从错误中恢复
时序 DelayDelayEachTimeoutThrottleTimeSampleTimeBufferWithTime 控制发射时序
副作用 Tap/DoTapOnNextTapOnErrorTapOnComplete 观察而不改变流
终端 CollectToSliceToChannelToMap 将流消费为 Go 类型

使用类型化的 Pipe2Pipe3 ... Pipe25 在操作符链中实现编译时类型安全。非类型化的 Pipe 使用 any 并失去类型检查。

有关完整操作符目录(150+ 操作符及签名),请参见 Operators Guide

常见错误

错误 失败原因 修复
使用 ro.OnNext() 而没有错误处理 错误被静默丢弃——bug 在生产中隐藏 使用 ro.NewObserver(onNext, onError, onComplete) 并包含所有 3 个回调
使用非类型化的 Pipe() 而不是 Pipe2/Pipe3 失去编译时类型安全,错误在运行时暴露 使用 Pipe2Pipe3...Pipe25 进行类型化的操作符链
在无限流上忘记 .Unsubscribe() goroutine 泄漏——observable 永远运行 使用 TakeUntil(signal)、上下文取消或显式 Unsubscribe()
在冷足够时使用 Share() 不必要的复杂性,难以推理生命周期 仅在多个消费者需要相同流时使用热 observable
使用 samber/ro 进行有限切片转换 同步操作带来流开销(goroutine、订阅) 使用 samber/lo——更简单、更快,专为切片构建
未传播上下文以进行取消 流忽略关闭信号,导致终止时资源泄漏 在管道中链式调用 ContextWithTimeoutThrowOnContextCancel

最佳实践

  1. 始终处理所有三个事件——使用 NewObserver(onNext, onError, onComplete),而不仅仅是 OnNext。未处理的错误会导致静默数据丢失
  2. 使用 Collect() 进行同步消费——当流有限且你需要 []T 时,Collect 阻塞直到完成并返回切片和错误
  3. 优先使用类型化的 Pipe 函数——Pipe2Pipe3...Pipe25 在编译时捕获类型不匹配。保留非类型化 Pipe 用于动态操作符链
  4. 限制无限流——使用 Take(n)TakeUntil(signal)Timeout(d) 或上下文取消。无界流会泄漏 goroutine
  5. 使用 Tap/Do 进行可观测性——在不改变流的情况下记录、追踪或计量发射。链式调用 TapOnError 进行错误监控
  6. 简单转换优先使用 samber/lo——如果数据是有限切片且你需要 Map/Filter/Reduce,使用 lo。当数据随时间到达、来自多个源或需要重试/超时/背压时,使用 ro

插件生态系统

40+ 插件通过领域特定操作符扩展 ro:

类别 插件 导入路径前缀
编码 JSON、CSV、Base64、Gob plugins/encoding/...
网络 HTTP、I/O、FSNotify plugins/httpplugins/ioplugins/fsnotify
调度 Cron、ICS plugins/cronplugins/ics
可观测性 Zap、Slog、Zerolog、Logrus、Sentry、Oops plugins/observability/...plugins/samber/oops
限流 Native、Ulule plugins/ratelimit/...
数据 Bytes、Strings、Sort、Strconv、Regexp、Template plugins/bytesplugins/strings
系统 Process、Signal plugins/procplugins/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 技能以在生产中监控响应式管道