本文目录导读:

在ThinkPHP项目中实现邮件队列与延迟发送,主要有以下几种方案:
使用ThinkPHP自带的队列功能(推荐)
配置队列驱动
// config/queue.php
return [
'default' => 'redis', // 或 'database', 'sync'
'connections' => [
'sync' => [
'type' => 'sync',
],
'database' => [
'type' => 'database',
'queue' => 'default',
'table' => 'jobs',
'connection' => null,
],
'redis' => [
'type' => 'redis',
'queue' => 'default',
'host' => '127.0.0.1',
'port' => 6379,
'password' => '',
'select' => 0,
'timeout' => 0,
'persistent' => false,
],
],
'failed' => [
'type' => 'none',
'table' => 'failed_jobs',
],
];
创建邮件任务类
namespace app\job;
use think\queue\Job;
use think\facade\Mail;
class SendEmailJob
{
/**
* 执行任务
* @param Job $job 任务对象
* @param array $data 任务数据
*/
public function fire(Job $job, $data)
{
try {
// 发送邮件
Mail::send($data['template'], $data['data'], function ($message) use ($data) {
$message->to($data['to'])->subject($data['subject']);
// 如果有附件
if (isset($data['attachments'])) {
foreach ($data['attachments'] as $attachment) {
$message->attach($attachment);
}
}
});
// 删除任务
$job->delete();
} catch (\Exception $e) {
// 记录异常,可重试
if ($job->attempts() > 3) {
// 超过3次尝试,删除任务
$job->delete();
Log::error('邮件发送失败: ' . $e->getMessage());
} else {
// 延迟重试(10秒后)
$job->release(10);
}
}
}
}
使用队列发送邮件
use think\facade\Queue;
use think\facade\Mail;
class EmailService
{
/**
* 即时发送邮件
*/
public function sendNow($to, $subject, $content)
{
// 直接推送任务到队列
Queue::push(SendEmailJob::class, [
'to' => $to,
'subject' => $subject,
'template' => 'email/default',
'data' => ['content' => $content]
]);
}
/**
* 延迟发送邮件
*/
public function sendLater($to, $subject, $content, $delay = 300)
{
// 延迟5分钟发送
Queue::later($delay, SendEmailJob::class, [
'to' => $to,
'subject' => $subject,
'template' => 'email/default',
'data' => ['content' => $content]
]);
}
/**
* 指定队列名发送
*/
public function sendToSpecificQueue($to, $subject, $content, $connection = 'redis')
{
Queue::connection($connection)->push(SendEmailJob::class, [
'to' => $to,
'subject' => $subject,
], 'mail_queue');
}
}
使用项目中的邮件发送方式配合队列
namespace app\service;
use think\facade\Queue;
use app\job\MailJob;
class MailService
{
protected $mail;
public function __construct()
{
// 初始化邮件配置
}
/**
* 同步发送邮件
*/
public function sendSync($to, $subject, $content)
{
// 即时发送,等待结果
Mail::send('email.template', ['content' => $content], function($message) use ($to, $subject) {
$message->to($to)->subject($subject);
});
}
/**
* 异步发送邮件(加入队列)
*/
public function sendAsync($to, $subject, $content)
{
// 推送到队列,立即返回
Queue::push(function($job) use ($to, $subject, $content) {
Mail::send('email.template', ['content' => $content], function($message) use ($to, $subject) {
$message->to($to)->subject($subject);
});
$job->delete();
}, null, 'mail');
}
/**
* 延迟+批量发送
*/
public function sendBulk(array $recipients, $subject, $content, $delay = 3600)
{
foreach ($recipients as $recipient) {
Queue::later($delay, MailJob::class, [
'to' => $recipient['email'],
'name' => $recipient['name'],
'subject' => $subject,
'content' => $content
], 'mail');
}
}
}
扩展:使用Redis延迟队列
namespace app\service;
use think\facade\Cache;
class RedisMailQueue
{
/**
* 推送到延迟队列
*/
public function pushDelayed($email, $subject, $content, $delay = 3600)
{
$key = 'mail:delayed';
$data = [
'email' => $email,
'subject' => $subject,
'content' => $content,
'created_at' => time(),
'execute_at' => time() + $delay
];
return Cache::store('redis')->zAdd($key, time() + $delay, json_encode($data));
}
/**
* 消费延迟队列
*/
public function consumeDelayed()
{
$key = 'mail:delayed';
$now = time();
// 获取到期的任务
$tasks = Cache::store('redis')->zRangeByScore($key, 0, $now);
foreach ($tasks as $task) {
$data = json_decode($task, true);
// 发送邮件
$this->sendEmail($data);
// 从队列移除
Cache::store('redis')->zRem($key, $task);
}
}
private function sendEmail($data)
{
// 发送邮件逻辑
}
}
消费队列任务
使用命令行消费
# 处理默认队列 php think queue:work # 指定队列名和连接 php think queue:work --queue mail --connection redis # 指定消费次数 php think queue:work --once # 指定排队超时时间 php think queue:work --timeout 60
守护进程消费(Linux)
# 使用 Supervisor 管理 [program:think-queue] process_name=%(program_name)s_%(process_num)02d command=php /path/to/project/think queue:work --queue mail --tries 3 --timeout 60 autostart=true autorestart=true user=www-data numprocs=1 redirect_stderr=true stdout_logfile=/var/log/think/queue.log
完整的邮件发送Service示例
namespace app\service;
use think\facade\Queue;
use think\facade\Log;
use app\job\SendEmailJob;
use think\facade\Config;
class MailSendService
{
/**
* 发送邮件到队列
* @param string $to 收件人
* @param string $subject 主题
* @param string $content 内容
* @param int $delay 延迟秒数
* @param string $queue 队列名称
*/
public function queueSend($to, $subject, $content, $delay = 0, $queue = 'mail')
{
$data = [
'to' => $to,
'subject' => $subject,
'content' => $content
];
if ($delay > 0) {
// 延迟任务
return Queue::later($delay, SendEmailJob::class, $data, $queue);
} else {
// 即时的异步任务
return Queue::push(SendEmailJob::class, $data, $queue);
}
}
/**
* 批量发送
*/
public function batchQueueSend(array $emails, $subject, $content, $delay = 0)
{
$result = [];
foreach ($emails as $email) {
$result[] = $this->queueSend($email, $subject, $content, $delay);
}
return $result;
}
/**
* 检查邮件队列状态
*/
public function getQueueStatus($queue = 'mail')
{
// 这里可以查询队列长度等
return [
'queue' => $queue,
'status' => 'active'
];
}
}
注意事项
-
数据库方式:需创建
jobs表CREATE TABLE `jobs` ( `id` int(11) NOT NULL AUTO_INCREMENT, `queue` varchar(255) NOT NULL, `payload` longtext NOT NULL, `attempts` tinyint(3) unsigned NOT NULL, `reserved` tinyint(3) unsigned NOT NULL, `reserved_at` int(10) unsigned DEFAULT NULL, `available_at` int(10) unsigned NOT NULL, `created_at` int(10) unsigned NOT NULL, PRIMARY KEY (`id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
-
配置文件:确保邮件和队列配置正确
-
错误处理:增加重试和日志记录
-
队列消费:使用 supervisor 或 cron 保证队列持续消费
这样就能在ThinkPHP项目中灵活实现邮件的异步发送和延迟发送了。