TKService 轻量级任务调度框架设计与实现
TKService 内置了一套高效的任务调度框架,帮助开发者轻松管理周期性作业与批处理任务。无需依赖外部定时服务,即可在应用内部实现数据库备份、报表生成、数据清洗等自动化运维需求。
一、内置调度方案的优势
传统方式通常依赖操作系统层面的定时任务(如 Linux crontab 或 Windows Task Scheduler),但在现代分布式应用中暴露诸多局限:
- 维护割裂:调度规则与业务代码分离,版本管理困难
- 权限壁垒:调整执行计划需系统管理员介入
- 观测盲区:执行过程黑盒化,故障定位效率低
- 容错缺失:异常场景缺乏自动恢复与通知机制
- 集群困境:多节点环境下任务重复执行或遗漏难以规避
内嵌式调度引擎将执行逻辑纳入应用生命周期,从根本上解决上述痛点。
二、核心架构概览
该框架围绕 JobEngine 类构建,具备以下能力矩阵:
| 特性 | 说明 |
|---|---|
| 多模式触发 | Cron 表达式、固定间隔、指定时刻、一次性延迟 |
| 依赖编排 | 支持 DAG 形式的前置任务依赖声明 |
| 并发防护 | 同一任务实例互斥执行,防止重叠 |
| 故障自愈 | 阶梯式重试与降级策略 |
| 全链路追踪 | 执行轨迹持久化,支持审计与复盘 |
核心调度器实现如下:
class JobEngine:
"""分布式友好的任务调度引擎"""
def __init__(self, worker_threads=16):
self._logger = logging.getLogger("JobEngine")
self._registry = {} # 任务注册表
self._executing = {} # 运行中任务句柄
self._archive = {} # 执行档案
self._executor = ThreadPoolExecutor(max_workers=worker_threads)
self._mutex = threading.RLock()
self._active = False
self._dispatch_thread = None
def enroll(self, job_name, handler, schedule=None, interval=None,
trigger_at=None, prerequisites=None,
max_attempts=0, backoff_secs=60):
"""注册调度任务"""
with self._mutex:
if job_name in self._registry:
self._logger.warning(f"任务 {job_name} 存在覆盖风险")
entry = {
"name": job_name,
"handler": handler,
"schedule": schedule, # Cron 字符串
"interval": interval, # 秒级固定间隔
"trigger_at": trigger_at, # 首次触发时刻
"last_fired": None,
"next_fire": None,
"prerequisites": prerequisites or [],
"max_attempts": max_attempts,
"backoff_secs": backoff_secs,
"state": "PENDING", # PENDING/RUNNING/ERROR
"active": True
}
self._compute_next_fire(entry)
self._registry[job_name] = entry
self._archive[job_name] = []
self._logger.info(f"任务注册成功: {job_name}, 预计触发: {entry['next_fire']}")
return job_name
def _compute_next_fire(self, entry):
"""推算下次触发时间戳"""
now = datetime.now()
if entry["trigger_at"] and now < entry["trigger_at"]:
entry["next_fire"] = entry["trigger_at"]
return
if entry["schedule"]:
base = entry["last_fired"] or now
itr = croniter(entry["schedule"], base)
entry["next_fire"] = itr.get_next(datetime)
return
if entry["interval"]:
delta = timedelta(seconds=entry["interval"])
entry["next_fire"] = (entry["last_fired"] or now) + delta
return
entry["next_fire"] = None
def ignite(self):
"""启动调度循环"""
with self._mutex:
if self._active:
return
self._active = True
self._dispatch_thread = threading.Thread(target=self._dispatch_loop)
self._dispatch_thread.daemon = True
self._dispatch_thread.start()
self._logger.info("调度引擎已启动")
def halt(self):
"""优雅停机"""
with self._mutex:
self._active = False
if self._dispatch_thread:
self._dispatch_thread.join(timeout=2)
self._executor.shutdown(wait=False)
self._logger.info("调度引擎已停止")
def _dispatch_loop(self):
"""调度主循环"""
while self._active:
try:
self._scan_and_fire()
time.sleep(0.5) # 500ms 精度轮询
except Exception as ex:
self._logger.error(f"调度循环异常: {ex}")
def _scan_and_fire(self):
"""扫描并触发到期任务"""
now = datetime.now()
with self._mutex:
for name, entry in self._registry.items():
if not entry["active"] or entry["state"] == "RUNNING":
continue
if entry["next_fire"] and now >= entry["next_fire"]:
if not self._check_prerequisites(entry):
continue
self._fire_job(entry)
三、任务执行与容错机制
任务提交至线程池后,框架负责状态流转、结果归档及异常恢复:
def _fire_job(self, entry):
"""投递任务至执行队列"""
entry["state"] = "RUNNING"
entry["last_fired"] = datetime.now()
trace = {
"fired_at": entry["last_fired"],
"settled_at": None,
"outcome": "RUNNING",
"return_val": None,
"exception": None,
"attempts": 0
}
self._archive[entry["name"]].append(trace)
self._executor.submit(self._wrap_handler, entry, trace)
self._logger.info(f"任务 {entry['name']} 开始执行")
def _wrap_handler(self, entry, trace):
"""包装执行逻辑,处理重试与归档"""
attempts = 0
result = error = None
try:
result = entry["handler"]()
status = "COMPLETED"
self._logger.info(f"任务 {entry['name']} 执行成功")
except Exception as ex:
error = str(ex)
status = "FAILED"
self._logger.error(f"任务 {entry['name']} 异常: {error}")
while attempts < entry["max_attempts"]:
attempts += 1
self._logger.info(f"任务 {entry['name']} 第 {attempts} 次重试")
time.sleep(entry["backoff_secs"] * attempts) # 指数退避
try:
result = entry["handler"]()
status = "RECOVERED"
error = None
self._logger.info(f"任务 {entry['name']} 重试成功")
break
except Exception as retry_ex:
error = str(retry_ex)
self._logger.error(f"重试失败: {error}")
finally:
with self._mutex:
entry["state"] = "PENDING" if status in ("COMPLETED", "RECOVERED") else "ERROR"
self._compute_next_fire(entry)
trace.update({
"settled_at": datetime.now(),
"outcome": status,
"return_val": result,
"exception": error,
"attempts": attempts
})
if status == "FAILED":
self._notify_failure(entry["name"], error)
关键设计点:
- 异步执行避免阻塞调度线程
- 渐进式退避重试,降低系统压力
- 执行档案完整记录入参、出参、异常堆栈
四、Cron 表达式实战
框架集成 croniter 解析器,支持标准 Unix Cron 语法:
engine = JobEngine()
# 每日 02:30 执行全量备份
engine.enroll(
job_name="full_backup",
handler=perform_backup,
schedule="30 2 * * *",
max_attempts=2,
backoff_secs=180
)
# 工作日 09:00-18:00 每 15 分钟健康探测
engine.enroll(
job_name="health_probe",
handler=check_service_health,
schedule="*/15 9-18 * * 1-5"
)
# 每月 1 日 00:00 执行账单结算
engine.enroll(
job_name="monthly_billing",
handler=generate_invoices,
schedule="0 0 1 * *",
prerequisites=["full_backup"] # 确保备份完成后执行
)
engine.ignite()
五、依赖链与执行准入
复杂业务流程中,任务间存在明确的先后约束。框架通过前置条件检查实现依赖管控:
def _check_prerequisites(self, entry):
"""验证前置任务状态"""
for prereq_name in entry["prerequisites"]:
if prereq_name not in self._registry:
self._logger.warning(f"前置任务 {prereq_name} 未注册")
return False
prereq = self._registry[prereq_name]
# 从未执行则阻塞
if not prereq["last_fired"]:
return False
# 执行中则等待
if prereq["state"] == "RUNNING":
return False
# 前置失败则传播失败
if prereq["state"] == "ERROR":
return False
# 确保获取的是最新结果
if entry["last_fired"] and prereq["last_fired"] <= entry["last_fired"]:
if not prereq["next_fire"] or prereq["next_fire"] > datetime.now():
continue
return False
return True
该机制支持构建任意深度的有向无环图(DAG),自动处理依赖就绪通知与失败传播。
六、运行时管控接口
除自动调度外,框架暴露丰富的运维接口:
def pause(self, job_name):
"""暂停任务调度"""
with self._mutex:
if job_name in self._registry:
self._registry[job_name]["active"] = False
return True
return False
def resume(self, job_name):
"""恢复任务调度"""
with self._mutex:
if job_name in self._registry:
self._registry[job_name]["active"] = True
self._compute_next_fire(self._registry[job_name])
return True
return False
def trigger_immediately(self, job_name):
"""强制立即触发"""
with self._mutex:
if job_name not in self._registry:
return False
entry = self._registry[job_name]
if entry["state"] == "RUNNING":
return False
self._fire_job(entry)
return True
def reconfigure(self, job_name, schedule=None, interval=None,
max_attempts=None, active=None):
"""热更新任务参数"""
with self._mutex:
if job_name not in self._registry:
return False
entry = self._registry[job_name]
if schedule: entry["schedule"] = schedule
if interval is not None: entry["interval"] = interval
if max_attempts is not None: entry["max_attempts"] = max_attempts
if active is not None: entry["active"] = active
self._compute_next_fire(entry)
return True
def inspect(self, job_name=None, history_limit=20):
"""查询任务状态与历史"""
with self._mutex:
if job_name:
if job_name not in self._registry:
return None
return {
"config": {k: v for k, v in self._registry[job_name].items()
if k not in ("handler",)},
"history": sorted(
self._archive.get(job_name, []),
key=lambda x: x["fired_at"],
reverse=True
)[:history_limit]
}
return {name: {
"state": entry["state"],
"last_fired": entry["last_fired"],
"next_fire": entry["next_fire"],
"active": entry["active"]
} for name, entry in self._registry.items()}
七、多渠道告警通知
任务持续失败时,框架自动触发分级告警:
def _notify_failure(self, job_name, error_detail):
"""发送故障通知"""
payload = f"""[任务告警] {job_name}
时间: {datetime.now().isoformat()}
异常: {error_detail}
节点: {socket.gethostname()}
"""
self._logger.critical(payload)
# 邮件通道
self._send_email(payload, subject=f"任务失败: {job_name}")
# 即时通讯通道
self._send_webhook(payload, channel="ops_alert")
def _send_email(self, body, subject, receivers=None):
try:
cfg = self._load_alert_config()
if not cfg.get("smtp_enabled"):
return
msg = MIMEText(body, _charset="utf-8")
msg["Subject"] = subject
msg["From"] = cfg["smtp_sender"]
msg["To"] = ", ".join(receivers or cfg["smtp_receivers"])
with smtplib.SMTP(cfg["smtp_host"], cfg["smtp_port"]) as srv:
srv.starttls()
srv.login(cfg["smtp_user"], cfg["smtp_pass"])
srv.send_message(msg)
except Exception as ex:
self._logger.error(f"邮件发送失败: {ex}")
def _send_webhook(self, body, channel):
try:
cfg = self._load_alert_config()
webhook_url = cfg.get(f"{channel}_url")
if not webhook_url:
return
requests.post(webhook_url, json={
"text": body,
"priority": "high"
}, timeout=10)
except Exception as ex:
self._logger.error(f"Webhook 发送失败: {ex}")
八、生产环境最佳实践
基于实际部署经验,建议遵循以下原则:
| 幂等设计 | 任务 handler 需支持多次执行结果一致,防止重试副作用 |
| 超时控制 | 长耗时任务内置超时机制,避免线程池耗尽 |
| 资源隔离 | CPU 密集型与 IO 密集型任务分池调度 |
| 依赖扁平 | 控制依赖链深度,降低级联故障影响面 |
| 监控覆盖 | 关键任务配置 P99 延迟、失败率告警 |
| 灰度发布 | 调度规则变更先小流量验证,再全量生效 |
九、完整业务示例
def main():
# 初始化引擎
orchestrator = JobEngine(worker_threads=24)
# 定义业务任务
def extract_oltp_data():
"""从 OLTP 库抽取变更数据"""
return etl_service.pull_incremental()
def transform_to_warehouse():
"""转换并加载至数据仓库"""
return etl_service.apply_transforms()
def refresh_dashboard():
"""刷新 BI 看板缓存"""
return cache_warmer.preheat_metrics()
# 注册任务流
orchestrator.enroll(
job_name="cdc_extract",
handler=extract_oltp_data,
schedule="0 */4 * * *", # 每 4 小时
max_attempts=3
)
orchestrator.enroll(
job_name="dw_load",
handler=transform_to_warehouse,
schedule="30 */4 * * *", # 抽取后 30 分钟
prerequisites=["cdc_extract"],
max_attempts=2
)
orchestrator.enroll(
job_name="cache_warm",
handler=refresh_dashboard,
interval=1800, # 每 30 分钟
prerequisites=["dw_load"]
)
# 启动服务
orchestrator.ignite()
# 暴露 HTTP 管理端点(示例)
@app.route("/jobs/status")
def api_status():
return jsonify(orchestrator.inspect())
@app.route("/jobs/<name>/trigger", methods=["POST"])
def api_trigger(name):
success = orchestrator.trigger_immediately(name)
return jsonify({"accepted": success})
# 优雅停机
atexit.register(orchestrator.halt)
if __name__ == "__main__":
main()