Files
2026-06-08 18:14:59 +08:00

715 lines
24 KiB
Python

#!/usr/bin/env python
# -*- coding: utf-8 -*-
"""
Scheduler API - 定时任务管理接口
提供定时任务的 CRUD 操作和管理功能
"""
from datetime import datetime, timedelta
from typing import List, Optional
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import select, func, or_
from sqlalchemy.ext.asyncio import AsyncSession
from app.database import get_db
from app.config import settings
from app.base_schema import PaginatedResponse, ResponseModel
from scheduler.model import SchedulerJob, SchedulerLog
from scheduler.schema import (
SchedulerJobCreate,
SchedulerJobUpdate,
SchedulerJobResponse,
SchedulerJobSimple,
SchedulerJobBatchDeleteIn,
SchedulerJobBatchDeleteOut,
SchedulerJobBatchUpdateStatusIn,
SchedulerJobBatchUpdateStatusOut,
SchedulerJobExecuteIn,
SchedulerJobExecuteOut,
SchedulerJobStatisticsOut,
SchedulerJobSearchRequest,
SchedulerLogResponse,
SchedulerLogBatchDeleteIn,
SchedulerLogBatchDeleteOut,
SchedulerLogCleanIn,
SchedulerLogCleanOut,
SchedulerStatusOut,
)
from scheduler.service import scheduler_service
router = APIRouter(prefix="/scheduler", tags=["定时任务管理"])
def _build_job_response(job: SchedulerJob) -> SchedulerJobResponse:
"""构建任务响应"""
return SchedulerJobResponse(
id=job.id,
application_id=job.application_id,
name=job.name,
code=job.code,
description=job.description,
group=job.group,
trigger_type=job.trigger_type,
trigger_type_display=job.get_trigger_type_display(),
cron_expression=job.cron_expression,
interval_seconds=job.interval_seconds,
run_date=job.run_date,
task_func=job.task_func,
task_args=job.task_args,
task_kwargs=job.task_kwargs,
status=job.status,
status_display=job.get_status_display(),
priority=job.priority,
max_instances=job.max_instances,
max_retries=job.max_retries,
timeout=job.timeout,
coalesce=job.coalesce,
allow_concurrent=job.allow_concurrent,
total_run_count=job.total_run_count,
success_count=job.success_count,
failure_count=job.failure_count,
success_rate=job.get_success_rate(),
last_run_time=job.last_run_time,
next_run_time=job.next_run_time,
last_run_status=job.last_run_status,
last_run_result=job.last_run_result,
remark=job.remark,
sort=job.sort,
sys_create_datetime=job.sys_create_datetime,
sys_update_datetime=job.sys_update_datetime,
)
def _build_log_response(log: SchedulerLog) -> SchedulerLogResponse:
"""构建日志响应"""
return SchedulerLogResponse(
id=log.id,
job_id=log.job_id,
job_name=log.job_name,
job_code=log.job_code,
status=log.status,
status_display=log.get_status_display(),
start_time=log.start_time,
end_time=log.end_time,
duration=log.duration,
result=log.result,
exception=log.exception,
traceback=log.traceback,
hostname=log.hostname,
process_id=log.process_id,
retry_count=log.retry_count,
sys_create_datetime=log.sys_create_datetime,
)
# ==================== SchedulerJob APIs ====================
@router.post("/job", response_model=SchedulerJobResponse, summary="创建定时任务")
async def create_scheduler_job(data: SchedulerJobCreate, db: AsyncSession = Depends(get_db)):
"""创建新的定时任务"""
# 检查任务编码是否已存在
result = await db.execute(
select(SchedulerJob).where(
SchedulerJob.code == data.code,
SchedulerJob.is_deleted == False # noqa: E712
)
)
if result.scalar_one_or_none():
raise HTTPException(status_code=400, detail=f"任务编码已存在: {data.code}")
# 创建任务
job = SchedulerJob(**data.model_dump())
db.add(job)
await db.commit()
await db.refresh(job)
# 如果任务是启用状态,添加到调度器
if job.is_enabled() and scheduler_service.is_running():
await scheduler_service.add_job(job)
return _build_job_response(job)
@router.get("/job/all", response_model=List[SchedulerJobSimple], summary="获取所有定时任务(简化版)")
async def get_all_scheduler_jobs(
application_id: str = Query(None, alias="applicationId", description="所属应用ID"),
db: AsyncSession = Depends(get_db),
):
"""获取所有定时任务(不分页,简化版)"""
conditions = [SchedulerJob.is_deleted == False] # noqa: E712
if application_id:
conditions.append(SchedulerJob.application_id == application_id)
else:
conditions.append(SchedulerJob.application_id.is_(None))
result = await db.execute(
select(SchedulerJob).where(*conditions).order_by(SchedulerJob.priority.desc(), SchedulerJob.name)
)
jobs = result.scalars().all()
return jobs
@router.get("/job", response_model=PaginatedResponse[SchedulerJobResponse], summary="获取定时任务列表")
async def get_scheduler_job_list(
page: int = Query(default=1, ge=1, description="页码"),
page_size: int = Query(default=settings.PAGE_SIZE, ge=1, le=settings.PAGE_MAX_SIZE, alias="pageSize", description="每页数量"),
application_id: str = Query(None, alias="applicationId", description="所属应用ID"),
name: Optional[str] = Query(default=None, description="任务名称"),
code: Optional[str] = Query(default=None, description="任务编码"),
group: Optional[str] = Query(default=None, description="任务分组"),
trigger_type: Optional[str] = Query(default=None, description="触发器类型"),
status: Optional[int] = Query(default=None, description="任务状态"),
db: AsyncSession = Depends(get_db)
):
"""获取定时任务列表(分页)"""
filters = [SchedulerJob.is_deleted == False] # noqa: E712
# 应用过滤
if application_id:
filters.append(SchedulerJob.application_id == application_id)
else:
filters.append(SchedulerJob.application_id.is_(None))
if name:
filters.append(SchedulerJob.name.ilike(f"%{name}%"))
if code:
filters.append(SchedulerJob.code.ilike(f"%{code}%"))
if group:
filters.append(SchedulerJob.group == group)
if trigger_type:
filters.append(SchedulerJob.trigger_type == trigger_type)
if status is not None:
filters.append(SchedulerJob.status == status)
# 查询总数
count_result = await db.execute(
select(func.count(SchedulerJob.id)).where(*filters)
)
total = count_result.scalar()
# 查询数据
offset = (page - 1) * page_size
result = await db.execute(
select(SchedulerJob).where(*filters)
.order_by(SchedulerJob.priority.desc(), SchedulerJob.sys_update_datetime.desc())
.offset(offset).limit(page_size)
)
jobs = result.scalars().all()
return PaginatedResponse(
items=[_build_job_response(job) for job in jobs],
total=total
)
@router.post("/job/batch/delete", response_model=SchedulerJobBatchDeleteOut, summary="批量删除定时任务")
async def batch_delete_scheduler_jobs(
data: SchedulerJobBatchDeleteIn,
db: AsyncSession = Depends(get_db)
):
"""批量删除定时任务"""
success_count = 0
failed_ids = []
for job_id in data.ids:
try:
result = await db.execute(
select(SchedulerJob).where(SchedulerJob.id == job_id)
)
job = result.scalar_one_or_none()
if job:
# 从调度器移除
if scheduler_service.is_running():
await scheduler_service.remove_job(job.code)
# 软删除
job.is_deleted = True
success_count += 1
else:
failed_ids.append(job_id)
except Exception:
failed_ids.append(job_id)
await db.commit()
return SchedulerJobBatchDeleteOut(count=success_count, failed_ids=failed_ids)
@router.post("/job/batch/update_status", response_model=SchedulerJobBatchUpdateStatusOut, summary="批量更新任务状态")
async def batch_update_scheduler_job_status(
data: SchedulerJobBatchUpdateStatusIn,
db: AsyncSession = Depends(get_db)
):
"""批量启用、禁用或暂停任务"""
result = await db.execute(
select(SchedulerJob).where(SchedulerJob.id.in_(data.ids))
)
jobs = result.scalars().all()
count = 0
for job in jobs:
job.status = data.status
count += 1
# 同步更新调度器
if scheduler_service.is_running():
if job.is_enabled():
await scheduler_service.add_job(job)
elif job.is_paused():
await scheduler_service.pause_job(job.code)
else:
await scheduler_service.remove_job(job.code)
await db.commit()
return SchedulerJobBatchUpdateStatusOut(count=count)
@router.post("/job/execute", response_model=SchedulerJobExecuteOut, summary="立即执行任务")
async def execute_scheduler_job(
data: SchedulerJobExecuteIn,
db: AsyncSession = Depends(get_db)
):
"""立即执行指定任务(不影响正常调度)"""
result = await db.execute(
select(SchedulerJob).where(SchedulerJob.id == data.job_id)
)
job = result.scalar_one_or_none()
if not job:
raise HTTPException(status_code=404, detail="任务不存在")
if not scheduler_service.is_running():
raise HTTPException(status_code=400, detail="调度器未运行")
# 立即执行任务
success = await scheduler_service.run_job_now(job.code)
if success:
return SchedulerJobExecuteOut(
success=True,
message=f"任务 {job.name} 将立即执行"
)
else:
return SchedulerJobExecuteOut(
success=False,
message=f"任务 {job.name} 执行失败,可能任务未在调度器中"
)
@router.post("/job/search", response_model=PaginatedResponse[SchedulerJobResponse], summary="搜索定时任务")
async def search_scheduler_jobs(
data: SchedulerJobSearchRequest,
page: int = Query(default=1, ge=1, description="页码"),
page_size: int = Query(default=settings.PAGE_SIZE, ge=1, le=settings.PAGE_MAX_SIZE, alias="pageSize", description="每页数量"),
db: AsyncSession = Depends(get_db)
):
"""搜索定时任务"""
keyword = data.keyword
filters = [
SchedulerJob.is_deleted == False, # noqa: E712
or_(
SchedulerJob.name.ilike(f"%{keyword}%"),
SchedulerJob.code.ilike(f"%{keyword}%"),
SchedulerJob.description.ilike(f"%{keyword}%"),
)
]
# 查询总数
count_result = await db.execute(
select(func.count(SchedulerJob.id)).where(*filters)
)
total = count_result.scalar()
# 查询数据
offset = (page - 1) * page_size
result = await db.execute(
select(SchedulerJob).where(*filters)
.order_by(SchedulerJob.priority.desc())
.offset(offset).limit(page_size)
)
jobs = result.scalars().all()
return PaginatedResponse(
items=[_build_job_response(job) for job in jobs],
total=total
)
@router.get("/job/statistics/data", response_model=SchedulerJobStatisticsOut, summary="获取任务统计信息")
async def get_scheduler_job_statistics(db: AsyncSession = Depends(get_db)):
"""获取任务统计信息"""
# 任务统计
total_result = await db.execute(
select(func.count(SchedulerJob.id)).where(SchedulerJob.is_deleted == False) # noqa: E712
)
total_jobs = total_result.scalar() or 0
enabled_result = await db.execute(
select(func.count(SchedulerJob.id)).where(
SchedulerJob.is_deleted == False, # noqa: E712
SchedulerJob.status == 1
)
)
enabled_jobs = enabled_result.scalar() or 0
disabled_result = await db.execute(
select(func.count(SchedulerJob.id)).where(
SchedulerJob.is_deleted == False, # noqa: E712
SchedulerJob.status == 0
)
)
disabled_jobs = disabled_result.scalar() or 0
paused_result = await db.execute(
select(func.count(SchedulerJob.id)).where(
SchedulerJob.is_deleted == False, # noqa: E712
SchedulerJob.status == 2
)
)
paused_jobs = paused_result.scalar() or 0
# 执行统计
total_exec_result = await db.execute(
select(func.count(SchedulerLog.id))
)
total_executions = total_exec_result.scalar() or 0
success_exec_result = await db.execute(
select(func.count(SchedulerLog.id)).where(SchedulerLog.status == 'success')
)
success_executions = success_exec_result.scalar() or 0
failed_exec_result = await db.execute(
select(func.count(SchedulerLog.id)).where(SchedulerLog.status == 'failed')
)
failed_executions = failed_exec_result.scalar() or 0
# 计算成功率
success_rate = round(success_executions / total_executions * 100, 2) if total_executions > 0 else 0
return SchedulerJobStatisticsOut(
total_jobs=total_jobs,
enabled_jobs=enabled_jobs,
disabled_jobs=disabled_jobs,
paused_jobs=paused_jobs,
total_executions=total_executions,
success_executions=success_executions,
failed_executions=failed_executions,
success_rate=success_rate,
)
@router.get("/job/{job_id}", response_model=SchedulerJobResponse, summary="获取定时任务详情")
async def get_scheduler_job(job_id: str, db: AsyncSession = Depends(get_db)):
"""获取单个定时任务的详细信息"""
result = await db.execute(
select(SchedulerJob).where(SchedulerJob.id == job_id)
)
job = result.scalar_one_or_none()
if not job:
raise HTTPException(status_code=404, detail="任务不存在")
return _build_job_response(job)
@router.put("/job/{job_id}", response_model=SchedulerJobResponse, summary="更新定时任务")
async def update_scheduler_job(
job_id: str,
data: SchedulerJobUpdate,
db: AsyncSession = Depends(get_db)
):
"""更新定时任务"""
result = await db.execute(
select(SchedulerJob).where(SchedulerJob.id == job_id)
)
job = result.scalar_one_or_none()
if not job:
raise HTTPException(status_code=404, detail="任务不存在")
# 检查任务编码是否已存在(排除自身)
if data.code:
code_result = await db.execute(
select(SchedulerJob).where(
SchedulerJob.code == data.code,
SchedulerJob.id != job_id,
SchedulerJob.is_deleted == False # noqa: E712
)
)
if code_result.scalar_one_or_none():
raise HTTPException(status_code=400, detail=f"任务编码已存在: {data.code}")
# 更新字段
update_data = data.model_dump(exclude_unset=True)
for key, value in update_data.items():
setattr(job, key, value)
await db.commit()
await db.refresh(job)
# 同步更新调度器
if scheduler_service.is_running():
await scheduler_service.modify_job(job)
return _build_job_response(job)
@router.delete("/job/{job_id}", response_model=ResponseModel, summary="删除定时任务")
async def delete_scheduler_job(
job_id: str,
hard: bool = Query(default=False, description="是否物理删除"),
db: AsyncSession = Depends(get_db)
):
"""删除定时任务"""
result = await db.execute(
select(SchedulerJob).where(SchedulerJob.id == job_id)
)
job = result.scalar_one_or_none()
if not job:
raise HTTPException(status_code=404, detail="任务不存在")
# 从调度器移除
if scheduler_service.is_running():
await scheduler_service.remove_job(job.code)
if hard:
await db.delete(job)
else:
job.is_deleted = True
await db.commit()
return ResponseModel(message="删除成功")
# ==================== SchedulerLog APIs ====================
@router.get("/log", response_model=PaginatedResponse[SchedulerLogResponse], summary="获取任务执行日志列表")
async def get_scheduler_log_list(
page: int = Query(default=1, ge=1, description="页码"),
page_size: int = Query(default=settings.PAGE_SIZE, ge=1, le=settings.PAGE_MAX_SIZE, alias="pageSize", description="每页数量"),
job_id: Optional[str] = Query(default=None, description="任务ID"),
job_code: Optional[str] = Query(default=None, description="任务编码"),
job_name: Optional[str] = Query(default=None, description="任务名称"),
status: Optional[str] = Query(default=None, description="执行状态"),
start_time_gte: Optional[datetime] = Query(default=None, alias="startTimeGte", description="开始时间>="),
start_time_lte: Optional[datetime] = Query(default=None, alias="startTimeLte", description="开始时间<="),
db: AsyncSession = Depends(get_db)
):
"""获取任务执行日志列表(分页)"""
filters = []
if job_id:
filters.append(SchedulerLog.job_id == job_id)
if job_code:
filters.append(SchedulerLog.job_code.ilike(f"%{job_code}%"))
if job_name:
filters.append(SchedulerLog.job_name.ilike(f"%{job_name}%"))
if status:
filters.append(SchedulerLog.status == status)
if start_time_gte:
filters.append(SchedulerLog.start_time >= start_time_gte)
if start_time_lte:
filters.append(SchedulerLog.start_time <= start_time_lte)
# 查询总数
count_result = await db.execute(
select(func.count(SchedulerLog.id)).where(*filters) if filters else select(func.count(SchedulerLog.id))
)
total = count_result.scalar()
# 查询数据
offset = (page - 1) * page_size
query = select(SchedulerLog).order_by(SchedulerLog.start_time.desc()).offset(offset).limit(page_size)
if filters:
query = select(SchedulerLog).where(*filters).order_by(SchedulerLog.start_time.desc()).offset(offset).limit(page_size)
result = await db.execute(query)
logs = result.scalars().all()
return PaginatedResponse(
items=[_build_log_response(log) for log in logs],
total=total
)
@router.get("/log/by/job/{job_id}", response_model=PaginatedResponse[SchedulerLogResponse], summary="获取指定任务的执行日志")
async def get_scheduler_logs_by_job(
job_id: str,
page: int = Query(default=1, ge=1, description="页码"),
page_size: int = Query(default=settings.PAGE_SIZE, ge=1, le=settings.PAGE_MAX_SIZE, alias="pageSize", description="每页数量"),
db: AsyncSession = Depends(get_db)
):
"""获取指定任务的所有执行日志"""
# 查询总数
count_result = await db.execute(
select(func.count(SchedulerLog.id)).where(SchedulerLog.job_id == job_id)
)
total = count_result.scalar()
# 查询数据
offset = (page - 1) * page_size
result = await db.execute(
select(SchedulerLog).where(SchedulerLog.job_id == job_id)
.order_by(SchedulerLog.start_time.desc())
.offset(offset).limit(page_size)
)
logs = result.scalars().all()
return PaginatedResponse(
items=[_build_log_response(log) for log in logs],
total=total
)
@router.get("/log/{log_id}", response_model=SchedulerLogResponse, summary="获取任务执行日志详情")
async def get_scheduler_log(log_id: str, db: AsyncSession = Depends(get_db)):
"""获取单个任务执行日志的详细信息"""
result = await db.execute(
select(SchedulerLog).where(SchedulerLog.id == log_id)
)
log = result.scalar_one_or_none()
if not log:
raise HTTPException(status_code=404, detail="日志不存在")
return _build_log_response(log)
@router.delete("/log/{log_id}", response_model=ResponseModel, summary="删除任务执行日志")
async def delete_scheduler_log(log_id: str, db: AsyncSession = Depends(get_db)):
"""删除任务执行日志"""
result = await db.execute(
select(SchedulerLog).where(SchedulerLog.id == log_id)
)
log = result.scalar_one_or_none()
if not log:
raise HTTPException(status_code=404, detail="日志不存在")
await db.delete(log)
await db.commit()
return ResponseModel(message="删除成功")
@router.post("/log/batch/delete", response_model=SchedulerLogBatchDeleteOut, summary="批量删除任务执行日志")
async def batch_delete_scheduler_logs(
data: SchedulerLogBatchDeleteIn,
db: AsyncSession = Depends(get_db)
):
"""批量删除任务执行日志"""
result = await db.execute(
select(SchedulerLog).where(SchedulerLog.id.in_(data.ids))
)
logs = result.scalars().all()
count = len(logs)
for log in logs:
await db.delete(log)
await db.commit()
return SchedulerLogBatchDeleteOut(count=count)
@router.post("/log/clean", response_model=SchedulerLogCleanOut, summary="清理旧日志")
async def clean_scheduler_logs(
data: SchedulerLogCleanIn,
db: AsyncSession = Depends(get_db)
):
"""清理旧日志"""
cutoff_date = datetime.now() - timedelta(days=data.days)
filters = [SchedulerLog.start_time < cutoff_date]
if data.status:
filters.append(SchedulerLog.status == data.status)
result = await db.execute(
select(SchedulerLog).where(*filters)
)
logs = result.scalars().all()
count = len(logs)
for log in logs:
await db.delete(log)
await db.commit()
return SchedulerLogCleanOut(count=count)
# ==================== Scheduler Control APIs ====================
@router.post("/start", response_model=ResponseModel, summary="启动调度器")
async def start_scheduler():
"""启动调度器(注意:调度器应在应用启动时自动启动)"""
if scheduler_service.is_running():
raise HTTPException(status_code=400, detail="调度器已在运行中")
# APScheduler 4.x 中调度器应在 lifespan 中启动
raise HTTPException(status_code=400, detail="请重启应用以启动调度器")
@router.post("/shutdown", response_model=ResponseModel, summary="关闭调度器")
async def shutdown_scheduler():
"""关闭调度器(注意:调度器会在应用关闭时自动关闭)"""
if not scheduler_service.is_running():
raise HTTPException(status_code=400, detail="调度器未运行")
# APScheduler 4.x 中调度器应在 lifespan 中关闭
raise HTTPException(status_code=400, detail="请关闭应用以停止调度器")
@router.post("/pause", response_model=ResponseModel, summary="暂停调度器")
async def pause_scheduler():
"""暂停调度器(暂不支持)"""
raise HTTPException(status_code=400, detail="APScheduler 4.x 暂不支持暂停整个调度器")
@router.post("/resume", response_model=ResponseModel, summary="恢复调度器")
async def resume_scheduler():
"""恢复调度器(暂不支持)"""
raise HTTPException(status_code=400, detail="APScheduler 4.x 暂不支持恢复整个调度器")
@router.get("/status", response_model=SchedulerStatusOut, summary="获取调度器状态")
async def get_scheduler_status():
"""获取调度器状态"""
is_running = scheduler_service.is_running()
jobs = await scheduler_service.get_all_jobs() if is_running else []
return SchedulerStatusOut(
is_running=is_running,
job_count=len(jobs),
jobs=jobs,
)
@router.get("/log/{log_id}/stream", summary="实时日志流")
async def stream_task_log(log_id: str):
"""
通过 SSE 实时推送任务执行日志
Args:
log_id: 日志记录 ID
Returns:
SSE 事件流,包含任务执行过程中的日志消息
"""
from fastapi.responses import StreamingResponse
from scheduler.task_log_service import TaskLogService
import json
async def event_generator():
try:
async for log_entry in TaskLogService.subscribe(log_id):
yield f"data: {json.dumps(log_entry, ensure_ascii=False)}\n\n"
except Exception as e:
yield f"data: {json.dumps({'level': 'error', 'message': str(e)}, ensure_ascii=False)}\n\n"
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no",
}
)