| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136 |
- <?php
- declare(strict_types=1);
- namespace SixShop\System\Command;
- use SixShop\System\Event\CrontabWorkerStartEvent;
- use think\console\Command;
- use think\console\Input;
- use think\console\input\Argument;
- use think\console\input\Option;
- use think\console\Output;
- use SixShop\System\Workerman\Worker as SixShopWorker;
- use Workerman\Crontab\Crontab;
- use Workerman\Timer;
- use Workerman\Worker;
- class WorkerCommand extends Command
- {
- public function __construct(private \think\queue\Worker $queueWorker)
- {
- parent::__construct();
- }
- public function configure(): void
- {
- $this->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,
- );
- }
- }
|