下面由laravel教程栏目给大家介绍laravel-swoole消息队列,希望对需要的朋友有所帮助!
这段时间用laravel8+laravel-swoole做项目,可发现laravel-swoole的扩展不兼容消息队列;
思来想去这咋办呢,这咋办呢.咋办那就自己写咯!还好thinkphp-swoole扩展已经兼容了,那不就嘿嘿嘿!
直接上修改的思路和代码!开干!
一种是增加另外启动的命令或者在swoole启动的时候一起启动消息队列进行消费,我这么懒的人一个命令能解决的,绝不写两命令.
首先重写swoole启动命令
<?phpnamespace crmeb\swoole\command;use illuminate\support\arr;use swoole\process;use swooletw\http\server\facades\server;use swooletw\http\server\manager;use crmeb\swoole\server\interactswithqueue;use crmeb\swoole\server\filewatcher;use swoole\runtime;class httpservercommand extends \swooletw\http\commands\httpservercommand{ use interactswithqueue; /** * the name and signature of the console command. * * @var string */ protected $signature = 'crmeb:http {action : start|stop|restart|reload|infos}'; /** * run swoole_http_server. */ protected function start() { if ($this->isrunning()) { $this->error('failed! swoole_http_server process is already running.'); return; } $host = arr::get($this->config, 'server.host'); $port = arr::get($this->config, 'server.port'); $hotreloadenabled = arr::get($this->config, 'hot_reload.enabled'); $queueenabled = arr::get($this->config, 'queue.enabled'); $accesslogenabled = arr::get($this->config, 'server.access_log'); $coroutineenable = arr::get($this->config, 'coroutine.enable'); $this->info('starting swoole http server...'); $this->info("swoole http server started: <http://{$host}:{$port}>"); if ($this->isdaemon()) { $this->info( '> (you can run this command to ensure the ' . 'swoole_http_server process is running: ps aux|grep "swoole")' ); } $manager = $this->laravel->make(manager::class); $server = $this->laravel->make(server::class); if ($accesslogenabled) { $this->registeraccesslog(); } //热更新重写 if ($hotreloadenabled) { $manager->addprocess($this->gethotreloadprocessnow($server)); } //启动消息队列进行消费 if ($queueenabled) { $this->preparequeue($manager); } if ($coroutineenable) { runtime::enablecoroutine(true, arr::get($this->config, 'coroutine.flags', swoole_hook_all)); } $manager->run(); } /** * @param server $server * @return process|void */ protected function gethotreloadprocessnow($server) { return new process(function () use ($server) { $watcher = new filewatcher( arr::get($this->config, 'hot_reload.include', []), arr::get($this->config, 'hot_reload.exclude', []), arr::get($this->config, 'hot_reload.name', []) ); $watcher->watch(function () use ($server) { $server->reload(); }); }, false, 0, true); }}
interactswithqueue 类
<?phpnamespace crmeb\swoole\server;use crmeb\swoole\queue\manager as queuemanager;use swooletw\http\server\manager;/** * trait interactswithqueue * @package crmeb\swoole\server */trait interactswithqueue{ public function preparequeue(manager $manager) { /** @var queuemanager $queuemanager */ $queuemanager = $this->laravel->make(queuemanager::class); $queuemanager->attachtoserver($manager, $this->output); }}
manager类
<?phpnamespace crmeb\swoole\queue;use illuminate\contracts\container\container;use swoole\constant;use swoole\process;use swoole\process\pool;use swoole\timer;use illuminate\support\arr;use illuminate\queue\events\jobfailed;use illuminate\queue\worker;use crmeb\swoole\server\withcontainer;use illuminate\queue\jobs\job;use function swoole\coroutine\run;use illuminate\queue\workeroptions;use swooletw\http\server\manager as servermanager;use illuminate\console\outputstyle;class manager{ use withcontainer; /** * container. * * @var \illuminate\contracts\container\container */ protected $container; /** * @var outputstyle */ protected $output; /** * @var closure[] */ protected $workers = []; /** * manager constructor. * @param container $container */ public function __construct(container $container) { $this->container = $container; } /** * @param servermanager $server */ public function attachtoserver(servermanager $server, outputstyle $output) { $this->output = $output; $this->listenforevents(); $this->createworkers(); foreach ($this->workers as $worker) { $server->addprocess(new process($worker, false, 0, true)); } } /** * 运行消息队列命令 */ public function run(): void { @cli_set_process_title("swoole queue: manager process"); $this->listenforevents(); $this->createworkers(); $pool = new pool(count($this->workers)); $pool->on(constant::event_worker_start, function (pool $pool, int $workerid) { $process = $pool->getprocess($workerid); run($this->workers[$workerid], $process); }); $pool->start(); } /** * 创建执行任务 */ protected function createworkers() { $workers = $this->getconfig('queue.workers', []); foreach ($workers as $queue => $options) { if (strpos($queue, '@') !== false) { [$queue, $connection] = explode('@', $queue); } else { $connection = null; } $this->workers[] = function (process $process) use ($options, $connection, $queue) { @cli_set_process_title("swoole queue: worker process"); /** @var worker $worker */ $worker = $this->container->make('queue.worker'); /** @var workeroptions $option */ $option = $this->container->make(workeroptions::class); $option->sleep = arr::get($options, "sleep", 3); $option->maxtries = arr::get($options, "tries", 0); $option->timeout = arr::get($options, "timeout", 60); $timer = timer::after($option->timeout * 1000, function () use ($process) { $process->exit(); }); $worker->runnextjob($connection, $queue, $option); timer::clear($timer); }; } } /** * 注册事件 */ protected function listenforevents() { $this->container->make('events')->listen(jobfailed::class, function (jobfailed $event) { $this->writeoutput($event->job); $this->logfailedjob($event); }); } /** * 记录失败任务 * @param jobfailed $event */ protected function logfailedjob(jobfailed $event) { $this->container['queue.failer']->log( $event->connection, $event->job->getqueue(), $event->job->getrawbody(), $event->exception ); } /** * write the status output for the queue worker. * * @param job $job * @param $status */ protected function writeoutput(job $job, $status) { switch ($status) { case 'starting': $this->writestatus($job, 'processing', 'comment'); break; case 'success': $this->writestatus($job, 'processed', 'info'); break; case 'failed': $this->writestatus($job, 'failed', 'error'); break; } } /** * format the status output for the queue worker. * * @param job $job * @param string $status * @param string $type * @return void */ protected function writestatus(job $job, $status, $type) { $this->output->writeln(sprintf( "<{$type}>[%s][%s] %s</{$type}> %s", date('y-m-d h:i:s'), $job->getjobid(), str_pad("{$status}:", 11), $job->getname() )); }}
增加crmebserviceprovider类
<?phpnamespace crmeb\swoole;use illuminate\contracts\debug\exceptionhandler;use illuminate\contracts\http\kernel;use crmeb\swoole\command\httpservercommand;use illuminate\queue\worker;use swooletw\http\httpserviceprovider;use swooletw\http\middleware\accesslog;use swooletw\http\server\manager;/** * class crmebserviceprovider * @package crmeb\swoole */class crmebserviceprovider extends httpserviceprovider{ /** * register manager. * * @return void */ protected function registermanager() { $this->app->singleton(manager::class, function ($app) { return new manager($app, 'laravel'); }); $this->app->alias(manager::class, 'swoole.manager'); $this->app->singleton('queue.worker', function ($app) { $isdownformaintenance = function () { return $this->app->isdownformaintenance(); }; return new worker( $app['queue'], $app['events'], $app[exceptionhandler::class], $isdownformaintenance ); }); } /** * boot websocket routes. * * @return void */ protected function bootwebsocketroutes() { require base_path('vendor/swooletw/laravel-swoole') . '/routes/laravel_routes.php'; } /** * register access log middleware to container. * * @return void */ protected function pushaccesslogmiddleware() { $this->app->make(kernel::class)->pushmiddleware(accesslog::class); } /** * register commands. */ protected function registercommands() { $this->commands([ httpservercommand::class, ]); } /** * merge configurations. */ protected function mergeconfigs() { $this->mergeconfigfrom(base_path('vendor/swooletw/laravel-swoole') . '/config/swoole_http.php', 'swoole_http'); $this->mergeconfigfrom(base_path('vendor/swooletw/laravel-swoole') . '/config/swoole_websocket.php', 'swoole_websocket'); } /** * publish files of this package. */ protected function publishfiles() { $this->publishes([ base_path('vendor/swooletw/laravel-swoole') . '/config/swoole_http.php' => base_path('config/swoole_http.php'), base_path('vendor/swooletw/laravel-swoole') . '/config/swoole_websocket.php' => base_path('config/swoole_websocket.php'), base_path('vendor/swooletw/laravel-swoole') . '/routes/websocket.php' => base_path('routes/websocket.php'), ], 'laravel-swoole'); }}
然后再把\crmeb\swoole\crmebserviceprovider::class放入config/app.php中的providers中加载重写了swoole的命令启动方式
配置config/swoole_http.php
return [ 'queue' => [ //是否开启自动消费队列 'enabled' => true, 'workers' => [ //队列名称 'crmeb' => [] ] ],];
输入命令:
php artisan crmeb:http restart
swoole启动后就可以自动消费队列了。
相关推荐:最新的五个laravel视频教程
以上就是laravel8中laravel-swoole的扩展不兼容消息队列怎么办?的详细内容。
