
定时任务与异步处理机制
大约 9 分钟
定时任务与异步处理机制
想象一下,你的测试平台就像一个勤劳的"机器人管家",能够在你睡觉的时候自动跑回归测试,在你喝咖啡的时候处理大批量的用例执行。今天我们要给平台装上"自动驾驶"系统,让它真正实现无人值守!
🎯 为什么需要异步处理?
痛点分析:同步处理的"三大噩梦"
作为一个经常被测试执行时间折磨的工程师,我深知同步处理的痛苦:
1. 用户体验噩梦 😱
- 执行测试套件时,页面卡死等待
- 用户点击按钮后,半天没反应
- 浏览器超时,前功尽弃
2. 资源浪费严重 💸
- 服务器线程被长时间占用
- 数据库连接池耗尽
- 内存占用居高不下
3. 扩展性极差 📈
- 并发用户增加,系统崩溃
- 无法处理大批量任务
- 系统响应越来越慢
异步处理的优势:让系统"飞"起来
异步处理就像给系统装上了"涡轮增压器":
- 响应迅速:用户操作立即响应,体验丝滑
- 资源高效:充分利用系统资源,提升吞吐量
- 可扩展性强:轻松应对高并发场景
- 容错性好:任务失败可重试,系统更稳定
🏗️ 异步架构设计
整体架构:分布式任务处理系统
┌─────────────────────────────────────────────────────────────┐
│ 🌐 Web应用层 │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ 任务提交 │ │ 状态查询 │ │ 结果展示 │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
├─────────────────────────────────────────────────────────────┤
│ 🔧 任务调度层 │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ Celery Beat│ │ 任务路由 │ │ 优先级队列 │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
├─────────────────────────────────────────────────────────────┤
│ ⚙️ 任务执行层 │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │
│ │ Worker 1 │ │ Worker 2 │ │ Worker N │ │
│ └─────────────┘ └─────────────┘ └─────────────┘ │
├─────────────────────────────────────────────────────────────┤
│ 💾 消息队列 │
│ Redis / RabbitMQ │
└─────────────────────────────────────────────────────────────┘技术选型:工具链的智慧选择
任务队列:
- Celery:Python生态最成熟的分布式任务队列
- Redis:高性能的消息代理,简单易用
- Flower:Celery的Web监控工具,实时查看任务状态
定时调度:
- Celery Beat:Celery内置的定时任务调度器
- Crontab:支持标准的Cron表达式
- Django-celery-beat:数据库存储的定时任务(可选)
⚙️ Celery配置与集成
Celery基础配置:让任务队列跑起来
app/celery_app.py:
from celery import Celery
from app.config import Config
import os
# 创建Celery实例
celery_app = Celery('test_platform')
# 配置Celery
celery_app.conf.update(
# 消息代理配置
broker_url=Config.CELERY_BROKER_URL,
result_backend=Config.CELERY_RESULT_BACKEND,
# 任务序列化
task_serializer='json',
accept_content=['json'],
result_serializer='json',
timezone='Asia/Shanghai',
enable_utc=True,
# 任务路由配置
task_routes={
'app.tasks.test_execution.*': {'queue': 'test_execution'},
'app.tasks.report_generation.*': {'queue': 'report_generation'},
'app.tasks.data_processing.*': {'queue': 'data_processing'},
},
# 任务优先级
task_default_priority=5,
worker_prefetch_multiplier=1,
# 任务重试配置
task_acks_late=True,
task_reject_on_worker_lost=True,
# 结果过期时间
result_expires=3600,
# 定时任务配置
beat_schedule={
'daily-statistics': {
'task': 'app.tasks.data_processing.calculate_daily_statistics',
'schedule': 60.0 * 60.0 * 24, # 每天执行一次
'args': (),
},
'cleanup-old-results': {
'task': 'app.tasks.maintenance.cleanup_old_execution_results',
'schedule': 60.0 * 60.0 * 6, # 每6小时执行一次
'args': (),
},
},
)
# 自动发现任务
celery_app.autodiscover_tasks(['app.tasks'])
# 配置日志
celery_app.conf.worker_log_format = '[%(asctime)s: %(levelname)s/%(processName)s] %(message)s'
celery_app.conf.worker_task_log_format = '[%(asctime)s: %(levelname)s/%(processName)s][%(task_name)s(%(task_id)s)] %(message)s'任务定义:各司其职的工作者
app/tasks/test_execution.py:
from celery import current_task
from app.celery_app import celery_app
from app.models.testcase import TestCase, TestSuite
from app.models.execution import ExecutionRecord, ExecutionStep
from app.services.execution_engine import ExecutorFactory
from app.utils.logger import logger
import uuid
from datetime import datetime
import traceback
@celery_app.task(bind=True, name='app.tasks.test_execution.execute_single_testcase')
def execute_single_testcase(self, testcase_id: int, executor: str = 'system', environment: str = 'test'):
"""执行单个测试用例"""
task_id = self.request.id
logger.info(f"开始执行测试用例 {testcase_id}, 任务ID: {task_id}")
try:
# 获取测试用例
testcase = TestCase.get_by_id(testcase_id)
if testcase.is_deleted or testcase.status != 'active':
raise ValueError(f"测试用例 {testcase_id} 不可执行")
# 创建执行记录
execution_record = ExecutionRecord.create(
testcase=testcase,
execution_id=str(uuid.uuid4()),
trigger_type='api',
executor=executor,
status='running',
start_time=datetime.now(),
environment=environment
)
# 更新任务状态
self.update_state(
state='PROGRESS',
meta={
'current': 0,
'total': 1,
'status': '正在执行测试用例...',
'execution_id': execution_record.execution_id
}
)
# 创建执行器并执行
executor_instance = ExecutorFactory.create_executor(testcase, execution_record)
result = await executor_instance.execute()
# 更新执行记录
execution_record.status = 'success' if result['success'] else 'failed'
execution_record.end_time = datetime.now()
execution_record.duration = int((execution_record.end_time - execution_record.start_time).total_seconds())
execution_record.set_result_summary_dict(result)
if not result['success']:
execution_record.error_message = result.get('error', '执行失败')
execution_record.save()
# 更新用例统计
testcase.total_executions += 1
if result['success']:
testcase.success_executions += 1
testcase.save()
logger.info(f"测试用例 {testcase_id} 执行完成,结果: {'成功' if result['success'] else '失败'}")
return {
'success': result['success'],
'execution_id': execution_record.execution_id,
'duration': execution_record.duration,
'result': result
}
except Exception as e:
logger.error(f"执行测试用例 {testcase_id} 失败: {str(e)}")
logger.error(traceback.format_exc())
# 更新执行记录为失败状态
if 'execution_record' in locals():
execution_record.status = 'failed'
execution_record.end_time = datetime.now()
execution_record.error_message = str(e)
execution_record.save()
# 重新抛出异常,让Celery处理重试逻辑
raise
@celery_app.task(bind=True, name='app.tasks.test_execution.execute_test_suite')
def execute_test_suite(self, suite_id: int, executor: str = 'system', environment: str = 'test'):
"""执行测试套件"""
task_id = self.request.id
logger.info(f"开始执行测试套件 {suite_id}, 任务ID: {task_id}")
try:
# 获取测试套件
suite = TestSuite.get_by_id(suite_id)
if suite.is_deleted:
raise ValueError(f"测试套件 {suite_id} 不存在")
# 获取套件中的测试用例
suite_cases = suite.suite_cases.join(TestCase).where(
TestCase.status == 'active'
).order_by(suite.suite_cases.c.order)
total_cases = suite_cases.count()
if total_cases == 0:
raise ValueError("测试套件中没有可执行的用例")
# 创建套件执行记录
suite_execution = ExecutionRecord.create(
testsuite=suite,
execution_id=str(uuid.uuid4()),
trigger_type='api',
executor=executor,
status='running',
start_time=datetime.now(),
environment=environment
)
results = []
success_count = 0
# 执行每个测试用例
for index, suite_case in enumerate(suite_cases):
testcase = suite_case.testcase
# 更新任务进度
self.update_state(
state='PROGRESS',
meta={
'current': index,
'total': total_cases,
'status': f'正在执行用例: {testcase.name}',
'execution_id': suite_execution.execution_id
}
)
try:
# 执行单个用例
case_result = execute_single_testcase.apply_async(
args=[testcase.id, executor, environment]
).get(timeout=300) # 5分钟超时
results.append({
'testcase_id': testcase.id,
'testcase_name': testcase.name,
'success': case_result['success'],
'duration': case_result['duration'],
'execution_id': case_result['execution_id']
})
if case_result['success']:
success_count += 1
except Exception as e:
logger.error(f"执行用例 {testcase.name} 失败: {str(e)}")
results.append({
'testcase_id': testcase.id,
'testcase_name': testcase.name,
'success': False,
'error': str(e)
})
# 更新套件执行记录
suite_execution.status = 'success' if success_count == total_cases else 'failed'
suite_execution.end_time = datetime.now()
suite_execution.duration = int((suite_execution.end_time - suite_execution.start_time).total_seconds())
suite_execution.set_result_summary_dict({
'total_cases': total_cases,
'success_cases': success_count,
'failed_cases': total_cases - success_count,
'success_rate': round(success_count / total_cases * 100, 2),
'results': results
})
suite_execution.save()
logger.info(f"测试套件 {suite_id} 执行完成,成功率: {success_count}/{total_cases}")
return {
'success': success_count == total_cases,
'execution_id': suite_execution.execution_id,
'total_cases': total_cases,
'success_cases': success_count,
'failed_cases': total_cases - success_count,
'duration': suite_execution.duration,
'results': results
}
except Exception as e:
logger.error(f"执行测试套件 {suite_id} 失败: {str(e)}")
logger.error(traceback.format_exc())
# 更新套件执行记录为失败状态
if 'suite_execution' in locals():
suite_execution.status = 'failed'
suite_execution.end_time = datetime.now()
suite_execution.error_message = str(e)
suite_execution.save()
raise
@celery_app.task(name='app.tasks.test_execution.scheduled_execution')
def scheduled_execution(schedule_id: int):
"""定时执行任务"""
logger.info(f"开始执行定时任务 {schedule_id}")
try:
from app.models.schedule import TestSchedule
# 获取定时任务配置
schedule = TestSchedule.get_by_id(schedule_id)
if not schedule.is_active:
logger.warning(f"定时任务 {schedule_id} 已禁用")
return
# 根据任务类型执行
if schedule.target_type == 'testcase':
result = execute_single_testcase.delay(
schedule.target_id,
'scheduler',
schedule.environment
)
elif schedule.target_type == 'testsuite':
result = execute_test_suite.delay(
schedule.target_id,
'scheduler',
schedule.environment
)
else:
raise ValueError(f"不支持的任务类型: {schedule.target_type}")
# 更新最后执行时间
schedule.last_execution_time = datetime.now()
schedule.last_execution_task_id = result.id
schedule.save()
logger.info(f"定时任务 {schedule_id} 提交成功,任务ID: {result.id}")
return {
'success': True,
'task_id': result.id,
'schedule_id': schedule_id
}
except Exception as e:
logger.error(f"定时任务 {schedule_id} 执行失败: {str(e)}")
raise定时任务管理:让系统自己跑
app/models/schedule.py:
from peewee import *
from .base import BaseModel
from .project import Project
from .testcase import TestCase, TestSuite
from datetime import datetime
class TestSchedule(BaseModel):
"""测试定时任务模型"""
name = CharField(max_length=200, verbose_name="任务名称")
description = TextField(null=True, verbose_name="任务描述")
# 关联目标
target_type = CharField(
max_length=20,
choices=[('testcase', '测试用例'), ('testsuite', '测试套件')],
verbose_name="目标类型"
)
target_id = IntegerField(verbose_name="目标ID")
# 调度配置
cron_expression = CharField(max_length=100, verbose_name="Cron表达式")
timezone = CharField(max_length=50, default='Asia/Shanghai', verbose_name="时区")
# 执行配置
environment = CharField(max_length=50, default='test', verbose_name="执行环境")
max_retry_count = IntegerField(default=3, verbose_name="最大重试次数")
timeout = IntegerField(default=3600, verbose_name="超时时间(秒)")
# 状态管理
is_active = BooleanField(default=True, verbose_name="是否启用")
# 执行记录
last_execution_time = DateTimeField(null=True, verbose_name="最后执行时间")
last_execution_task_id = CharField(max_length=100, null=True, verbose_name="最后执行任务ID")
next_execution_time = DateTimeField(null=True, verbose_name="下次执行时间")
# 统计信息
total_executions = IntegerField(default=0, verbose_name="总执行次数")
success_executions = IntegerField(default=0, verbose_name="成功执行次数")
class Meta:
table_name = 'test_schedules'
@property
def success_rate(self):
"""成功率"""
if self.total_executions == 0:
return 0
return round(self.success_executions / self.total_executions * 100, 2)
def get_target_object(self):
"""获取目标对象"""
if self.target_type == 'testcase':
return TestCase.get_by_id(self.target_id)
elif self.target_type == 'testsuite':
return TestSuite.get_by_id(self.target_id)
else:
return None🎨 前端任务监控
任务状态监控组件:实时掌控执行状态
src/components/TaskMonitor/TaskMonitor.tsx:
import React, { useState, useEffect } from 'react';
import { Card, Table, Tag, Progress, Button, Space, Modal, message } from 'antd';
import {
PlayCircleOutlined,
PauseCircleOutlined,
StopOutlined,
ReloadOutlined,
EyeOutlined
} from '@ant-design/icons';
import { ColumnsType } from 'antd/es/table';
import { taskApi } from '../../services/task';
import './TaskMonitor.less';
interface TaskInfo {
id: string;
name: string;
status: 'pending' | 'running' | 'success' | 'failed' | 'cancelled';
progress: number;
start_time: string;
duration?: number;
result?: any;
error?: string;
}
const TaskMonitor: React.FC = () => {
const [tasks, setTasks] = useState<TaskInfo[]>([]);
const [loading, setLoading] = useState(false);
const [selectedTask, setSelectedTask] = useState<TaskInfo | null>(null);
const [detailVisible, setDetailVisible] = useState(false);
// 获取任务列表
const fetchTasks = async () => {
setLoading(true);
try {
const response = await taskApi.getTaskList();
setTasks(response.data.data);
} catch (error) {
message.error('获取任务列表失败');
} finally {
setLoading(false);
}
};
// 取消任务
const cancelTask = async (taskId: string) => {
try {
await taskApi.cancelTask(taskId);
message.success('任务已取消');
fetchTasks();
} catch (error) {
message.error('取消任务失败');
}
};
// 重试任务
const retryTask = async (taskId: string) => {
try {
await taskApi.retryTask(taskId);
message.success('任务已重新提交');
fetchTasks();
} catch (error) {
message.error('重试任务失败');
}
};
// 查看任务详情
const viewTaskDetail = (task: TaskInfo) => {
setSelectedTask(task);
setDetailVisible(true);
};
// 状态标签渲染
const renderStatus = (status: string) => {
const statusConfig = {
pending: { color: 'default', text: '等待中' },
running: { color: 'processing', text: '执行中' },
success: { color: 'success', text: '成功' },
failed: { color: 'error', text: '失败' },
cancelled: { color: 'warning', text: '已取消' },
};
const config = statusConfig[status] || statusConfig.pending;
return <Tag color={config.color}>{config.text}</Tag>;
};
// 进度条渲染
const renderProgress = (progress: number, status: string) => {
if (status === 'pending') {
return <Progress percent={0} size="small" />;
}
if (status === 'running') {
return <Progress percent={progress} size="small" status="active" />;
}
if (status === 'success') {
return <Progress percent={100} size="small" status="success" />;
}
if (status === 'failed') {
return <Progress percent={progress} size="small" status="exception" />;
}
return <Progress percent={progress} size="small" />;
};
// 表格列配置
const columns: ColumnsType<TaskInfo> = [
{
title: '任务名称',
dataIndex: 'name',
key: 'name',
ellipsis: true,
},
{
title: '状态',
dataIndex: 'status',
key: 'status',
render: renderStatus,
},
{
title: '进度',
key: 'progress',
render: (_, record) => renderProgress(record.progress, record.status),
},
{
title: '开始时间',
dataIndex: 'start_time',
key: 'start_time',
render: (text) => new Date(text).toLocaleString(),
},
{
title: '执行时长',
dataIndex: 'duration',
key: 'duration',
render: (duration) => duration ? `${duration}秒` : '-',
},
{
title: '操作',
key: 'action',
render: (_, record) => (
<Space size="middle">
<Button
type="text"
icon={<EyeOutlined />}
onClick={() => viewTaskDetail(record)}
>
详情
</Button>
{record.status === 'running' && (
<Button
type="text"
danger
icon={<StopOutlined />}
onClick={() => cancelTask(record.id)}
>
取消
</Button>
)}
{record.status === 'failed' && (
<Button
type="text"
icon={<ReloadOutlined />}
onClick={() => retryTask(record.id)}
>
重试
</Button>
)}
</Space>
),
},
];
// 定时刷新
useEffect(() => {
fetchTasks();
const interval = setInterval(fetchTasks, 5000); // 每5秒刷新一次
return () => clearInterval(interval);
}, []);
return (
<div className="task-monitor">
<Card
title="任务监控"
extra={
<Button
icon={<ReloadOutlined />}
onClick={fetchTasks}
loading={loading}
>
刷新
</Button>
}
>
<Table
columns={columns}
dataSource={tasks}
loading={loading}
rowKey="id"
pagination={{
showSizeChanger: true,
showQuickJumper: true,
showTotal: (total) => `共 ${total} 个任务`,
}}
/>
</Card>
{/* 任务详情模态框 */}
<Modal
title="任务详情"
open={detailVisible}
onCancel={() => setDetailVisible(false)}
footer={null}
width={800}
>
{selectedTask && (
<div className="task-detail">
<div className="detail-item">
<span className="label">任务ID:</span>
<span className="value">{selectedTask.id}</span>
</div>
<div className="detail-item">
<span className="label">任务名称:</span>
<span className="value">{selectedTask.name}</span>
</div>
<div className="detail-item">
<span className="label">状态:</span>
<span className="value">{renderStatus(selectedTask.status)}</span>
</div>
<div className="detail-item">
<span className="label">进度:</span>
<span className="value">
{renderProgress(selectedTask.progress, selectedTask.status)}
</span>
</div>
{selectedTask.error && (
<div className="detail-item">
<span className="label">错误信息:</span>
<pre className="error-message">{selectedTask.error}</pre>
</div>
)}
{selectedTask.result && (
<div className="detail-item">
<span className="label">执行结果:</span>
<pre className="result-content">
{JSON.stringify(selectedTask.result, null, 2)}
</pre>
</div>
)}
</div>
)}
</Modal>
</div>
);
};
export default TaskMonitor;🎯 下一步预告
今天我们为测试平台装上了强大的"自动驾驶"系统,实现了:
- 完整的异步任务处理架构
- 灵活的定时任务调度机制
- 实时的任务监控界面
- 可靠的错误处理和重试机制
下一篇我们将进入AI时代,学习如何集成AutoGen和RAG技术,打造智能化的测试平台。让AI成为你的测试助手,自动生成测试用例、分析测试结果!
💡 异步处理小贴士:异步处理就像多线程的生活,要学会合理分配任务,避免资源竞争。记住,并发不是银弹,设计合理的任务队列比盲目追求高并发更重要!
