外部调度平台接入(Aegis.Scheduler)
Aegis.Scheduler 是 Aegis 对接外部工作流调度平台的能力,用于在代码里以编程方式管理定时任务(创建、触发、查询、删除)。
组成结构
由三个 NuGet 包组成:
| 包 | 定位 | 何时安装 |
|---|---|---|
Aegis.Scheduler | 抽象层。定义 ISchedulerClient 接口、中立任务类型(ShellTask/SqlTask/HttpTask)、数据模型。业务代码只依赖它 | 总是(被实现层自动引入) |
Aegis.Scheduler.DolphinScheduler | 对接 Apache DolphinScheduler 的实现层。提供 DsShellTask 等海豚专有类型和 DAG 扩展 | 生产环境使用海豚调度时 |
Aegis.Scheduler.Stub | 空实现。所有方法静默成功,不真正调用调度平台 | 开发/测试环境没有海豚集群时,让应用能正常启动 |
核心原则:业务代码只依赖抽象层的中立类型。后续替换调度平台时业务代码无需修改,只需替换实现层包。
何时使用
适合场景:
- 已经部署了 DolphinScheduler 集群,需要在代码里以编程方式管理定时任务
- 希望用一行代码完成原来在调度平台界面上需要多步操作的工作流创建
- 需要把 ETL 任务、数据加工任务的调度纳入应用代码管理
如果任务是「业务方法在进程内延后执行」,用 任务调度(Aegis.Jobs),不需要外部调度平台。
与 Aegis.Jobs 的区别
| 对比项 | Aegis.Jobs | Aegis.Scheduler |
|---|---|---|
| 执行位置 | 应用进程内 | 外部调度平台集群 |
| 持久化 | 应用自身数据库 | 调度平台自身存储 |
| 任务编排 | 单任务延迟执行 | DAG 工作流(海豚专有) |
| 典型场景 | 业务回调、延迟通知、补偿逻辑 | 大数据 ETL、跨系统数据加工 |
两者互补,可以在同一个项目里同时使用。
引入与配置
方式一:从 appsettings.json 读取配置(推荐)
appsettings.json:
{
"DolphinScheduler": {
"BaseUrl": "http://192.168.1.10:12345/dolphinscheduler",
"Token": "你的 Access Token",
"DefaultProjectCode": 1234567890,
"TimeoutSeconds": 30
}
}
这里 AddDolphinScheduler(IConfiguration) 重载会从传入配置的 DolphinScheduler 节读取上表各项。下面给出两种宿主模型的完整写法。
.NET 6+ 最小宿主(Program.cs) —— 配置来自 builder.Configuration:
var builder = WebApplication.CreateBuilder(args);
// builder.Configuration 已自动加载 appsettings.json,
// 方法内部读取其中的 "DolphinScheduler" 节
builder.Services.AddDolphinScheduler(builder.Configuration);
var app = builder.Build();
app.Run();
传统 Startup.cs —— 配置通过构造函数注入 IConfiguration:
public class Startup
{
private readonly IConfiguration _configuration;
// 框架自动注入 IConfiguration(已合并 appsettings.json、环境变量、命令行等)
public Startup(IConfiguration configuration)
{
_configuration = configuration;
}
public void ConfigureServices(IServiceCollection services)
{
// 传入 IConfiguration,方法内部读取 "DolphinScheduler" 节
services.AddDolphinScheduler(_configuration);
}
}
节名必须与
appsettings.json里的键一致(默认DolphinScheduler)。如需自定义,用第三个参数:AddDolphinScheduler(configuration, sectionName: "MyScheduler")。
方式二:配置回调
不想把配置写在 appsettings.json,或需要根据条件动态计算配置时,用回调(写在方式一的 ConfigureServices / var builder = ... 处):
services.AddDolphinScheduler(options =>
{
options.BaseUrl = "http://192.168.1.10:12345/dolphinscheduler";
options.Token = "你的 Access Token";
options.DefaultProjectCode = 1234567890;
});
方式三:直接传配置实例
已有 DolphinSchedulerOptions 对象时直接传入(同上写在注册处):
var options = new DolphinSchedulerOptions
{
BaseUrl = "http://192.168.1.10:12345/dolphinscheduler",
Token = "你的 Access Token",
DefaultProjectCode = 1234567890
};
services.AddDolphinScheduler(options);
方式四:用 Stub 接入(无海豚环境)
开发/测试环境没有海豚集群时,用 Stub 让应用正常启动、注入 ISchedulerClient 不报错:
// 只安装 Aegis.Scheduler.Stub 包
services.AddSchedulerStub();
Stub 所有方法静默返回假数据(SUCCESS),不发起任何网络请求。通过环境判断切换生产/Stub 实现(最小宿主示例,builder.Configuration / builder.Environment 均由 WebApplication 提供):
var builder = WebApplication.CreateBuilder(args);
if (builder.Environment.IsDevelopment())
builder.Services.AddSchedulerStub(); // 开发环境:静默成功
else
builder.Services.AddDolphinScheduler(builder.Configuration); // 生产:连真实海豚
var app = builder.Build();
配置项说明
| 配置项 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
BaseUrl | string | 是 | — | DolphinScheduler API 基础地址 |
Token | string | 是 | — | Access Token,用于 API 鉴权 |
DefaultProjectCode | long | 是 | — | 默认项目编码,所有操作在该项目下进行 |
TimeoutSeconds | int | 否 | 30 | HTTP 请求超时秒数 |
Token 获取与过期处理
创建 Token
- 登录 DolphinScheduler 控制台
- 进入「安全中心 → 令牌管理」
- 点击「创建令牌」,选择过期时间(或不设置以创建长期令牌)
- 复制生成的 Token 字符串,填入配置
过期处理
Token 有有效期,过期后所有调度 API 调用会返回 401 认证失败。当前 Token 是静态配置,过期后需要:
- 在海豚控制台重新创建 Token
- 更新
appsettings.json中的Token值 - 重启应用使配置生效
安全建议
Token 是敏感凭据,不要明文提交到代码仓库。推荐做法:
- 开发环境:用 .NET User Secrets(
dotnet user-secrets set "DolphinScheduler:Token" "xxx") - 生产环境:用环境变量或配置中心(Vault / Nacos / Consul)注入,
appsettings.json只留占位符 DefaultProjectCode在 DolphinScheduler 的「项目管理」中查看对应项目的编码
核心接口 ISchedulerClient
ISchedulerClient 只覆盖所有调度平台都具备的通用能力,共 5 个方法。注入后即可调用:
public class ReportJobService
{
private readonly ISchedulerClient _scheduler;
public ReportJobService(ISchedulerClient scheduler)
{
_scheduler = scheduler;
}
}
| 方法 | 能力 |
|---|---|
AddTaskAsync | 创建单个定时任务(一键完成创建工作流、上线、创建调度、上线) |
TriggerAsync | 立即手动触发一次执行,可带运行时参数 |
TriggerAtAsync | 在指定的未来时间点触发一次执行 |
GetLastExecutionAsync | 查询最近一次执行状态 |
RemoveTaskAsync | 删除定时任务(下线 Schedule 并下线 Workflow) |
TaskId 主键与流转
AddTaskAsync 返回的 ScheduledTaskResult.TaskId 是后续所有操作的唯一凭证:
调用方应保存 TaskId(如存入业务表),后续触发、查询、删除都传它。在海豚实现中 TaskId 对应 Workflow Code(数字字符串),由实现层内部解析。
任务类型选择指引
抽象层提供中立的 ShellTask/SqlTask/HttpTask,海豚实现层提供 DsShellTask/DsSqlTask/DsHttpTask/DsSubProcessTask。选择依据:
| 类型 | 适用场景 | 跨平台 |
|---|---|---|
中立 ShellTask/SqlTask/HttpTask | 普通单任务,只用脚本/SQL/HTTP | 是,换调度平台不用改 |
DsShellTask/DsSqlTask/DsHttpTask | 需要海豚专有配置(如 SQL 数据源类型、延迟时间)或 DAG 编排 | 否,依赖海豚实现层 |
DsSubProcessTask | 调用另一个已存在的海豚工作流 | 否,海豚专有 |
中立任务通过翻译器自动映射到海豚任务并补默认值(如 SQL 默认数据源类型 MYSQL)。需要覆盖默认值时直接用 Ds* 类型。
HTTP 任务详解(最常用)
实际业务中大多数调度任务都是「在某个时间点回调一个 HTTP 接口」——比如定时拉取数据、定时触发报表生成、定时通知下游系统。因此 HttpTask 是日常使用频率最高的任务类型。
何时用中立 HttpTask vs 海豚 DsHttpTask
| 情况 | 选哪个 | 原因 |
|---|---|---|
| 普通定时回调一个接口 | 中立 HttpTask | 平台无关,换调度平台不用改业务代码 |
| 需要在 DAG 里编排(A 完成后调 B) | DsHttpTask | 需要用 DependsOn 建立依赖,DAG 是海豚专有能力 |
| 需要海豚 Worker 分组等专有配置 | DsHttpTask | 中立类型不暴露平台特有字段 |
简单记:单任务定时回调用中立
HttpTask,多任务编排或需要海豚特性时才用DsHttpTask。 两者 API 一致(WithUrl/WithMethod/WithBody),切换成本低。
字段说明
| 字段 | 链式方法 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|---|
Url | WithUrl(url) | string | 是 | — | 请求地址,海豚 Worker 能访问的内网/公网 URL |
HttpMethod | WithMethod(method) | string | 否 | GET | GET / POST / PUT / DELETE |
Body | WithBody(body) | string? | 否 | null | 请求体,常用于 POST/PUT 传 JSON |
WorkerGroup | WithWorkerGroup(g) | string | 否 | default | 执行该任务的 Worker 分组 |
Timeout | WithTimeout(sec) | int | 否 | 0(不限) | 单次执行超时秒数 |
RetryTimes | 直接赋值属性 | int | 否 | 0 | 失败重试次数(继承自 TaskDefinitionBase) |
RetryInterval | 直接赋值属性 | int | 否 | 1 | 重试间隔秒数(继承自 TaskDefinitionBase) |
Url为空会在创建时抛InvalidOperationException,请在AddTaskAsync之前用WithUrl设置。
创建一个 HTTP 定时任务(最小写法)
public class SyncJobService
{
private readonly ISchedulerClient _scheduler;
public SyncJobService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task<string> CreateSyncTaskAsync()
{
var task = new HttpTask("nightly-sync") // 中立类型
.WithUrl("http://internal-api/data/sync") // 要回调的接口
.WithMethod("POST"); // POST 触发同步
// 每天 02:30 执行一次
var created = await _scheduler.AddTaskAsync(
"nightly-sync", task, "0 30 2 * * ?");
return created.TaskId; // 持久化到业务表,后续触发/查询/删除都用它
}
}
带请求体(POST JSON)
定时推送数据到下游系统的典型写法——基于上面的 SyncJobService,改任务配置:
public class PushReportService
{
private readonly ISchedulerClient _scheduler;
public PushReportService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task<string> CreatePushTaskAsync()
{
var payload = """{"batchId":"B20260701","type":"full"}""";
var task = new HttpTask("push-report")
.WithUrl("http://internal-api/report/push")
.WithMethod("POST")
.WithBody(payload); // POST 请求体
var created = await _scheduler.AddTaskAsync("push-report", task, "0 0 6 * * ?");
return created.TaskId;
}
}
重试与超时配置
接口偶发失败时,在任务对象上配置重试(重试字段继承自 TaskDefinitionBase,是公开属性,链式调用之后直接赋值):
public async Task<string> CreateFlakyTaskAsync()
{
var task = new HttpTask("flaky-sync")
.WithUrl("http://internal-api/data/sync")
.WithMethod("POST");
// 链式方法只覆盖 Url/Method/Body,重试与超时直接赋值
task.RetryTimes = 3; // 失败后最多重试 3 次
task.RetryInterval = 60; // 每次重试间隔 60 秒
task.Timeout = 300; // 单次执行超时 5 分钟
var created = await _scheduler.AddTaskAsync("flaky-sync", task, "0 0 2 * * ?");
return created.TaskId;
}
在 DAG 中编排 HTTP 任务(用 DsHttpTask)
需要「数据准备 → 调接口 A → 调接口 B」这类串行/并行编排时,改用 DsHttpTask 并通过 DependsOn 建立依赖。此时需直接依赖 Aegis.Scheduler.DolphinScheduler 实现层:
using Aegis.Scheduler.DolphinScheduler;
using Aegis.Scheduler.DolphinScheduler.Tasks;
public class HttpPipelineService
{
private readonly ISchedulerClient _scheduler;
public HttpPipelineService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task<string> CreateAsync()
{
var prepare = new DsHttpTask("prepare")
.WithUrl("http://internal-api/data/prepare")
.WithMethod("POST");
var report = new DsHttpTask("report")
.WithUrl("http://internal-api/report/generate")
.WithMethod("POST")
.DependsOn(prepare); // prepare 完成后才执行 report
// AddDagTaskAsync 是海豚实现层扩展方法
var created = await _scheduler.AddDagTaskAsync(
"http-pipeline",
new[] { prepare, report },
"0 0 3 * * ?");
return created.TaskId;
}
}
执行顺序:prepare → report。
常见坑
- URL 必须海豚 Worker 能访问:
Url是由海豚 Worker 节点发起请求,不是应用服务器。内网接口要确保 Worker 与其在同一网络可达。 - GET 不要带 Body:
WithBody主要给 POST/PUT 用;GET 带 Body 可能被目标服务忽略。 - 方法名大写:
WithMethod("post")虽不报错,但建议传大写POST,与 HTTP 规范一致,避免某些网关区分大小写。 - 超时与重试:HTTP 任务对象的
WithTimeout控制单次执行上限;失败重试由任务基类的RetryTimes/RetryInterval属性控制,直接在任务对象上赋值即可(如task.RetryTimes = 3;)。
完整端到端示例
场景一:创建 Shell 定时任务并管理全生命周期
用中立 ShellTask,全程只依赖抽象层类型:
public class ReportJobService
{
private readonly ISchedulerClient _scheduler;
public ReportJobService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task RunAsync(string savedTaskId = null)
{
// 1. 创建(拿到 TaskId,业务方应持久化)
var task = new ShellTask("daily-report").WithScript("sh /opt/jobs/generate-report.sh");
var created = await _scheduler.AddTaskAsync("daily-report", task, "0 0 2 * * ?");
var taskId = created.TaskId;
Console.WriteLine($"创建成功,TaskId = {taskId}");
// 2. 立即触发一次
var exec = await _scheduler.TriggerAsync(taskId);
Console.WriteLine($"触发成功,ExecutionId = {exec.ExecutionId},状态 = {exec.State}");
// 3. 查询最近一次执行
var last = await _scheduler.GetLastExecutionAsync(taskId);
Console.WriteLine($"最近执行状态 = {last.State}");
// 4. 删除
await _scheduler.RemoveTaskAsync(taskId);
}
}
场景二:定点执行(TriggerAtAsync vs TriggerAsync)
两者区别:
TriggerAsync(taskId):立即执行一次TriggerAtAsync(taskId, runAt):在指定的未来时间点执行一次(海豚实现通过scheduleTime参数下发)
前提:调用前必须先用
AddTaskAsync创建任务并拿到TaskId,触发操作都基于这个TaskId。
public class OneShotJobService
{
private readonly ISchedulerClient _scheduler;
public OneShotJobService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task RunAsync()
{
// 1. 先创建一个定时任务(这里 cron 只是占位,定点触发不依赖它)
var task = new HttpTask("ad-hoc-report")
.WithUrl("http://internal-api/report/generate")
.WithMethod("POST");
var created = await _scheduler.AddTaskAsync("ad-hoc-report", task, "0 0 0 * * ?");
var taskId = created.TaskId; // 后续触发/查询都靠它
// 2a. 立即触发一次
var now = await _scheduler.TriggerAsync(taskId);
Console.WriteLine($"立即触发,ExecutionId = {now.ExecutionId},状态 = {now.State}");
// now.ScheduleTime 为 null(立即触发没有计划时间)
// 2b. 定点触发:明天 10:00 执行一次
var runAt = DateTime.Now.AddDays(1).Date.AddHours(10);
var scheduled = await _scheduler.TriggerAtAsync(taskId, runAt);
Console.WriteLine($"定点触发,ExecutionId = {scheduled.ExecutionId},状态 = {scheduled.State}");
// scheduled.State: WAITING(等待到点执行)
// scheduled.ScheduleTime: 传入的 runAt
}
}
注意:
scheduleTime按调度平台服务器时区。若应用服务器与海豚服务器时区不同,调用方需自行换算。
场景三:海豚 DAG 工作流(多任务编排)
DAG 不是所有调度平台的通用能力,已下沉到海豚实现层。需要 DAG 编排时直接依赖实现层,用 Ds*Task + DependsOn + AddDagTaskAsync 扩展方法:
using Aegis.Scheduler.DolphinScheduler;
using Aegis.Scheduler.DolphinScheduler.Tasks;
public class EtlPipelineService
{
private readonly ISchedulerClient _scheduler;
public EtlPipelineService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task<string> CreatePipelineAsync()
{
var extract = new DsSqlTask("extract")
.WithSql("INSERT INTO dwd_orders SELECT * FROM ods_orders", sqlType: 1)
.WithDatasource("postgres-prod")
.WithDatasourceType("POSTGRESQL");
var transform = new DsShellTask("transform")
.WithScript("python /opt/jobs/transform.py")
.DependsOn(extract);
var notify = new DsHttpTask("notify")
.WithUrl("http://internal-api/report/done")
.WithMethod("POST")
.DependsOn(transform);
// AddDagTaskAsync 是海豚实现层扩展方法(不在抽象层 ISchedulerClient 上)
var created = await _scheduler.AddDagTaskAsync(
"etl-pipeline",
new[] { extract, transform, notify },
"0 30 3 * * ?");
return created.TaskId;
}
}
执行顺序:extract → transform → notify。
场景四:无海豚环境用 Stub 接入
// Startup.cs
public class Startup
{
private readonly IConfiguration _configuration;
private readonly IWebHostEnvironment _env;
// 框架自动注入这两个对象
public Startup(IConfiguration configuration, IWebHostEnvironment env)
{
_configuration = configuration;
_env = env;
}
public void ConfigureServices(IServiceCollection services)
{
if (_env.IsDevelopment())
services.AddSchedulerStub(); // 静默成功,不调海豚
else
services.AddDolphinScheduler(_configuration); // 读取 "DolphinScheduler" 节
}
}
Stub 注入的 ISchedulerClient 所有方法返回假数据:AddTaskAsync 返回 stub-xxxxxxxx 格式的 TaskId,TriggerAsync 返回 SUCCESS。业务代码在开发环境能完整跑通流程,不依赖真实海豚集群。
返回值字段说明
ScheduledTaskResult(创建结果)
| 字段 | 类型 | 说明 |
|---|---|---|
TaskId | string | 任务唯一标识,后续操作凭证。海豚实现中为 Workflow Code |
Name | string | 任务名称 |
CronExpression | string | cron 表达式 |
Status | string | 状态,通常为 ONLINE |
CreateTime | DateTime | 创建时间 |
TaskExecutionInfo(执行状态)
| 字段 | 类型 | 说明 |
|---|---|---|
ExecutionId | string | 执行实例唯一标识。海豚实现中为 ProcessInstance Id |
TaskId | string | 所属任务唯一标识 |
Name | string | 任务名称 |
State | string | 执行状态:SUCCESS / FAILURE / RUNNING / WAITING / UNKNOWN |
ScheduleTime | DateTime? | 计划执行时间(定点触发时由 TriggerAtAsync 指定,立即触发为 null) |
StartTime | DateTime | 实际开始时间 |
EndTime | DateTime? | 结束时间 |
Parameters | Dictionary | 运行时参数 |
错误处理
所有调度操作失败抛 SchedulerException(抽象层基类),海豚实现抛 DolphinSchedulerException(继承自 SchedulerException):
public class SafeJobService
{
private readonly ISchedulerClient _scheduler;
public SafeJobService(ISchedulerClient scheduler) => _scheduler = scheduler;
public async Task<string?> CreateAsync()
{
var task = new HttpTask("daily-job")
.WithUrl("http://internal-api/job/run")
.WithMethod("POST");
try
{
var result = await _scheduler.AddTaskAsync("daily-job", task, "0 0 8 * * ?");
return result.TaskId;
}
catch (SchedulerException ex)
{
// ex.Message 含具体原因
// 常见场景:
// - "认证失败: Token 无效或已过期"(HTTP 401,检查 Token)
// - "DolphinScheduler API 错误: ..."(海豚业务错误)
// - "无效的 taskId"(taskId 无法解析为工作流编码)
return null;
}
}
}
常见错误场景:
| 场景 | 原因 | 处理 |
|---|---|---|
| HTTP 401 认证失败 | Token 无效或已过期 | 重新创建 Token 并更新配置 |
| 创建时报重复 | 同名 Workflow 已存在 | 先 RemoveTaskAsync 或换名 |
| 无效的 taskId | 传入的 taskId 无法解析 | 确认使用 AddTaskAsync 返回的 TaskId |
cron 表达式速查
DolphinScheduler 使用 6 位 cron(秒 分 时 日 月 周),常用示例:
| 含义 | 表达式 |
|---|---|
| 每天凌晨 2:00 | 0 0 2 * * ? |
| 每小时整点 | 0 0 0/1 * * ? |
| 每天 8:30 | 0 30 8 * * ? |
| 每月 1 号 0:00 | 0 0 0 1 * ? |
| 工作日 9:00 | 0 0 9 ? * MON-FRI |
第一位是「秒」,第 6 位是「周」(用
?或MON-FRI等)。注意和 Linux 的 5 位 cron 不同。AddTaskAsync只做非空校验,格式错误的表达式会提交到 DolphinScheduler 由服务端报错,提交前建议自行校验。
扩展点
实现自定义海豚任务类型
继承 DsTaskDefinitionBase,实现 TaskType 和 ToDsParams:
using System.Text.Json;
using Aegis.Scheduler.DolphinScheduler.Tasks;
public class DsPythonTask : DsTaskDefinitionBase
{
public override string TaskType => "PYTHON";
public string Script { get; private set; } = string.Empty;
public DsPythonTask(string name) : base(name) { }
public DsPythonTask WithScript(string script)
{
Script = script;
return this;
}
public override string ToDsParams()
{
var paramObj = new
{
resourceList = new object[] { },
localParams = new object[] { },
rawScript = Script
};
return JsonSerializer.Serialize(paramObj);
}
}
ToDsParams 返回的 JSON 原样写入 DolphinScheduler 的 taskParams 字段,结构需匹配对应任务类型在海豚中的定义。
对接新的调度平台
抽象层只依赖 .NET BCL。对接其他调度平台时,实现 ISchedulerClient 即可。业务代码只依赖 ISchedulerClient,无需修改:
- 中立
ShellTask/SqlTask/HttpTask通过你的实现层翻译成目标平台的任务格式 - 不支持的能力(如该平台没有 DAG)不需要实现——DAG 本就不在抽象层契约里
DAG 设计说明
DAG 多任务编排是海豚、Airflow 等高级调度平台的能力,不是所有调度平台的通用契约(例如 XXL-Job 只支持单任务)。因此:
- 抽象层
ISchedulerClient不包含 DAG 方法,ITaskDefinition不包含依赖关系 - DAG 编排作为海豚实现层的扩展方法
AddDagTaskAsync提供 - 海豚实现层的
DsTaskDefinitionBase重新引入DependsOn,供 DAG 编排使用
这样接入简单调度平台时,实现方只需实现 5 个通用方法,不会被 DAG 概念拖累。
边界与限制
- TaskId 精确定位:所有操作基于
TaskId(Workflow Code),精确高效,不再依赖名称模糊匹配 - cron 不做格式校验:表达式只做非空校验,格式错误由服务端报错
- 单项目作用域:所有操作基于
DefaultProjectCode指定的项目,不支持跨项目 - 无日志拉取:当前不支持获取任务执行日志,需要到 DolphinScheduler 控制台查看
- Token 静态配置:Token 不支持动态刷新,过期需手动更新配置并重启