【消息队列学习一】TP6 基于 redis 实现消息队列和延迟队列

前言

本文中主要记录TP6 中使用 think-queue 来实现redis的消息队列和延迟队列的过程以及其中出现的问题

think-queue:是thinkphp 官方提供的一个消息队列服务,它支持消息队列的一些基本特性:

  • 消息的发布,获取,执行,删除,重发,失败处理,延迟执行,超时控制等
  • 队列的多队列, 内存限制 ,启动,停止,守护等
  • 消息队列可降级为同步执行

环境准备(以下是本人的环境)

  • WAMP(win10主要是为了方便本地测试使用)
  • PHP 7.3.5  thinkphp 6.0.5
  • mysql 5.7.26
  • apache 2.4.39
  • Centos7 redis 6.0.5(部署在线上)

 

安装

在此我就不再略过TP6的项目创建过程了,大致就是安装composer工具,安装成功以后,再使用composer去创建项目即可。

think-queue 安装

composer require topthink/think-queue

我所安装的是目前最新的 3.0 版本

# 安装需要在项目的根目录下

安装完成以后可以在项目根目录下 vendor > topthink > think-queue

项目中添加驱动配置

我们需要在安装好的 think-queue > src 下找到 config.php 复制里面的内容,然后在根目录下 config 目录下创建 queue.php文件,将复制的内容粘贴进去

<?php
/**
 * 消息队列配置
 */


return [
    'default'     => 'redis',
    'connections' => [
        'sync'     => [
            'type' => 'sync',
        ],
        'database' => [
            'type'       => 'database',
            'queue'      => 'default',
            'table'      => 'jobs',
            'connection' => null,
        ],
        'redis'    => [
            'type'       => 'redis',
            'queue'      => 'default',
            'host'       => env('redis.host', '127.0.0.1'),
            'port'       => env('redis.port', 6379),
            'password'   => env('redis.password', ''),
            'select'     => 0,
            'timeout'    => 0,
            'persistent' => false,
        ],
    ],
    'failed'      => [
        'type'  => 'none',
        'table' => 'failed_jobs',
    ],
];

# env 指的就是我在外部的资源文件中设置了值

除此之外还需要安装redis扩展,没有安装的就自行去下载安装即可。

消息队列实现过程流程图

1、通过生产者推送消息到消息队列服务中

2、消息队列服务将收到的消息存入redis队列中(zset)

3、消费者进行监听队列,当监听到队列有新的消息时,获取队列第一条

4、处理获取下来的消息调用业务类进行处理相关业务

5、业务处理后,需要从队列中删除消息

功能实现

创建一个生产者

<?php
namespace appapicontroller;

use appBaseController;
use thinkfacadeQueue;

class Index extends BaseController
{
    public function index()
    {
        // echo phpinfo();exit();
        // 1.当前任务由哪个类来负责处理
        // 当轮到该任务时,系统将生成该类的实例,并调用其fire方法
        $jobHandlerClassName = 'appapicontrollerJob1';

        // 2.当任务归属的队列名称,如果为新队列,会自动创建
        $jobQueueName = "helloJobQueue";

        // 3.当前任务所需业务数据,不能为resource类型,其他类型最终将转化为json形式的字符串
        $jobData = ['ts' => time(), 'bizId' => uniqid(), 'a' => 1];

        // 4.将该任务推送到消息列表,等待对应的消费者去执行
        // 入队列,later延迟执行,单位秒,push立即执行
        $isPushed = Queue::later(10, $jobHandlerClassName, $jobData, $jobQueueName);

        // database 驱动时,返回值为 1|false  ;   redis 驱动时,返回值为 随机字符串|false
        if ($isPushed !== false) {
            echo '推送成功';
        } else {
            echo '推送失败';
        }
    }
}

创建一个消费者

<?php
namespace appapicontroller;

use thinkfacadeLog;
use thinkqueueJob;

class Job1
{
    /**
     * fire方法是消息队列默认调用的方法
     * @param Job $job 当前的任务对象
     * @param array $data 发布任务时自定义的数据
     */
    public function fire(Job $job, array $data)
    {
        // 有些任务在到达消费者时,可能已经不再需要执行了
        $isJobStillNeedToBeDone = $this->checkDatabaseToSeeIfJobNeedToBeDone($data);
        if (!$isJobStillNeedToBeDone) {
            $job->delete();
            return;
        }


        $isJobDone = $this->doHelloJob($data);
        if ($isJobDone){
            $job->delete();
            echo "删除任务" . $job->attempts() . '
';
        }else{
            if ($job->attempts() > 3){
                $job->delete();
                echo "超时任务删除" . $job->attempts() . '
';
            }
        }

    }

