CRMEB 系统采用基于 Redis 的消息队列机制,通过 ThinkPHP 框架的队列组件实现异步任务处理。队列系统在提高系统性能、优化用户体验、保证数据一致性方面发挥着重要作用。
return [
'default' => 'redis', // 默认驱动
'prefix' => 'crmeb_', // 队列前缀
'connections' => [
'redis' => [
'driver' => 'redis',
'queue' => 'CRMEB', // 队列名称
'host' => '127.0.0.1',
'port' => 6379,
'password' => '',
'select' => 0,
],
],
'failed' => [
'type' => 'database',
'table' => 'failed_jobs', // 失败任务表
],
];
abstract class BaseJobs implements JobInterface
{
// 任务执行入口
public function fire(Job $job, $data): void
// 任务重试机制
protected function runJob(string $action, Job $job, array $infoData, int $errorCount = 3)
}
trait QueueTrait
{
// 立即推送任务
public static function dispatch($action, array $data = [], string $queueName = null)
// 延迟推送任务
public static function dispatchSecs(int $secs, $action, array $data = [], string $queueName = null)
}
功能: 订单支付成功后的后续处理 触发时机: 订单支付成功时 主要处理:
// 使用示例
OrderJob::dispatch('doJob', [$orderInfo]);
功能: 订单创建后的数据处理 触发时机: 订单创建完成后 主要处理:
功能: 自动取消超时未支付订单 触发时机: 定时任务或延迟队列 处理逻辑:
// 延迟30分钟后执行
UnpaidOrderCancelJob::dispatchSecs(1800, 'doJob', [$orderId]);
功能: 记录商品相关行为日志 触发时机: 商品浏览、加购、下单等 日志类型:
visit: 商品访问cart: 加入购物车order: 下单pay: 支付collect: 收藏refund: 退款ProductLogJob::dispatch('doJob', ['visit', $productData]);
功能: 商品库存计算和更新 触发时机: 商品操作后 处理功能:
功能: 异步复制商品信息 触发时机: 商品复制操作 处理内容:
功能: 异步发送短信通知 触发时机: 各种业务事件 支持场景:
SmsJob::dispatch('doJob', [$phone, $data, $template]);
功能: 发送微信模板消息 触发时机: 业务状态变更 消息类型:
功能: 多端消息同步 触发时机: 消息发送时 同步范围:
功能: 拼团活动处理 触发时机: 拼团相关操作 处理内容:
功能: 直播相关处理 触发时机: 直播事件 处理内容:
功能: 自动生成商品评论 触发时机: 订单完成后 处理逻辑:
功能: 系统定时维护 执行频率: 每日定时 维护内容:
// 清理昨日海报
TaskJob::dispatch('emptyYesterdayAttachment');
功能: 系统升级处理 触发时机: 手动触发或定时 处理内容:
功能: 监控队列健康状态 执行频率: 定时检查 检查项目:
功能: 用户相关异步处理 触发时机: 用户操作 处理内容:
功能: 代理商相关处理 触发时机: 代理业务 处理内容:
功能: 物流信息处理 触发时机: 物流状态变更 处理内容:
功能: 订单接收处理 触发时机: 订单创建 处理内容:
功能: 发票处理 触发时机: 发票申请 处理内容:
功能: 数据外部推送 触发时机: 数据变更 推送内容:
基于 Workerman 实现的定时任务调度器
# 启动定时任务
php think timer start
# 守护进程方式启动
php think timer start -d
# 停止定时任务
php think timer stop
# 重启定时任务
php think timer reload
数据库表: system_crontab
字段说明:
name: 任务名称task_type: 任务类型mark: 任务标识content: 任务内容max_execution_time: 最大执行时间execution_cycle: 执行周期is_open: 是否开启next_execution_time: 下次执行时间last_execution_time: 上次执行时间// 检查队列状态
Queue::instance()->getQueueInfo();
// 获取队列长度
Queue::instance()->getQueueLength();
// 获取失败任务
Queue::instance()->getFailedJobs();
// 重新执行失败任务
Queue::instance()->retryFailed($jobId);
// 清空队列
Queue::instance()->clearQueue($queueName);
// 暂停队列
Queue::instance()->pauseQueue($queueName);
所有队列任务都有详细的日志记录:
根据业务重要性分组:
// 高优先级队列
'order_pay' => 'CRMEB_ORDER_PAY'
'sms' => 'CRMEB_SMS'
'email' => 'CRMEB_EMAIL'
// 普通优先级队列
'statistics' => 'CRMEB_STATISTICS'
'log' => 'CRMEB_LOG'
// 低优先级队列
'report' => 'CRMEB_REPORT'
'cleanup' => 'CRMEB_CLEANUP'
// 不同任务的重试次数
'order_pay' => 3 // 支付任务重试3次
'sms' => 5 // 短信任务重试5次
'email' => 3 // 邮件任务重试3次
'statistics' => 1 // 统计任务重试1次
// 常见延迟时间
UnpaidOrderCancelJob::dispatchSecs(1800, 'doJob', [$orderId]); // 30分钟
OrderRemindJob::dispatchSecs(86400, 'doJob', [$orderId]); // 24小时
CleanupJob::dispatchSecs(604800, 'doJob', [$data]); // 7天
现象: 队列长度持续增长 原因: 消费速度小于生产速度 解决:
现象: 同一任务被执行多次 原因: 任务重复入队或幂等性问题 解决:
现象: 进程内存持续增长 原因: 任务中存在资源未释放 解决:
# 暂停所有队列
php think queue:pause
# 暂停指定队列
php think queue:pause --queue=order_pay
# 清空指定队列
php think queue:clear --queue=statistics
# 清空所有队列
php think queue:clear --all
# 重建队列表
php think queue:table
# 重新发布失败任务
php think queue:retry all
<?php
namespace app\jobs;
use crmeb\basic\BaseJobs;
use crmeb\traits\QueueTrait;
class CustomJob extends BaseJobs
{
use QueueTrait;
/**
* 执行任务
* @param array $data
* @return bool
*/
public function doJob(array $data): bool
{
try {
// 业务逻辑处理
$this->processData($data);
return true;
} catch (\Exception $e) {
// 记录错误日志
Log::error('任务执行失败: ' . $e->getMessage());
return false;
}
}
private function processData(array $data): void
{
// 具体业务逻辑
}
}
// 立即执行
CustomJob::dispatch('doJob', [$data]);
// 延迟执行
CustomJob::dispatchSecs(300, 'doJob', [$data]);
// 指定队列
CustomJob::dispatch('doJob', [$data], 'custom_queue');
public function testCustomJob()
{
$job = new CustomJob();
$result = $job->doJob($testData);
$this->assertTrue($result);
}
public function testQueueDispatch()
{
// 模拟队列调度
Queue::fake();
CustomJob::dispatch('doJob', [$testData]);
Queue::assertPushed(CustomJob::class, function ($job) use ($testData) {
return $job->data[0] === $testData;
});
}
public function doJob(array $data): bool
{
Log::info('任务开始执行', ['data' => $data]);
try {
// 业务逻辑
Log::info('任务执行成功');
return true;
} catch (\Exception $e) {
Log::error('任务执行失败', [
'error' => $e->getMessage(),
'trace' => $e->getTraceAsString()
]);
return false;
}
}
public function doJob(array $data): bool
{
$startTime = microtime(true);
try {
// 业务逻辑
$endTime = microtime(true);
$executeTime = ($endTime - $startTime) * 1000; // 毫秒
Log::info('任务执行完成', [
'execute_time' => $executeTime . 'ms',
'memory_usage' => memory_get_usage(true)
]);
return true;
} catch (\Exception $e) {
Log::error('任务执行失败', ['error' => $e->getMessage()]);
return false;
}
}
CRMEB 的队列任务系统是系统高性能运行的关键组件,通过合理的任务分类、优先级管理、错误处理和监控机制,确保了系统的稳定性和可扩展性。在实际使用中,需要根据业务特点合理设计任务粒度,优化执行效率,并建立完善的监控告警机制。
队列系统的正确使用可以显著提升用户体验,降低系统负载,提高数据处理的一致性。建议开发团队充分理解队列机制,在合适场景下合理应用,并持续优化和改进队列性能。
提示:该文档由AI生成,仅供参考。