# jobflux **Repository Path**: vipkwds/jobflux ## Basic Information - **Project Name**: jobflux - **Description**: Package jobflux 提供通用的任务流管道管理器 设计理念:每个业务管道是一个自治单元,包含队列、协程池、限流、重试等完整能力 JobHub 仅负责任务流路由和生命周期管理 - **Primary Language**: Go - **License**: MIT - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-04-29 - **Last Updated**: 2026-05-22 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # jobflux [![Go Version](https://img.shields.io/badge/Go-1.18+-blue.svg)](https://golang.org/) [![License](https://img.shields.io/badge/License-MIT-green.svg)](LICENSE) **jobflux** 是一个轻量级的 Go 任务流管道管理器。 每个业务管道是一个自治单元,包含队列、协程池、限流、重试等完整能力。JobHub 仅负责任务流路由和生命周期管理。 ## 特性 - 🚀 **轻量级** - 零外部依赖,仅使用 Go 标准库 - 🎯 **高内聚** - 每个 Pipe 自治,包含完整的任务处理能力 - 🔄 **自动重试** - 支持指数退避重试策略 - ⏱️ **超时控制** - 每个任务可独立设置超时时间 - 🚦 **限流保护** - 支持每秒请求数限流 - 📊 **内置统计** - 开箱即用的监控指标 - 🛑 **优雅关闭** - 支持 context 超时和队列排空 - 🔌 **接口驱动** - 只需实现 `JobHandler` 接口即可接入业务 ## 安装 ```bash go get gitee.com/vipkwds/jobflux ``` ## 快速开始 ### 1. 实现业务处理器 ```go package main import ( "context" "fmt" "time" "gitee.com/vipkwds/jobflux" ) // ImageHandler 图片处理业务 type ImageHandler struct{} func (h *ImageHandler) JobType() string { return "image" } func (h *ImageHandler) Handle(ctx context.Context, payload interface{}) error { data := payload.(map[string]interface{}) fmt.Printf("处理图片: %s\n", data["path"]) time.Sleep(100 * time.Millisecond) return nil } // EmailHandler 邮件发送业务 type EmailHandler struct{} func (h *EmailHandler) JobType() string { return "email" } func (h *EmailHandler) Handle(ctx context.Context, payload interface{}) error { email := payload.(map[string]string) fmt.Printf("发送邮件: to=%s, subject=%s\n", email["to"], email["subject"]) return nil } ``` ### 2. 创建管道并注册到中心 ```go func main() { // 创建管道中心 hub := jobflux.NewJobHub() // 创建图片处理管道(使用默认配置) imagePipe := jobflux.NewJobPipe(&ImageHandler{}, nil) hub.Register(imagePipe) // 创建邮件发送管道(自定义配置) emailConfig := &jobflux.JobPipeConfig{ Name: "email_sender", MaxWorkers: 10, // 最大并发数 QueueSize: 500, // 队列缓冲大小 MaxRetry: 3, // 最大重试次数 RetryDelay: time.Second, // 重试间隔 RateLimit: 100, // 每秒限流 100 个 EnableStats: true, // 启用统计 } emailPipe := jobflux.NewJobPipe(&EmailHandler{}, emailConfig) hub.Register(emailPipe) // 启动所有管道 ctx := context.Background() if err := hub.StartAll(ctx); err != nil { panic(err) } defer hub.StopAll(ctx) // 提交任务 hub.Submit(&jobflux.JobFlow{ ID: "flow_001", Type: "image", Payload: map[string]interface{}{ "path": "/uploads/photo.jpg", }, }) hub.Submit(&jobflux.JobFlow{ ID: "flow_002", Type: "email", Payload: map[string]string{ "to": "user@example.com", "subject": "欢迎注册", "body": "感谢您的注册", }, }) // 等待任务完成 time.Sleep(2 * time.Second) } ``` ## 配置说明 ### JobPipeConfig | 字段 | 类型 | 默认值 | 说明 | |-----|------|--------|------| | Name | string | - | 管道名称 | | MaxWorkers | int | 5 | 最大并发工作协程数 | | QueueSize | int | 1000 | 队列缓冲大小 | | MaxRetry | int | 3 | 最大重试次数 | | RetryDelay | time.Duration | 1s | 重试间隔基数(指数退避) | | RateLimit | int | 0 | 每秒限流数(0=不限流) | | EnableStats | bool | true | 是否启用统计 | ## API 文档 ### JobHandler 接口 ```go type JobHandler interface { // JobType 返回该处理器支持的任务流类型 JobType() string // Handle 执行具体的业务逻辑 Handle(ctx context.Context, payload interface{}) error } ``` ### JobPipe 接口 ```go type JobPipe interface { Name() string JobType() string Start(ctx context.Context) error Stop(ctx context.Context) error Submit(flow *JobFlow) error Stats() JobPipeStats IsRunning() bool } ``` ### JobHub 方法 | 方法 | 说明 | |-----|------| | `NewJobHub()` | 创建管道中心 | | `Register(pipe JobPipe) error` | 注册管道 | | `Unregister(jobType string) error` | 注销管道 | | `Submit(flow *JobFlow) error` | 提交任务流 | | `SubmitBatch(flows []*JobFlow) []error` | 批量提交 | | `StartAll(ctx context.Context) error` | 启动所有管道 | | `StopAll(ctx context.Context) error` | 停止所有管道 | | `GetPipe(jobType string) (JobPipe, bool)` | 获取指定管道 | | `GetAllStats() map[string]JobPipeStats` | 获取所有统计 | ## 运行测试 ```bash go test -v ``` 运行基准测试: ```bash go test -bench=. -benchmem ``` ## 适用场景 ✅ **适合:** - 单 Go 进程内的多业务任务隔离 - 需要队列、并发控制、限流、重试,但不想引入 Redis/消息队列 - 比原生 channel 更丰富,比 machinery 更简单的场景 ❌ **不适合:** - 分布式任务调度 - 需要持久化队列 - 跨进程/跨服务任务分发 ## 设计原则 1. **内聚**:每个 Pipe 自治,包含队列、工作池、限流、重试等完整能力 2. **简单**:零依赖,只使用标准库 3. **可观测**:内置统计收集器 4. **安全**:支持超时控制、优雅关闭