WorkerCommand.php 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136
  1. <?php
  2. declare(strict_types=1);
  3. namespace SixShop\System\Command;
  4. use SixShop\System\Event\CrontabWorkerStartEvent;
  5. use think\console\Command;
  6. use think\console\Input;
  7. use think\console\input\Argument;
  8. use think\console\input\Option;
  9. use think\console\Output;
  10. use SixShop\System\Workerman\Worker as SixShopWorker;
  11. use Workerman\Crontab\Crontab;
  12. use Workerman\Timer;
  13. use Workerman\Worker;
  14. class WorkerCommand extends Command
  15. {
  16. public function __construct(private \think\queue\Worker $queueWorker)
  17. {
  18. parent::__construct();
  19. }
  20. public function configure(): void
  21. {
  22. $this->setName('system:worker')
  23. ->addArgument('action', Argument::OPTIONAL, 'start|stop|restart|reload|status|connections', 'start')
  24. ->addOption('daemon', 'd', Option::VALUE_NONE, 'daemon mode')
  25. ->addOption('grace', 'g', Option::VALUE_NONE, 'graceful shutdown/reload')
  26. ->setDescription('Workerman Server for ThinkPHP (custom entry)');
  27. }
  28. protected function execute(Input $input, Output $output): void
  29. {
  30. SixShopWorker::setArgv($this->inputToArgv($input));
  31. SixShopWorker::setApp($this->app);
  32. $this->addCronWorker();
  33. $this->addQueueWorker();
  34. SixShopWorker::runAll();
  35. }
  36. /**
  37. * @param Input $input
  38. * @return array
  39. */
  40. private function inputToArgv(Input $input): array
  41. {
  42. $argv = [PHP_BINARY, 'worker', $input->getArgument('action')];
  43. $daemon = $input->getOption('daemon');
  44. if ($daemon) {
  45. $argv[] = '-d';
  46. }
  47. $grace = $input->getOption('grace');
  48. if ($grace) {
  49. $argv[] = '-g';
  50. }
  51. return $argv;
  52. }
  53. /**
  54. * @return void
  55. */
  56. private function addCronWorker(): void
  57. {
  58. SixShopWorker::addWorker(
  59. onWorkerStart: function (Worker $worker) {
  60. $this->app->event->trigger(CrontabWorkerStartEvent::class, $worker);
  61. },
  62. name: 'crontab',
  63. count: 1,
  64. onWorkerStop: function (Worker $worker) {
  65. foreach (Crontab::getAll() as $crontab) {
  66. /** @var Crontab $crontab */
  67. $crontab->destroy();
  68. }
  69. }
  70. );
  71. }
  72. /**
  73. * @return void
  74. */
  75. private function addQueueWorker(): void
  76. {
  77. SixShopWorker::addWorker(
  78. onWorkerStart: function (Worker $worker) {
  79. $shouldExit = false;
  80. $logger = $this->app->log;
  81. $connection = $this->app->config->get('queue.default');
  82. $queue = $this->app->config->get("queue.connections.{$connection}.queue", 'default');
  83. $delay = 0;
  84. $sleep = 2;
  85. $tries = 3;
  86. $memoryLimit = 128; // 内存上限MB
  87. $queueWorker = $this->queueWorker;
  88. $run = function () use (&$run, &$shouldExit, $queueWorker, $connection, $queue, $delay, $sleep, $tries, $memoryLimit, $logger, $worker) {
  89. if ($shouldExit) {
  90. $logger->info(sprintf("[pid:%d] 收到退出标记,停止消费调度", $worker->pid));
  91. return;
  92. }
  93. try {
  94. $queueWorker->runNextJob($connection, $queue, $delay, $sleep, $tries);
  95. } catch (\Throwable $e) {
  96. $logger->error(sprintf(
  97. "[pid:%d] 队列底层异常:%s \n堆栈:%s",
  98. $worker->pid,
  99. $e->getMessage(),
  100. $e->getTraceAsString()
  101. ));
  102. }
  103. $usedMb = memory_get_usage(true) / 1024 / 1024;
  104. if ($usedMb > $memoryLimit) {
  105. $logger->info(sprintf(
  106. "[pid:%d] 内存超限 %.2fMB > %dMB,进程标记退出自愈",
  107. $worker->pid,
  108. $usedMb,
  109. $memoryLimit
  110. ));
  111. $shouldExit = true;
  112. return;
  113. }
  114. Timer::add(1, $run, [], false);
  115. };
  116. Timer::add(1, $run, [], false);
  117. },
  118. name: 'queue',
  119. count: 1,
  120. );
  121. }
  122. }