    /**
     * 有些消息在到达消费者时,可能已经不再需要执行了
     * @param array $data
     * @return bool
     */
    private function checkDatabaseToSeeIfJobNeedToBeDone(array $data) {
        return true;
    }

    /**
     * 根据消息中的数据进行实际的业务处理...
     * @param array $data
     * @return bool
     */
    private function doHelloJob(array $data)
    {
        echo '执行业务逻辑:' . $data['bizId'] . '
';

        return true;
    }
}

通过浏览器访问

 # 我这里是在本地配置了一个域名解析,实际访问的是127.0.0.1

访问后,可以看到已经向消息队列服务推送了消息,此时我们需要在项目根目录下运行命令创建工作进程来处理队列中的消息

 通过上述可以看到,当我们开启了work进程时,就会从队列中获取任务,然后找到消费者执行后续的业务逻辑。

因为这里我采用的push 表示立即执行,所以只要队列中有就会立马执行,如果我们需要使用到延时场景,例如订单支付超时,这时我们就可以使用later即可

多模块多功能实现

修改生产者代码

<?php
namespace appapicontroller;

use appBaseController;
use thinkfacadeQueue;

class Index extends BaseController
{
    public function index()
    {
        // echo phpinfo();exit();
        // 1.当前任务由哪个类来负责处理
        // 当轮到该任务时,系统将生成该类的实例,并调用其fire方法
        $jobHandlerClassName = 'appapicontrollerJob1';

        // 2.当任务归属的队列名称,如果为新队列,会自动创建
        $jobQueueName = "helloJobQueue";

        // 3.当前任务所需业务数据,不能为resource类型,其他类型最终将转化为json形式的字符串
        $jobData = ['ts' => time(), 'bizId' => uniqid(), 'a' => 1];

        // 4.将该任务推送到消息列表,等待对应的消费者去执行
        // 入队列,later延迟发送,单位秒,push立即发送
        $isPushed = Queue::later(10, $jobHandlerClassName, $jobData, $jobQueueName);

        // database 驱动时,返回值为 1|false  ;   redis 驱动时,返回值为 随机字符串|false
        if ($isPushed !== false) {
            echo '推送成功';
        } else {
            echo '推送失败';
        }
    }

    /**
     * 多模块延迟队列实现
     */
    public function pay(){
        $orderData = [
            "orderId" => uniqid()
        ];
        $isPushed = Queue::later(60, "appapicontrollerPayMessage", json_encode($orderData), "helloJobQueue");
        if ($isPushed)echo "
 订单支付成功 
";

        $email = [
            "email" => "1234567890@qq.com"
        ];
        $isPushed = Queue::later(120, "appapicontrollerEmailMessage", json_encode($email), "helloJobQueue");
        if ($isPushed)echo "
 邮件发送成功 
";
    }
}

新增支付消息消费者

<?php

namespace appapicontroller;


use thinkqueueJob;

class PayMessage
{
    public function fire(Job $job, $data){
        $data = json_decode($data, true);
        if ($this->doJob($data)){
            $job->delete();
        }else{
            if ($job->attempts() > 3){
                print_r("订单超时:" . $data['orderId']);
                $job->delete();
            }
        }
    }

    public function doJob($data){
        print_r("发送支付成功通知:" . $data['orderId'] );
        return true;
    }
}

新增邮箱发送消费者

<?php

namespace appapicontroller;


use thinkqueueJob;

class EmailMessage
{
    public function fire(Job $job, $data){
        $data = json_decode($data, true);
        if ($this->doJob($data)){
            $job->delete();
        }else{
            if ($job->attempts() > 3){
                print_r("
 邮件发送超时:" . $data['orderId'] . '
 ');
                $job->delete();
            }
        }
    }

    public function doJob($data){
        print_r("
 发送邮件:" . $data['email'] .'
 ');
        return true;
    }
}

通过浏览器模拟访问

 

 因为本次我们使用的是延时队列所以我们可以到redis中查看

127.0.0.1:6379> keys *
1) "{queues:helloJobQueue}:delayed"

当延时时间到了后,我们可以继续看到工作进程及时的进行消费

参考地址

think-queue官网文档:https://github.com/tp5er/think-queue/tree/master/doc

redis参考地址:https://www.runoob.com/redis/redis-tutorial.html

 

原文地址:https://www.cnblogs.com/lxd-ld/p/14012029.html