setName('system:worker') ->addArgument('action', Argument::OPTIONAL, 'start|stop|restart|reload|status|connections', 'start') ->addOption('daemon', 'd', Option::VALUE_NONE, 'daemon mode') ->addOption('grace', 'g', Option::VALUE_NONE, 'graceful shutdown/reload') ->setDescription('Workerman Server for ThinkPHP (custom entry)'); } protected function execute(Input $input, Output $output): void { SixShopWorker::setArgv($this->inputToArgv($input)); SixShopWorker::setApp($this->app); $this->addCronWorker(); $this->addQueueWorker(); SixShopWorker::runAll(); } /** * @param Input $input * @return array */ private function inputToArgv(Input $input): array { $argv = [PHP_BINARY, 'worker', $input->getArgument('action')]; $daemon = $input->getOption('daemon'); if ($daemon) { $argv[] = '-d'; } $grace = $input->getOption('grace'); if ($grace) { $argv[] = '-g'; } return $argv; } /** * @return void */ private function addCronWorker(): void { SixShopWorker::addWorker( onWorkerStart: function (Worker $worker) { $this->app->event->trigger(CrontabWorkerStartEvent::class, $worker); }, name: 'crontab', count: 1, onWorkerStop: function (Worker $worker) { foreach (Crontab::getAll() as $crontab) { /** @var Crontab $crontab */ $crontab->destroy(); } } ); } /** * @return void */ private function addQueueWorker(): void { SixShopWorker::addWorker( onWorkerStart: function (Worker $worker) { $shouldExit = false; $logger = $this->app->log; $connection = $this->app->config->get('queue.default'); $queue = $this->app->config->get("queue.connections.{$connection}.queue", 'default'); $delay = 0; $sleep = 2; $tries = 3; $memoryLimit = 128; // 内存上限MB $queueWorker = $this->queueWorker; $run = function () use (&$run, &$shouldExit, $queueWorker, $connection, $queue, $delay, $sleep, $tries, $memoryLimit, $logger, $worker) { if ($shouldExit) { $logger->info(sprintf("[pid:%d] 收到退出标记,停止消费调度", $worker->pid)); return; } try { $queueWorker->runNextJob($connection, $queue, $delay, $sleep, $tries); } catch (\Throwable $e) { $logger->error(sprintf( "[pid:%d] 队列底层异常:%s \n堆栈:%s", $worker->pid, $e->getMessage(), $e->getTraceAsString() )); } $usedMb = memory_get_usage(true) / 1024 / 1024; if ($usedMb > $memoryLimit) { $logger->info(sprintf( "[pid:%d] 内存超限 %.2fMB > %dMB,进程标记退出自愈", $worker->pid, $usedMb, $memoryLimit )); $shouldExit = true; return; } Timer::add(1, $run, [], false); }; Timer::add(1, $run, [], false); }, name: 'queue', count: 1, ); } }