# Chanx **Repository Path**: vipkwds/chanx ## Basic Information - **Project Name**: Chanx - **Description**: Chanx 是一个 Go 语言增强版 Channel 封装库,提供多种 Channel 模式和数据流处理能力。简化 Go 并发编程,让通道使用更加灵活便捷。 - **Primary Language**: Go - **License**: MIT - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-05-28 - **Last Updated**: 2026-06-10 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # Chanx - Go 增强版 Channel 封装库 [![Go Version](https://img.shields.io/badge/Go-1.18%2B-blue)](https://golang.org/) [![License](https://img.shields.io/badge/License-MIT-green)](LICENSE) ## 简介 Chanx 是一个 Go 语言增强版 Channel 封装库,提供多种 Channel 模式和数据流处理能力。简化 Go 并发编程,让通道使用更加灵活便捷。 **设计目标**: - 提供阻塞/非阻塞/无限容量三种模式 - 支持订阅发布、管道操作、FanIn/FanOut 等并发模式 - 内置指标统计,便于监控 - 零外部依赖,仅使用标准库 --- ## 目录 - [快速开始](#快速开始) - [核心类型](#核心类型) - [创建通道](#创建通道) - [发送与接收](#发送与接收) - [批量操作](#批量操作) - [订阅发布](#订阅发布) - [管道操作](#管道操作) - [FanIn/FanOut](#faninfanout) - [工作池](#工作池) - [指标监控](#指标监控) - [完整示例](#完整示例) --- ## 快速开始 ```bash go get gitee.com/vipkwds/chanx ``` ```go package main import ( "fmt" "gitee.com/vipkwds/chanx" ) func main() { // 创建一个非阻塞通道,容量 100 ch := chanx.NewNonBlocking[int](100) // 发送数据 ch.Send(1) ch.Send(2) ch.Send(3) // 接收数据 for { val, ok := ch.Receive() if !ok { break } fmt.Println("收到:", val) } ch.Close() } ``` --- ## 核心类型 ### Config - 通道配置 ```go type Config struct { Size int // 缓冲区大小,0=无缓冲,<0=无限(自动扩容) Block bool // 发送时是否阻塞,false=满了就丢弃 DropOldest bool // 满了丢旧数据(true)还是丢新数据(false) Timeout time.Duration // 默认超时 } ``` ### Metrics - 指标统计 ```go type Metrics struct { sent int64 // 成功发送的消息数 recv int64 // 成功接收的消息数 drop int64 // 丢弃的消息数 } ``` --- ## 创建通道 ### 1. 阻塞模式 (NewBlocking) 发送时如果通道满,会阻塞等待直到有空间。 ```go // 容量 10 的阻塞通道 ch := chanx.NewBlocking[int](10) // 发送时会等待有空间才返回 ch.Send(1) // 立即返回 ch.Send(2) // 立即返回 // 如果通道已满,会阻塞等待 // ch.Send(3) // 等待... ``` ### 2. 非阻塞模式 (NewNonBlocking) 发送时如果通道满,立即返回 `false` 并丢弃数据。 ```go // 容量 2 的非阻塞通道 ch := chanx.NewNonBlocking[int](2) ch.Send(1) // 成功 ch.Send(2) // 成功 ch.Send(3) // 失败返回 false,数据被丢弃 ``` ### 3. 无限容量模式 (NewInfinite) 数据存储在环形缓冲区中,不会因容量问题丢失数据。 ```go // 无限容量通道 ch := chanx.NewInfinite[string]() // 发送不会阻塞 for i := 0; i < 1000000; i++ { ch.Send(fmt.Sprintf("message-%d", i)) } // 数据通过后台循环转移到 channel for msg := range ch.ReceiveAll() { fmt.Println(msg) } ``` ### 4. 自定义配置 (New) 通过 `Config` 精细控制通道行为。 ```go // 使用 Config 创建 ch := chanx.New[int](chanx.Config{ Size: 100, // 缓冲区大小 Block: false, // 非阻塞 DropOldest: true, // 满了丢旧数据 }) // 无限容量 ch := chanx.New[string](chanx.Config{Size: -1}) ``` --- ## 发送与接收 ### 发送方法 | 方法 | 行为 | 返回值 | |------|------|--------| | `Send(val)` | 根据配置发送,阻塞或丢弃 | `true`=成功,`false`=失败 | | `SendTimeout(val, timeout)` | 带超时的发送 | `true`=成功,`false`=超时 | | `TrySend(val)` | 非阻塞立即发送 | `true`=成功,`false`=满或关闭 | ```go ch := chanx.NewNonBlocking[int](2) // Send - 阻塞模式等待,非阻塞模式丢弃 ch.Send(1) // SendTimeout - 1秒超时 if !ch.SendTimeout(2, time.Second) { fmt.Println("发送超时") } // TrySend - 立即返回 if ch.TrySend(3) { fmt.Println("发送成功") } ``` ### 接收方法 | 方法 | 行为 | 返回值 | |------|------|--------| | `Receive()` | 阻塞接收 | `(val, ok)` | | `ReceiveTimeout(timeout)` | 带超时的接收 | `(val, ok)` | | `TryReceive()` | 非阻塞立即接收 | `(val, ok)` | ```go // Receive - 阻塞接收 val, ok := ch.Receive() if ok { fmt.Println("收到:", val) } // ReceiveTimeout - 1秒超时 val, ok := ch.ReceiveTimeout(time.Second) if ok { fmt.Println("收到:", val) } else { fmt.Println("超时或已关闭") } // TryReceive - 立即返回 val, ok := ch.TryReceive() if ok { fmt.Println("收到:", val) } ``` ### 生命周期 ```go ch := chanx.NewBlocking[int](10) // 检查是否关闭 if !ch.IsClosed() { ch.Send(1) } // 关闭通道 ch.Close() // 关闭后发送返回 false if ch.Send(1) { fmt.Println("发送成功") } else { fmt.Println("发送失败 - 通道已关闭") } ``` --- ## 批量操作 ### BatchSend - 批量发送 ```go ch := chanx.NewNonBlocking[int](5) vals := []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10} sent := ch.BatchSend(vals) fmt.Printf("成功发送 %d 个\n", sent) // 最多发送 5 个 ``` ### BatchReceive - 批量接收 ```go ch := chanx.NewNonBlocking[int](100) // 发送 10 个数据 for i := 0; i < 10; i++ { ch.Send(i) } // 最多接收 5 个 vals := ch.BatchReceive(5) fmt.Printf("接收了 %d 个: %v\n", len(vals), vals) ``` ### ReceiveAll - 接收所有可用数据 ```go ch := chanx.NewNonBlocking[int](100) for i := 0; i < 10; i++ { ch.Send(i) } // 接收所有数据 vals := ch.ReceiveAll() fmt.Printf("接收了 %d 个: %v\n", len(vals), vals) ``` --- ## 订阅发布 一个通道可以有多人订阅,每个订阅者都会收到消息。 ### 基本订阅 ```go ch := chanx.NewNonBlocking[string](100) // 订阅函数 unsub := ch.Subscribe(func(msg string) { fmt.Println("订阅者A收到:", msg) }) // 发送消息 ch.Send("hello") // 订阅者收到: hello // 取消订阅 unsub() ch.Send("world") // 不再收到 ``` ### 多订阅者 ```go ch := chanx.NewNonBlocking[string](100) count1, count2 := 0, 0 // 订阅者1 ch.Subscribe(func(msg string) { count1++ }) // 订阅者2 ch.Subscribe(func(msg string) { count2++ }) // 发送 3 条消息 ch.Send("msg1") ch.Send("msg2") ch.Send("msg3") // 两个订阅者都会收到 fmt.Printf("订阅者1收到 %d 条\n", count1) // 3 fmt.Printf("订阅者2收到 %d 条\n", count2) // 3 ``` ### 取消订阅 ```go ch := chanx.NewNonBlocking[string](100) unsub1 := ch.Subscribe(func(msg string) { fmt.Println("订阅者1:", msg) }) unsub2 := ch.Subscribe(func(msg string) { fmt.Println("订阅者2:", msg) }) ch.Send("消息1") // 两者都收到 unsub2() // 取消订阅者2 ch.Send("消息2") // 只有订阅者1收到 ``` --- ## 管道操作 ### Pipe - 通道连接 将一个通道的数据转发到另一个通道。 ```go source := chanx.New[int](chanx.Config{Size: 10, Block: false}) dest := chanx.New[string](chanx.Config{Size: 10, Block: false}) source.Pipe(dest) source.Send(1) source.Send(2) source.Close() // dest 会收到转换后的数据 for { val, ok := dest.Receive() if !ok { break } fmt.Println(val) } ``` ### Filter - 数据过滤 只保留满足条件的数据。 ```go ch := chanx.New[int](chanx.Config{Size: 10, Block: false}) // 只保留偶数 evenCh := ch.Filter(func(n int) bool { return n%2 == 0 }) ch.Send(1) ch.Send(2) ch.Send(3) ch.Send(4) ch.Close() // evenCh 只有 2, 4 for { val, ok := evenCh.Receive() if !ok { break } fmt.Println(val) // 输出 2, 4 } ``` ### Map - 数据类型转换 将一种类型的通道转换成另一种类型。 ```go numCh := chanx.New[int](chanx.Config{Size: 10, Block: false}) // 转换为字符串通道 strCh := chanx.Map(numCh, func(n int) string { return strconv.Itoa(n) }) numCh.Send(100) numCh.Send(200) numCh.Close() for { val, ok := strCh.Receive() if !ok { break } fmt.Println(val) // 输出 "100", "200" } ``` ### 链式调用 ```go source := chanx.New[int](chanx.Config{Size: 100, Block: false}) // 过滤偶数 -> 乘以 10 -> 转换为字符串 result := source. Filter(func(n int) bool { return n%2 == 0 }). Map(func(n int) int { return n * 10 }). Map(func(n int) string { return strconv.Itoa(n) }) for i := 1; i <= 10; i++ { source.Send(i) } source.Close() // 输出: "20", "40", "60", "80", "100" for { val, ok := result.Receive() if !ok { break } fmt.Println(val) } ``` --- ## FanIn/FanOut ### FanIn - 多合一 将多个通道合并成一个。 ```go ch1 := chanx.New[int](chanx.Config{Size: 10, Block: false}) ch2 := chanx.New[int](chanx.Config{Size: 10, Block: false}) // 合并两个通道 merged := chanx.FanIn(ch1, ch2) ch1.Send(1) ch2.Send(2) ch1.Close() ch2.Close() // 从合并通道接收 for { val, ok := merged.Receive() if !ok { break } fmt.Println(val) // 输出 1, 2(顺序不保证) } ``` ### FanOut - 一对多广播 将一个通道的数据广播到多个输出通道。 ```go ch := chanx.New[int](chanx.Config{Size: 10, Block: false}) // 创建 3 个输出通道 outputs := ch.FanOut(3) ch.Send(100) ch.Send(200) ch.Close() // 每个输出通道都能收到所有数据 for i, out := range outputs { fmt.Printf("输出通道%d:\n", i) for { val, ok := out.Receive() if !ok { break } fmt.Printf(" 收到 %d\n", val) } } // 输出: // 输出通道0: // 收到 100 // 收到 200 // 输出通道1: // 收到 100 // 收到 200 // 输出通道2: // 收到 100 // 收到 200 ``` --- ## 工作池 ### 基本使用 ```go // 创建工作池,4个worker,任务类型int,结果类型string pool := chanx.NewWorkerPool(4, func(job int) string { time.Sleep(100 * time.Millisecond) return fmt.Sprintf("处理结果-%d", job) }) pool.Start() // 提交任务 for i := 1; i <= 10; i++ { pool.Submit(i) } // 收集结果 for result := range pool.Results() { fmt.Println(result) } pool.Stop() ``` ### 完整示例 ```go package main import ( "fmt" "time" "gitee.com/vipkwds/chanx" ) func main() { // 模拟图片处理 pool := chanx.NewWorkerPool(3, func(path string) string { // 模拟处理时间 time.Sleep(200 * time.Millisecond) return fmt.Sprintf("processed-%s", path) }) pool.Start() // 提交图片处理任务 images := []string{"img1.jpg", "img2.jpg", "img3.jpg", "img4.jpg", "img5.jpg"} for _, img := range images { pool.Submit(img) } // 收集处理结果 done := make(chan struct{}) go func() { for result := range pool.Results() { fmt.Printf("完成: %s\n", result) } done <- struct{}{} }() // 等待完成 <-done pool.Stop() } ``` ### 队列已满的场景 ```go pool := chanx.NewWorkerPool(2, func(job int) int { return job * 2 }) pool.Start() // 队列满后 Submit 返回 false for i := 0; i < 150; i++ { if !pool.Submit(i) { fmt.Printf("任务 %d 提交失败,队列已满\n", i) } } ``` --- ## 指标监控 ### 获取统计数据 ```go ch := chanx.NewNonBlocking[int](3) // 发送数据 ch.Send(1) // 成功 ch.Send(2) // 成功 ch.Send(3) // 成功 ch.Send(4) // 失败,丢弃 ch.Send(5) // 失败,丢弃 // 接收数据 ch.Receive() ch.Receive() ch.Receive() ch.Receive() // 获取统计 sent, recv, drop := ch.Metrics() fmt.Printf("发送: %d, 接收: %d, 丢弃: %d\n", sent, recv, drop) // 输出: 发送: 3, 接收: 3, 丢弃: 2 ``` ### 监控使用场景 ```go ch := chanx.NewNonBlocking[string](1000) go func() { for { sent, recv, drop := ch.Metrics() fmt.Printf("Metrics{sent:%d, recv:%d, drop:%d}\n", sent, recv, drop) time.Sleep(5 * time.Second) } }() // ... 生产代码 ... ``` --- ## 完整示例 ### 事件驱动架构 ```go // 订单通知系统 orderCh := chanx.NewNonBlocking[Order](1000) // 订阅者:物流系统 orderCh.Subscribe(func(o Order) { sendToLogistics(o) }) // 订阅者:短信通知 orderCh.Subscribe(func(o Order) { sendSMSNotification(o) }) // 订阅者:更新库存 orderCh.Subscribe(func(o Order) { updateInventory(o) }) // 处理订单 for order := range orders { orderCh.Send(order) } ``` ### 流量削峰 ```go // 接收请求 requestCh := chanx.NewInfinite[Request]() // 后台处理 go func() { for req := range requestCh.ReceiveAll() { processRequest(req) } }() // HTTP 处理函数 func handleRequest(w http.ResponseWriter, r *http.Request) { requestCh.Send(parseRequest(r)) w.WriteHeader(http.StatusAccepted) } ``` ### 数据清洗管道 ```go // 日志处理管道 errorLogs := logCh. Filter(func(l LogEntry) bool { return l.Level == "ERROR" }). Map(func(l LogEntry) string { return l.Timestamp.Format("2006-01-02") + " " + l.Message }). Map(func(s string) string { return strings.ToUpper(s) }) for log := range errorLogs.ReceiveAll() { saveToES(log) } ``` ### 并发任务处理 ```go // 图片批量处理 pool := chanx.NewWorkerPool(10, func(img Image) Result { return resizeAndCompress(img) }) pool.Start() for _, img := range images { pool.Submit(img) } for range images { result := <-pool.Results() save(result) } pool.Stop() ``` --- ## 性能注意事项 ### 1. 无限模式的 ticker `NewInfinite` 使用后台 ticker 定期检查队列,在数据量极大时可能有轻微延迟。建议在延迟敏感场景使用带缓冲的 channel。 ### 2. 非阻塞模式的数据丢弃 非阻塞模式满时自动丢弃数据,请根据业务场景选择合适的模式: - **阻塞模式** (`NewBlocking`): 不能丢数据 - **非阻塞模式** (`NewNonBlocking`): 可以丢数据,关心最新数据 ### 3. 订阅发布的 goroutine 每个订阅者启动一个独立 goroutine,订阅者执行时间会影响后续消息的传递时机。 --- ## API 索引 ### 创建函数 | 函数 | 说明 | |------|------| | `New[T](Config)` | 自定义配置创建 | | `NewBlocking[T](size)` | 创建阻塞通道 | | `NewNonBlocking[T](size)` | 创建非阻塞通道 | | `NewInfinite[T]()` | 创建无限容量通道 | | `FanIn[T](...*Chan[T])` | 合并多个通道 | | `NewWorkerPool[T,R](workers, fn)` | 创建工作池 | ### Chan 方法 | 方法 | 说明 | |------|------| | `Send(T) bool` | 发送数据 | | `SendTimeout(T, Duration) bool` | 带超时发送 | | `TrySend(T) bool` | 非阻塞发送 | | `Receive() (T, bool)` | 接收数据 | | `ReceiveTimeout(Duration) (T, bool)` | 带超时接收 | | `TryReceive() (T, bool)` | 非阻塞接收 | | `BatchSend([]T) int` | 批量发送 | | `BatchReceive(int) []T` | 批量接收 | | `ReceiveAll() []T` | 接收所有数据 | | `Subscribe(func(T)) func()` | 订阅消息 | | `Pipe(*Chan[T]) *Chan[T]` | 管道连接 | | `Filter(func(T) bool) *Chan[T]` | 数据过滤 | | `FanOut(int) []*Chan[T]` | 一对多广播 | | `Close()` | 关闭通道 | | `IsClosed() bool` | 检查是否关闭 | | `Len() int` | 当前长度 | | `Cap() int` | 容量 | | `Metrics() (sent, recv, drop int64)` | 获取指标 | ### WorkerPool 方法 | 方法 | 说明 | |------|------| | `Start()` | 启动 worker | | `Submit(T) bool` | 提交任务 | | `Results() <-chan R` | 获取结果 channel | | `Stop()` | 停止工作池 | --- ## License MIT License - 详见 LICENSE 文件