# WorkScheduler **Repository Path**: mengyun8/WorkScheduler ## Basic Information - **Project Name**: WorkScheduler - **Description**: No description available - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2025-07-31 - **Last Updated**: 2025-08-03 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # WorkScheduler - 分布式任务调度系统 WorkScheduler 是一个高性能的分布式任务调度系统,支持多种任务类型、自动负载均衡、任务清理和实时监控。 ## 🚀 特性 - **多种任务类型支持**:General、Compute、IO、Batch、Real-time - **分布式架构**:Registry + Server + Client 模式 - **自动负载均衡**:智能任务分配和服务器选择 - **任务清理机制**:自动清理过期任务,防止内存溢出 - **实时监控**:任务状态跟踪和性能指标 - **高可用性**:心跳检测和故障恢复 - **可扩展性**:支持自定义任务类型和处理逻辑 ## 📋 目录 - [快速开始](#快速开始) - [系统架构](#系统架构) - [任务类型详解](#任务类型详解) - [自定义任务开发](#自定义任务开发) - [API 参考](#api-参考) - [配置说明](#配置说明) - [监控和日志](#监控和日志) - [故障排除](#故障排除) ## 🏃‍♂️ 快速开始 ### 1. 编译项目 ```bash make clean && make ``` ### 2. 启动注册中心 ```bash ./bin/ws_register -addr 0.0.0.0 -port 50051 -db-host localhost -db-port 3306 -db-user root -db-password password -db-name workscheduler ``` ### 3. 启动服务器 ```bash ./bin/ws_server -id server-1 -addr 0.0.0.0 -port 50052 -registry-addr localhost -registry-port 50051 -worker-pool-size 5 -heartbeat-interval 5s -log-dir logs ``` ### 4. 运行演示程序 ```bash go run task_demo.go ``` ## 🏗️ 系统架构 ``` ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ Client │ │ Registry │ │ Server │ │ │◄──►│ │◄──►│ │ │ 提交任务 │ │ 任务调度 │ │ 执行任务 │ └─────────────┘ └─────────────┘ └─────────────┘ ``` ### 组件说明 - **Registry(注册中心)**:负责服务器注册、任务调度和负载均衡 - **Server(执行服务器)**:接收并执行任务,支持多种任务类型 - **Client(客户端)**:提交任务和查询状态 ## 📝 任务类型详解 ### 1. General Task(通用任务) 适用于简单的命令执行任务。 ```go task := &task.Task{ ID: "general-task-001", Name: "Echo Task", Content: "echo 'Hello World'", TaskType: proto.TaskType_GENERAL, TaskParams: map[string]string{ "command": "echo", "arguments": `["Hello", "World"]`, "environment": `{"PATH": "/usr/bin"}`, "working_dir": "/tmp", }, Priority: 5, Timeout: 30 * time.Second, MaxRetries: 3, State: proto.TaskState_UNEXECUTED, } ``` ### 2. Compute Task(计算任务) 适用于CPU密集型计算任务。 ```go task := &task.Task{ ID: "compute-task-001", Name: "Fibonacci Calculation", Content: "Calculate fibonacci numbers", TaskType: proto.TaskType_COMPUTE, TaskParams: map[string]string{ "iterations": "1000000", "complexity": "high", "cpu_intensity": "0.8", "memory_limit": "1073741824", // 1GB "allow_parallel": "true", }, Priority: 8, Timeout: 60 * time.Second, MaxRetries: 2, State: proto.TaskState_UNEXECUTED, } ``` ### 3. IO Task(IO任务) 适用于文件操作和网络IO任务。 ```go task := &task.Task{ ID: "io-task-001", Name: "File Processing", Content: "Process large file", TaskType: proto.TaskType_IO_BOUND, TaskParams: map[string]string{ "delay_ms": "500", "file_size": "10485760", // 10MB "io_pattern": "sequential", "buffer_size": "8192", "concurrency": "4", }, Priority: 6, Timeout: 120 * time.Second, MaxRetries: 3, State: proto.TaskState_UNEXECUTED, } ``` ### 4. Batch Task(批处理任务) 适用于批量数据处理任务。 ```go task := &task.Task{ ID: "batch-task-001", Name: "Data Processing Batch", Content: "Process data batch", TaskType: proto.TaskType_BATCH, TaskParams: map[string]string{ "batch_size": "100", "batch_timeout": "30", "retry_policy": "exponential", "max_failures": "5", "continue_on_error": "true", }, Priority: 7, Timeout: 300 * time.Second, MaxRetries: 3, State: proto.TaskState_UNEXECUTED, } ``` ### 5. Real-time Task(实时任务) 适用于对延迟敏感的任务。 ```go task := &task.Task{ ID: "realtime-task-001", Name: "Real-time Processing", Content: "Process real-time data", TaskType: proto.TaskType_REAL_TIME, TaskParams: map[string]string{ "deadline_ms": "100", "priority": "9", "jitter_tolerance": "10", "latency_target": "50", "throughput": "1000", }, Priority: 9, Timeout: 10 * time.Second, MaxRetries: 1, State: proto.TaskState_UNEXECUTED, } ``` ## 🔧 自定义任务开发 ### 1. 创建自定义任务类型 #### 步骤1:定义任务参数结构 ```go // 在 task/task.go 中添加新的参数结构 type CustomTaskParams struct { CustomField1 string `json:"custom_field1"` CustomField2 int `json:"custom_field2"` CustomField3 bool `json:"custom_field3"` } // 在 TaskParameters 结构体中添加 type TaskParameters struct { // ... 现有字段 CustomParams *CustomTaskParams `json:"custom_params,omitempty"` } ``` #### 步骤2:实现参数解析方法 ```go // 在 task/task.go 中添加 func (tm *TaskModel) parseCustomParams() *CustomTaskParams { return &CustomTaskParams{ CustomField1: tm.TaskParams["custom_field1"], CustomField2: func() int { if val, err := strconv.Atoi(tm.TaskParams["custom_field2"]); err == nil { return val } return 0 }(), CustomField3: func() bool { if val, err := strconv.ParseBool(tm.TaskParams["custom_field3"]); err == nil { return val } return false }(), } } // 在 ParseParameters 方法中添加 func (tm *TaskModel) ParseParameters() (*TaskParameters, error) { params := &TaskParameters{} switch tm.TaskType { // ... 现有case case proto.TaskType_CUSTOM: // 需要先在proto中定义 params.CustomParams = tm.parseCustomParams() default: return nil, errors.New("unsupported task type: " + tm.TaskType.String()) } return params, nil } ``` #### 步骤3:实现任务处理逻辑 ```go // 在 task/processor.go 中添加 func (tp *TaskProcessor) executeCustomTask(task *Task) error { log.Printf("Executing custom task %s", task.ID) // 解析参数 taskModel := NewTaskModel(task) params, err := taskModel.ParseParameters() if err != nil { return err } // 执行自定义逻辑 customParams := params.CustomParams log.Printf("Custom task parameters: %+v", customParams) // 模拟处理 time.Sleep(100 * time.Millisecond) // 设置结果 task.Result = fmt.Sprintf("Custom task completed with field1=%s, field2=%d", customParams.CustomField1, customParams.CustomField2) return nil } // 在 ProcessTask 方法中添加 func (tp *TaskProcessor) ProcessTask(task *Task) error { // ... 现有代码 switch taskModel.TaskType { // ... 现有case case proto.TaskType_CUSTOM: result = tp.executeCustomTask(task) } // ... 现有代码 } ``` ### 2. 完整自定义任务示例 ```go package main import ( "context" "fmt" "log" "time" "WorkScheduler/internal/client" "WorkScheduler/task" proto "WorkScheduler/proto" ) func main() { // 创建客户端 c, err := client.NewClient("localhost:50051") if err != nil { log.Fatalf("Failed to create client: %v", err) } defer c.Close() // 创建自定义任务 customTask := &task.Task{ ID: fmt.Sprintf("custom-task-%d", time.Now().Unix()), Name: "Custom Processing Task", Content: "Process custom data with specific parameters", TaskType: proto.TaskType_GENERAL, // 使用GENERAL类型作为示例 TaskParams: map[string]string{ "command": "python3", "arguments": `["custom_processor.py", "--input", "data.csv", "--output", "result.json"]`, "environment": `{"PYTHONPATH": "/app", "CUSTOM_CONFIG": "production"}`, "working_dir": "/app/scripts", "custom_field1": "important_value", "custom_field2": "42", "custom_field3": "true", }, Priority: 8, Timeout: 300 * time.Second, MaxRetries: 3, State: proto.TaskState_UNEXECUTED, } // 创建TaskModel taskModel := task.NewTaskModel(customTask) taskModel.ServerID = "custom-server" // 验证任务 if err := taskModel.ValidateParameters(); err != nil { log.Fatalf("Task validation failed: %v", err) } // 解析参数 params, err := taskModel.ParseParameters() if err != nil { log.Fatalf("Failed to parse parameters: %v", err) } log.Printf("Custom task parameters: %+v", params.GeneralParams) // 提交任务到服务器 req := &proto.SubmitTaskRequest{ TaskId: customTask.ID, TaskName: customTask.Name, TaskContent: customTask.Content, TaskType: customTask.TaskType, TaskParams: customTask.TaskParams, Priority: customTask.Priority, Timeout: &proto.Duration{Seconds: int64(customTask.Timeout.Seconds())}, MaxRetries: int32(customTask.MaxRetries), } resp, err := c.SubmitTask(context.Background(), req) if err != nil { log.Fatalf("Failed to submit task: %v", err) } log.Printf("Custom task submitted successfully: %s", resp.TaskId) log.Printf("Scheduled to server: %s", resp.ServerId) // 等待任务完成 time.Sleep(5 * time.Second) // 查询任务状态 statusReq := &proto.QueryTaskRequest{TaskId: customTask.ID} statusResp, err := c.QueryTaskStatus(context.Background(), statusReq) if err != nil { log.Printf("Failed to query task status: %v", err) return } log.Printf("Task status: %s", statusResp.State.String()) log.Printf("Task result: %s", statusResp.Result) } ``` ### 3. 高级自定义任务处理 ```go // 自定义任务处理器 type CustomTaskProcessor struct { *task.TaskProcessor customHandlers map[string]func(*task.Task) error } func NewCustomTaskProcessor() *CustomTaskProcessor { return &CustomTaskProcessor{ TaskProcessor: task.NewTaskProcessor(), customHandlers: make(map[string]func(*task.Task) error), } } // 注册自定义处理器 func (ctp *CustomTaskProcessor) RegisterHandler(taskType string, handler func(*task.Task) error) { ctp.customHandlers[taskType] = handler } // 重写ProcessTask方法 func (ctp *CustomTaskProcessor) ProcessTask(task *Task) error { // 检查是否有自定义处理器 if handler, exists := ctp.customHandlers[task.TaskType.String()]; exists { return handler(task) } // 使用默认处理器 return ctp.TaskProcessor.ProcessTask(task) } // 使用示例 func main() { processor := NewCustomTaskProcessor() // 注册自定义处理器 processor.RegisterHandler("CUSTOM_DATA_PROCESSING", func(task *task.Task) error { log.Printf("Processing custom data task: %s", task.ID) // 自定义处理逻辑 taskParams := task.TaskParams inputFile := taskParams["input_file"] outputFile := taskParams["output_file"] // 模拟数据处理 time.Sleep(2 * time.Second) task.Result = fmt.Sprintf("Processed %s -> %s successfully", inputFile, outputFile) return nil }) // 创建自定义任务 customTask := &task.Task{ ID: "custom-data-task-001", Name: "Custom Data Processing", TaskType: proto.TaskType_GENERAL, TaskParams: map[string]string{ "input_file": "/data/input.csv", "output_file": "/data/output.json", "command": "custom_processor", }, Priority: 7, Timeout: 60 * time.Second, MaxRetries: 2, } // 处理任务 err := processor.ProcessTask(customTask) if err != nil { log.Printf("Task failed: %v", err) } else { log.Printf("Task completed: %s", customTask.Result) } } ``` ## 📚 API 参考 ### TaskProcessor API ```go // 创建任务处理器 processor := task.NewTaskProcessor() // 处理任务 err := processor.ProcessTask(task) // 设置清理配置 processor.SetCleanupConfig(&task.CleanupConfig{ EnableAutoCleanup: true, CleanupInterval: 5 * time.Second, MaxCompletedTaskAge: 10 * time.Second, MaxCompletedTasks: 5, MemoryThreshold: 52428800, // 50MB EnableCleanupLogging: true, }) // 启动自动清理 processor.StartAutoCleanup() // 停止自动清理 processor.StopAutoCleanup() // 获取清理统计 stats := processor.GetCleanupStats() ``` ### TaskModel API ```go // 创建任务模型 taskModel := task.NewTaskModel(task) // 验证参数 err := taskModel.ValidateParameters() // 解析参数 params, err := taskModel.ParseParameters() // 序列化为JSON jsonData, err := taskModel.ToJSON() // 从JSON反序列化 taskModel, err := task.FromJSON(jsonData) ``` ## ⚙️ 配置说明 ### 服务器配置 ```bash ./bin/ws_server \ -id server-1 \ -addr 0.0.0.0 \ -port 50052 \ -registry-addr localhost \ -registry-port 50051 \ -worker-pool-size 5 \ -heartbeat-interval 5s \ -log-dir logs ``` ### 注册中心配置 ```bash ./bin/ws_register \ -addr 0.0.0.0 \ -port 50051 \ -db-host localhost \ -db-port 3306 \ -db-user root \ -db-password password \ -db-name workscheduler ``` ### 清理配置 ```go cleanupConfig := &task.CleanupConfig{ EnableAutoCleanup: true, // 启用自动清理 CleanupInterval: 5 * time.Second, // 清理间隔 MaxCompletedTaskAge: 10 * time.Second, // 最大完成任务年龄 MaxCompletedTasks: 5, // 最大完成任务数量 MemoryThreshold: 52428800, // 内存阈值 (50MB) EnableCleanupLogging: true, // 启用清理日志 } ``` ## 📊 监控和日志 ### 日志文件 - `logs/server_YYYYMMDD_HHMMSS.log` - 服务器日志 - `logs/register_YYYYMMDD_HHMMSS.log` - 注册中心日志 ### 监控指标 ```go // 获取任务处理器统计 stats := processor.GetCleanupStats() fmt.Printf("Active tasks: %d\n", stats["active_tasks_count"]) fmt.Printf("Completed tasks: %d\n", stats["completed_tasks_count"]) fmt.Printf("Memory usage: %d bytes\n", stats["memory_usage"]) ``` ### 性能监控 ```go // 监控任务执行时间 startTime := time.Now() err := processor.ProcessTask(task) executionTime := time.Since(startTime) log.Printf("Task %s executed in %v", task.ID, executionTime) ``` ## 🔧 故障排除 ### 常见问题 1. **服务器启动失败** ```bash # 检查端口是否被占用 netstat -tlnp | grep 50052 # 检查注册中心是否运行 ps aux | grep ws_register ``` 2. **任务提交失败** ```bash # 检查服务器连接 telnet localhost 50052 # 查看服务器日志 tail -f logs/server_*.log ``` 3. **内存使用过高** ```go // 调整清理配置 processor.SetCleanupConfig(&task.CleanupConfig{ MaxCompletedTasks: 10, // 减少最大任务数 MemoryThreshold: 26214400, // 降低内存阈值 }) ``` 4. **任务执行超时** ```go // 增加任务超时时间 task.Timeout = 300 * time.Second // 调整重试策略 task.MaxRetries = 5 ``` ### 调试技巧 1. **启用详细日志** ```go log.SetLevel(log.DebugLevel) ``` 2. **监控任务状态** ```go // 定期查询任务状态 ticker := time.NewTicker(1 * time.Second) defer ticker.Stop() for range ticker.C { status, err := client.QueryTaskStatus(ctx, taskID) if err != nil { log.Printf("Status query failed: %v", err) continue } log.Printf("Task status: %s", status.State) } ``` 3. **性能分析** ```go import "runtime" // 监控内存使用 var m runtime.MemStats runtime.ReadMemStats(&m) log.Printf("Memory usage: %d MB", m.Alloc/1024/1024) ``` ## 📄 许可证 MIT License ## 🤝 贡献 欢迎提交 Issue 和 Pull Request! ## 📞 支持 如有问题,请提交 Issue 或联系开发团队。