ThinkPHP项目邮件队列与延迟

wen PHP项目 5

本文目录导读:

ThinkPHP项目邮件队列与延迟

  1. 使用ThinkPHP自带的队列功能(推荐)
  2. 使用项目中的邮件发送方式配合队列
  3. 扩展:使用Redis延迟队列
  4. 消费队列任务
  5. 完整的邮件发送Service示例
  6. 注意事项

在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'
        ];
    }
}

注意事项

  1. 数据库方式:需创建 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;
  2. 配置文件:确保邮件和队列配置正确

  3. 错误处理:增加重试和日志记录

  4. 队列消费:使用 supervisor 或 cron 保证队列持续消费

这样就能在ThinkPHP项目中灵活实现邮件的异步发送和延迟发送了。

抱歉,评论功能暂时关闭!