refactor: 调度服务端改为控制台架构,借助现有H ost架构和分布式缓存架构
大石头 authored at 2023-06-10 10:59:44
11.74 KiB
AntJob
# 架构设计 — AntJob 蚂蚁调度 > 版本:v4.4 | 日期:2026-07-20 | 从代码逆向整理 本文档描述 AntJob 的分层架构、核心组件、关键流程和设计决策。需求定义见[需求文档](/NewLife/AntJob/Blob/master/Doc/需求文档.md),实现状态见[功能清单](/NewLife/AntJob/Blob/master/Doc/功能清单.md)。 --- ## 1. 项目分层 AntJob 遵循 [NewLife 架构分层](/NewLife/AntJob/Blob/master/Doc/../.github/instructions/development.instructions.md) 的务实渐进哲学: - **阶段 2 — 数据层独立**:AntJob.Data 作为独立类库,被 Server/Web/Agent 三个入口共享 - **按需引入服务层**:Server 端 AntService 和 Web 端共享 AppService/JobService,通过 `<Compile Include Link>` 避免重复 - **混合分层**:Agent 端无独立服务层,Handler 直接对接 Provider 通信 ``` ┌──────────────────────────────────────────────┐ │ AntJob.Web (net10.0) │ ← 表现层 │ Cube MVC + REST API + Startup │ │ Link: AppService, JobService, Setting │ ├──────────────────────────────────────────────┤ │ AntJob.Server (net10.0) │ ← 服务层+表现层 │ ApiServer(TCP) + AppService + JobService │ │ AntJob.Agent (net10.0) │ ← 客户端 ├──────────────────────────────────────────────┤ │ AntJob.Extensions (netstandard2.x) │ ← 扩展层 │ DataHandler / SqlHandler / SqlMessage │ ├──────────────────────────────────────────────┤ │ AntJob.Data (netstandard2.1) │ ← 独立数据层 │ XCode Entity × 6: Model.xml → App/Job/... │ ├──────────────────────────────────────────────┤ │ AntJob (netstandard2.1/2.0/net461/net45) │ ← 核心SDK │ Scheduler / Handler / IJobProvider / Models │ │ Providers: Network / File / Http │ └──────────────────────────────────────────────┘ ``` ### 依赖关系 ```mermaid graph TD subgraph NuGet NC[NewLife.Core] NX[NewLife.XCode] NR[NewLife.Remoting] NS[NewLife.Stardust] NCC[NewLife.Cube.Core] end AJ[AntJob 核心SDK] --> NC AJ --> NR AJ --> NS AJD[AntJob.Data] --> NX AJD --> AJ AJE[AntJob.Extensions] --> NX AJE --> AJ AJS[AntJob.Server] --> AJD AJS --> MySQL[NewLife.MySql] AJS --> Redis[NewLife.Redis] AJW[AntJob.Web] --> AJ AJW --> AJD AJW --> NCC AJW -.->|Link 源文件| AJS AJA[AntJob.Agent] --> AJ AJA --> AJE ``` > **Link 源文件**:AntJob.Web 通过 `<Compile Include="..\AntJob.Server\Services\AppService.cs" Link="Services\AppService.cs" />` 共享 AntJob.Server 的服务逻辑,避免代码重复。NuGet 打包时 AntJob.Web 内嵌 `AntJob.Data.dll`,用户只需引用一个包。 --- ## 2. 核心组件 ### 2.1 SYS — 核心调度引擎 | 组件 | 类型 | 职责 | |------|------|------| | `Scheduler` | class | 调度引擎:管理 `List<Handler>`,从 DI 发现处理器,TimerX 定时连接服务端 | | `Handler` | abstract class | 处理器基类:生命周期 Init→Start→Acquire→Process→OnProcess→Execute→OnFinish,支持同步/异步双模 | | `JobContext` | class | 任务上下文:Handler/Task/Result/Data/Error/Cost/Speed/Stopwatch | | `IJobProvider` | interface | 作业提供者接口:Acquire/Produce/Report/Finish/GetJobs/SetJob | | `JobProvider` | abstract class | 提供者基类:DisposeBase + IJobProvider + ITracerFeature + ILogFeature | | `NetworkJobProvider` | class | 网络提供者:AntClient RPC → 调度中心,Peers 邻居发现 | | `FileJobProvider` | class | 文件提供者:JobFile XML 持久化,离线单机调试 | | `HttpJobProvider` | class | HTTP 提供者:ApiHttpClient 接入(编译排除,A/B 测试) | | `AntClient` | class | RPC 客户端:ClientBase 子类,Login/GetJobs/AddJobs/Acquire/Report/Finish | | `AntSetting` | Config | 客户端配置:Server(多地址主备)/AppID/Secret/Debug | | `AntJobExtensions` | static class | DI 扩展:AddAntJob() → 注册 Scheduler/AntJobWorker/AntSetting | | `TemplateHelper` | static class | 模板引擎:{dt}/{End}/{Message} 变量替换,Pool.StringBuilder 优化 | | `TimeExpression` | class | 时间表达式:{dt+1M+5d:yyyyMMdd} 解析执行,y/M/d/H/m/s/w | | `MessageHandler` | abstract class | 消息处理器:Topic 订阅,JSON 解码,逐条 ProcessItem | | `CSharpHandler` | class | C# 脚本处理器:Mode=Time,Execute 读 code 执行(主体待实现) | | `MessageOption` | class | 消息选项:DelayTime/Unique/AppId | ### 2.2 DATA — 数据持久化 | 实体 | 主键 | 关键字段 | |------|------|----------| | `App` | Int32 自增 | Name(唯一)/Secret/Enable/Version/JobCount/MessageCount/ManagerId(Map→User) | | `AppOnline` | Int32 自增 | Instance(唯一)/AppID(Map→App)/Client/统计(Tasks/Total/Success/Error/Cost/Speed) | | `Job` | Int32 自增 | AppID+Name(联合唯一)/Mode/Cron/Step/Offset/MaxTask/控制参数/QuietTime/Data | | `JobTask` | Int32 自增 | JobID/DataTime(DataScale=time)/Status/Client/统计/Key/TraceId | | `JobError` | Int32 自增 | AppID/JobID/TaskID/DataTime/Data/Server/TraceId | | `AppHistory` | Int64 雪花Id | AppID/Action/Success/TraceId(DataScale=time 分表) | | `AppMessage` | Int64 雪花Id | AppID/JobID/Topic/Data/DelayTime(DataScale=time 分表) | ### 2.3 EXT — 扩展调度 | 组件 | 类型 | 职责 | |------|------|------| | `DataHandler` | abstract class | 数据窗口基类:Factory/Field/Where/OrderBy/Selects/KeepFirstPage,Mode=Data | | `DataHandler<TEntity>` | abstract class | 泛型数据处理器:构造器自动识别 Field(雪花Id→MasterTime→UpdateTime→CreateTime) | | `SqlHandler` | class | SQL 执行器:TemplateHelper→SqlSection.ParseAll→事务内 Query/Execute/Insert | | `SqlMessage` | class | SQL 消息处理器:继承 MessageHandler,Topic="Sql",topic_ 列自动生产消息 | | `SqlSection` | class | SQL 片段解析:/*use connName*/ 指定连接,双换行分隔,Query/Execute/Insert | ### 2.4 SRV — 调度中心服务 | 组件 | 类型 | 职责 | |------|------|------| | `AntService` | class | RPC API:[Api(null)] + IActionFilter,Session 绑定 App/AppOnline,Login/GetJobs/AddJobs/SetJob/Acquire/Report/Finish | | `AppService` | class | 应用认证:Login(自动注册+SaltPasswordProvider md5)/Logout/Ping/GetPeers/WriteHistory | | `JobService` | class | 核心调度:GetJobs/AddJobs/SetJob/Acquire(延迟重试→错误重试→时间切片→消息消费)/全局Redis锁 | | `Worker` | IHostedService | 后台服务:ApiServer(9999)→Register<AntService>→星尘注册→ClearOnline(10s)/ClearItems(1h) | | `AntJobSetting` | Config | 服务配置:[Config("AntJob")],Port/TokenSecret/AutoRegistry/SessionTimeout | ### 2.5 WEB — 可视化管理后台 | 组件 | 类型 | 职责 | |------|------|------| | `AntJobController` | ControllerBase | HTTP REST API:[ApiController]+[Route],JWT鉴权,/AntJob/Login、/AntJob/GetJobs 等 | | `AntArea` | AreaBase | Cube 区域:DisplayName="蚂蚁调度",RegisterArea<AntArea> | | `AntEntityController<T>` | EntityController<T> | 基类:OnActionExecuting 识别 appId 导航,OnGetFields 条件移除 AppName | | `AppController` | AntEntityController<App> | 应用管理:列表含在线/作业/任务/消息/错误/历史链接列 | | `JobController` | AntEntityController<Job> | 作业管理:Name列链接到JobTask,模式彩色,MyTextField自定义下一时间/Cron | | `JobTaskController` | AntEntityController<JobTask> | 任务管理:MyTitleField按模式彩色渲染时间,Status颜色类,TraceUrl | | `Startup` | class | 启动:AddStardust→AddCube→EntityFactory预热4连接→RegisterService | --- ## 3. 关键流程 ### 3.1 任务调度生命周期 ```mermaid sequenceDiagram participant Agent as 执行节点 participant Provider as NetworkJobProvider participant Server as 调度中心 participant DB as 数据库 Agent->>Provider: Start() Provider->>Server: AntClient.Login() Server->>DB: App/AppOnline 写入 Server-->>Provider: LoginResponse loop 调度循环 Provider->>Server: GetJobs() Server-->>Provider: Job[] Provider->>Agent: Handler.Init() → Start() Agent->>Provider: Acquire(count) Provider->>Server: Acquire(job, count) Server->>DB: 延迟→错误→切片→消息 Server-->>Provider: Task[] Agent->>Agent: Process/ProcessAsync Note over Agent: OnProcess→Execute→OnFinish Agent->>Provider: Report+Finish Provider->>Server: 上报 Server->>DB: JobTask 状态更新 end ``` ### 3.2 任务切片策略 ``` JobService.Acquire(app, model, online): 1. 检查应用启用 && 作业启用 && 非免打扰时段 2. 获取全局 Redis 锁 antjob:lock:{jobId}(15s 超时) 3. CheckDelayTask → 优先分配延迟任务 4. CheckOldTask → 其次分配错误/超时重试 5. 不足时按 Mode 生成新切片: - Time: Cron.GetNext() 或 DataTime+Step - Data: DataTime+Step,窗口 [DataTime, DataTime+Step) - Message: 从 AppMessage 拉取待消费消息 ``` ### 3.3 DataHandler 数据抽取 ``` Init: 自动识别 Field: DataScale含"time"/"timeShard" → 雪花Id主键 → MasterTime → UpdateTime → CreateTime Job.DataTime < 2000年 → GetMinDataTime()(按Field升序取首条记录时间) OnProcess(ctx): while true: Fetch(ctx, ref row): WHERE Field >= DataTime AND Field < End ORDER BY Field ASC, 分页 BatchSize KeepFirstPage=true 时 row 始终=0(每次只查第一页) if 无数据: break ctx.Data = list → Execute(ctx) → Report(ctx, 处理中) ``` --- ## 4. 关键设计决策 | 决策 | 说明 | |------|------| | UTC 通信 | 客户端↔服务端作业时间字段使用 UTC 传输,本地存储 LocalTime | | 全局 Redis 锁 | `antjob:lock:{jobId}` 防止多服务端并发分配同一作业 | | 双模处理器 | 自动检测是否重写 ExecuteAsync/ProcessAsync 等方法,选择同步/异步通道 | | 僵死检测 | MaxInactiveTime + ConcurrentDictionary<Int32, DateTime> 跟踪,超时标记不存活 | | 自动注册 | AutoRegistry=true,首次登录自动创建 App 记录并分配 Secret | | 免打扰时段 | QuietTime="09:00-12:00,13:00-18:00" 多时段,含跨天 "23:00-02:00" | | 多 TFM | 核心SDK: net45/net461/netstandard2.0/netstandard2.1;入口: net10.0 | | 内嵌 Data.dll | AntJob.Web NuGet 包通过 `<BuildOutputInPackage>` 内嵌 AntJob.Data.dll | | 密码保护 | SaltPasswordProvider(Algorithm=md5, SaltTime=60) 避免明文传输 | | 雪花Id分表 | AppHistory/AppMessage 使用 Int64+DataScale=time,大数据量自动按时间分表 | --- ## 5. 通信矩阵 | 通道 | 协议 | 端口 | 用途 | |------|------|------|------| | Agent ↔ Server | NewLife.Remoting TCP | 9999 | 作业上报/任务申请/进度报告 | | Agent ↔ Web | HTTP REST (JSON) | Kestrel 动态 | JWT 鉴权 API | | 浏览器 ↔ Web | HTTP | Kestrel 动态 | Cube MVC 管理页面 | --- ## 6. 待补充 - [ ] 消息调度 Message 模式的完整时序图 - [ ] SQL 调度 topic_ 列自动消息生产的详细机制 - [ ] 多调度中心主备切换 Failover 逻辑细节 - [ ] 星尘注册中心 Stardust 集成细节 - [ ] 性能测试基准数据(吞吐/延迟/资源)